---
title: GCP Pub/Sub Integration
slug: litmusedge/gcp-pubsub-integration
docTags: 
createdAt: 2024-08-13T15:37:27.421Z
---

To configure the GCP Pub/Sub 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 Google Cloud Pub/Sub 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 **Google Cloud Pub/Sub&#x20;**&#x70;rovider from the drop-down list.
4. Complete the following information for the *Google Cloud Pub/Sub* connector as shown in the screenshot and click **Update**.
   ::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/oKVDocR1dTli-tufu3gt4_image.png" size="80" width="944" height="933" position="center" caption="Edit a connector dialog box" showCaption="true"}
5. Click the **Google Cloud Pub/Sub&#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/H665vloLztpoeg-1udojP_image.png "Topics tab in Google Cloud Pub/Sub connector")

# Step 4: Set Up Databricks for GCP Pub/Sub Streaming&#x20;

In Databricks, you can configure *GCP Pub/Sub&#x20;*&#x63;onnector parameters from a **Python notebook** file. You need to authenticate using service account credentials and specify Pub/Sub topic.

**To set up Databricks for GCP Pub/Sub Streaming:**

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

2\. Create a new notebook in *Databricks*. Follow the [Databricks Pub/Sub Streaming Guide](https://docs.databricks.com/en/connect/streaming/pub-sub.html) to set up *Google Pub/Sub* 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/li1qCy_FJ3DEx6-r1P0qB_image.png" size="50" width="308" height="120" position="center" caption="Dataframe format" showCaption="true"}

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

# Example Notebook for GCP Pub/Sub 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 *Pub/Sub* data 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
```

2\. Define the credentials and *GCP Pub/Sub* information that will be used to authenticate and connect to the *Pub/Sub service*.

:::hint{type="info"}
**Note:** End-users should use secure methods of passing their *Pub/Sub&#x20;*&#x63;redentials.
:::

```python
# for both topic and subscription topic, you only need to pass the topic and sub-topic name. Not the full path
# example = databricks-test
# not the full path in GCP PubSub, which is "projects/litmus-sales-enablement/topics/databricks-test"

authOptions = {"<GCP_PUBSUB_SA_KEY>"}

pubsub_topic = "<GCP_PUBSUB_TOPIC>" 
pubsub_subscription = "<GCP_PUBSUB_SUB_TOPIC>"
projectId = "<GCP_PROJECT_ID>"
```

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

:::hint{type="info"}
**Note:&#x20;**&#x45;nsure to include the *subscription ID, topic ID, project ID,&#x20;*&#x61;n&#x64;*&#x20;authentication options*.
:::

```python
df = (spark.readStream
  .format("pubsub")
  .option("subscriptionId", pubsub_subscription) 
  .option("topicId", pubsub_topic)
  .option("projectId", projectId)
  .options(**authOptions)
  .load()
)
```

This will output the following:

::Image[]{src="https://api.archbee.com/api/optimize/SSUUxKZUk9bFTEPNn_6Zo/7JE6rz2qe81SAa3I4urn5_image.png" size="50" width="308" height="120" 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 *Pub/Sub stream* data accordingly.&#x20;
This involves decoding the payload, extracting attributes, and structuring the parsed data into a DataFrame.&#x20;

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

table = "<DATABRICKS_TABLE_NAME>"

query = (df
  .withColumn("payload", F.decode(F.col("payload"), "UTF-8"))
  .withColumn("attributes", F.from_json(F.col("attributes"), T.MapType(T.StringType(), T.StringType())))
  .withColumn('parsed_value', F.from_json(F.col('payload').cast('string'), schema))
  .withColumn("LE_timestamp", from_unixtime(F.col('parsed_value.timestamp')/1000, "yyyy-MM-dd HH:mm:ss.SSS"))
  .select("messageId","attributes","publishTimestampInMillis",'parsed_value.*', "LE_timestamp")
  .writeStream
  .format("delta")
  .outputMode("append")
  .option("checkpointLocation", "/tmp/delta/"+ table + "/_checkpoints/").toTable(table)
)
```

5\. 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/BOMapqxwCitLv9cEieMUv_image.png "Databricks table")

