---
title: Kafka Integration
slug: litmusedge/kafka-integration
docTags: 
createdAt: 2024-08-13T15:35:15.335Z
---

To configure the Kafka integration, complete the following steps.

# Step 1: Add Device

Follow the steps to [Connect a device](docId\:pal6ABPZbrimdU9LvGJ30). The device will be used to store tags that will be eventually used to create outbound topics in the connector. Make sure to select the **Enable Data Store** checkbox.&#x20;

# Step 2: Add Tags

After connecting the device in Litmus Edge, you can [Add Tags](docId\:XgWOkQbTPevII7OR82LL0) to the device. Create tags that you want to use to create outbound topics for the connector.&#x20;

# Step 3: Add Kafka Connector

**To add the Kafka connector in Litmus Edge: &#x20;**

1. Navigate to **Integration**.
2. Click **Add a connector&#x20;**&#x69;con. 
   The *Add a connector&#x20;*&#x64;ialog box appears.
3. Select the **Kafka SSL&#x20;**&#x70;rovider from the drop-down menu.&#x20;
   In this guide, we are using Confluent Kafka as the Kafka broker, so we selected *Kafka SSL Integration*.
4. Complete the following information for the *Kafka SSL* connector as shown in the screenshot and click **Add/Update**.
   ::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/Lg2qbmz1NEZbuUCcaIZLM_image.png" size="80" width="734" height="770" position="center" caption="Edit a connector" showCaption="true"}
5. Click the **kafka SSL&#x20;**&#x63;onnector tile.&#x20;
   The connector *Dashboard* appears.
6. Click the **Topics&#x20;**&#x74;ab.
7. Click the **Import from DeviceHub** tags icon and select all the tags to send data to *Databricks*.
   ![](https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/S5iC2Vx6ao8YdpIv45-MR_image.png "Topics tab in Kafka Connector")

# Step 4: Set Up Databricks for Kafka Streaming&#x20;

In Databricks, you can configure *Kafka&#x20;*&#x63;onnector parameters from a **Python notebook** file. You will need to input information about the Broker address, Topic name, and API key and secret.

**To set up Databricks for Kafka streaming:**

1\. Navigate to **Databricks&#x20;**&#x77;orkspace using your log in credentials.

2\. Create a new notebook in *Databricks*. Follow the [Databricks Kafka Streaming Guide](https://docs.databricks.com/en/connect/streaming/kafka.html) to set up *Kafka* streaming.&#x20;

3\. Once the data is loaded, it will be in the following format:

::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/73AHR-ihz_gJP_IQd_P8q_image.png" size="50" width="340" height="172" position="center" caption="Dataframe format" showCaption="true"}

The payload sent by LE will be in the **value&#x20;**&#x66;ield in binary format.

# Example Notebook for Kafka Streams

:::hint{type="info"}
**Note:&#x20;**&#x54;he user can create their own notebooks and customize them as needed. This example notebook will explain how to ingest the standard DeviceHub JSON Payload.
:::

**Open a new notebook and follow the steps:**

1\. Import the necessary libraries and dependencies to handle *Kafka&#x20;*&#x64;ata streams and process data using *PySpark*.&#x20;

```python
import pyspark.sql.types as T
import pyspark.sql.functions as F
from pyspark.sql.types import StructType, StructField, StringType, IntegerType, DoubleType, BooleanType, TimestampType, LongType
from pyspark.sql.functions import from_unixtime, from_json, col, cast
```

2\. Define the credentials and *Kafka broker* information that will be used to authenticate and connect to the *Kafka server*.&#x20;

:::hint{type="info"}
**Note:&#x20;**&#x45;nd-users should use secure methods of passing their *Kafka&#x20;*&#x63;redentials.
:::

```python
confluentBootstrapserver = "<KAFKA_BROKER_SERVER>"
KafkaTopic = "<KAFKA_TOPIC>"
apiKey = "<API_KEY>"
apiSecret = "<API_SECRET>"
```

3\. Configure **Spark&#x20;**&#x74;o read data streams from the *Kafka topic.*&#x20;

:::hint{type="info"}
**Note:&#x20;**&#x45;nsure to include the *Kafka server*, *security protocol*, and *authentication&#x20;*&#x64;etails.
:::

```python
df = spark.readStream \
        .format("kafka") \
        .option("kafka.bootstrap.servers", confluentBootstrapserver) \
        .option("kafka.security.protocol", "SASL_SSL") \
        .option("kafka.sasl.mechanism", "PLAIN") \
        .option("kafka.sasl.jaas.config", f'kafkashaded.org.apache.kafka.common.security.plain.PlainLoginModule required username="{apiKey}" password="{apiSecret}";') \
        .option("subscribe", KafkaTopic) \
        .option("startingOffsets", "earliest") \
        .load()
```

This will output the following:

::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/aMSqpwcWJaTFZpiqVHN6J_image.png" size="50" width="340" height="172" position="center" caption="Dataframe format" showCaption="true"}

:::hint{type="info"}
**Note:&#x20;**&#x54;his guide covers the DeviceHub JSON Payload. Please update the notebook if you are sending a different payload to Databricks.
:::

4\. Define the schema that matches the expected structure of the incoming data and parse the *Kafka stream* data accordingly.&#x20;

```python
# Define the schema
schema = StructType([
    StructField("deviceName", StringType()),
    StructField("tagName", StringType()),
    StructField("deviceID", StringType()),
    StructField("success", BooleanType()),
    StructField("datatype", StringType()),
    StructField("timestamp", LongType()),
    StructField("value", DoubleType()),
    StructField("metadata", StringType()),
    StructField("registerId", StringType()),
    StructField("description", StringType())
])

# Convert the binary 'value' column to string and then parse the JSON data
parsed_df = (df
             .withColumnsRenamed({'timestamp': 'kafka_time'})
             .withColumn("parsed_value", from_json(col("value").cast("string"), schema))
             .withColumn("LE_timestamp", from_unixtime(F.col('parsed_value.timestamp')/1000, "yyyy-MM-dd HH:mm:ss.SSS"))
             .select("key","topic", "kafka_time","partition","offset","parsed_value.*","LE_timestamp", )
)
```

This will output the following:

::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/vWXgtE2Cn9QF4JvhbXxvR_image.png" size="50" width="355" height="334" position="center" caption="Dataframe format" showCaption="true"}

5\. Specify the **Databricks&#x20;**&#x74;able where the parsed data will be stored.&#x20;

```python
table = "<DATABRICKS_TABLE_NAME>"


query = (parsed_df 
    .writeStream
    .format("delta")
    .outputMode("append")
    .option("checkpointLocation", "/tmp/delta/"+ table + "/_checkpoints/").toTable(table)
)
```

6\. Once the data is published to the *Databricks&#x20;*&#x74;able, you can query the table to verify that the data ingestion process is successful.&#x20;

![](https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/YMEkUvLPO_dkgIlfJ_-cr_image.png "Databricks table")

