Skip to main content

Use MQTT with Zerobus Ingest

Beta

This feature is in Beta. To use it, a workspace admin must turn on Zerobus Ingest MQTT Endpoint from the Previews page. See Manage Databricks previews.

Zerobus Ingest provides an MQTT v5 endpoint for applications that already publish MQTT messages. Each message contains one JSON object and writes directly to an existing Unity Catalog Delta table. This interface is a good fit for Internet of Things (IoT), edge, and telemetry applications that use an MQTT v5 client and don't require broker subscription features.

Prerequisites​

Before you connect, make sure you have the following resources and access:

Confirm that Zerobus Ingest is available in your workspace region. See Zerobus Ingest regional availability. MQTT availability can differ from general Zerobus Ingest availability.

Connection model​

Each MQTT connection targets one table. Set the table when the client sends CONNECT, and use the same three-level table name as the topic for every PUBLISH packet:

Text
catalog.schema.table

Each PUBLISH payload must be one UTF-8 encoded JSON object that matches the target table schema. A connection can't switch tables. Open another connection to publish to another table.

The MQTT endpoint uses TLS on port 8883 with the following hostname format:

<workspace-id>.zerobus.<region>.cloud.databricks.com

This is the same hostname as the Zerobus Ingest server endpoint. Configure your MQTT client with this hostname and port 8883, without an https:// scheme.

Authentication​

Authenticate during CONNECT with these MQTT v5 User Properties:

User Property

Value

Authorization

Bearer <token>, where <token> is a Databricks OAuth token scoped to the target table.

x-databricks-zerobus-table-name

The full target table name, such as main.default.air_quality.

User Property

Value

Authorization

Bearer <token>, where <token> is a Databricks OAuth token scoped to the target table.

x-databricks-zerobus-table-name

The full target table name, such as main.default.air_quality.

Obtain the OAuth token with the service principal client credentials flow. Include authorization_details for the catalog, schema, and table, and use the zerobusDirectWriteApi resource for your workspace. The Python example demonstrates this exchange.

OAuth credentials apply only when a connection starts. Connections have a bounded server-side lifetime, after which Zerobus Ingest disconnects the client. When you reconnect, obtain a fresh token and send it in the new CONNECT properties. MQTT enhanced authentication and in-connection reauthentication aren't supported.

Durability and acknowledgments​

Zerobus Ingest supports the following MQTT quality of service (QoS) levels:

  • QoS 0 is best-effort. Zerobus Ingest attempts to send valid records through the normal ingestion pipeline but sends neither a success acknowledgment nor a per-record failure response. A record can be rejected or lost without client notification, so an open connection doesn't confirm delivery. QoS 0 doesn't provide at-least-once delivery. Use it only when your application can tolerate record loss.
  • QoS 1 returns a PUBACK after Zerobus Ingest makes the record durable. Treat the record as acknowledged only when the PUBACK matches the PUBLISH packet ID and has a successful reason code.

A successful PUBACK confirms durability, not immediate visibility in the Delta table. Zerobus Ingest materializes durable records into the table asynchronously. See Latency.

QoS 1 delivery can contain duplicates. If a connection closes before your client receives PUBACK, the record might already be durable. Retrying that record can create a duplicate. Applications that retry unacknowledged records must tolerate duplicates or deduplicate them in the lakehouse.

A QoS 1 record is unacknowledged if its PUBACK has an unsuccessful reason code or if the connection closes before the PUBACK arrives. Apply your application's retry policy to unacknowledged records. A PUBLISH packet larger than the packet size limit closes the connection and is always rejected, so don't retry it.

Example: Publish a record​

The following example uses Paho MQTT 2.x and requests.

note

This introductory example disables Paho's automatic reconnect behavior. In a long-running application, reconnect with a newly fetched OAuth token. Track records without a successful PUBACK and apply your own retry policy, accounting for the duplicate risk described in Durability and acknowledgments.

The example expects a table with this schema:

SQL
CREATE TABLE main.default.air_quality (
device_name STRING,
temp INT,
humidity INT
);

Install the dependencies:

Bash
pip install "paho-mqtt>=2,<3" requests

Set the required environment variables. Set ZEROBUS_MQTT_HOST to the MQTT hostname described in Connection model.

Bash
export DATABRICKS_WORKSPACE_URL="https://<databricks-instance>"
export DATABRICKS_WORKSPACE_ID="<workspace-id>"
export DATABRICKS_CLIENT_ID="<service-principal-client-id>"
export DATABRICKS_CLIENT_SECRET="<service-principal-client-secret>"
export DATABRICKS_TABLE_NAME="main.default.air_quality"
export ZEROBUS_MQTT_HOST="<cloud-specific-zerobus-hostname>"

Run the following script:

Python
import json
import os
import ssl
import threading
import uuid

import paho.mqtt.client as mqtt
import requests

CONNECT_TIMEOUT_SECONDS = 10
CONNACK_TIMEOUT_SECONDS = 30
PUBACK_TIMEOUT_SECONDS = 30


def required_env(name):
value = os.environ.get(name)
if not value:
raise RuntimeError(f"Set the {name} environment variable")
return value


workspace_url = required_env("DATABRICKS_WORKSPACE_URL").rstrip("/")
workspace_id = required_env("DATABRICKS_WORKSPACE_ID")
client_id = required_env("DATABRICKS_CLIENT_ID")
client_secret = required_env("DATABRICKS_CLIENT_SECRET")
table_name = required_env("DATABRICKS_TABLE_NAME")
mqtt_host = required_env("ZEROBUS_MQTT_HOST")


def fetch_zerobus_token():
catalog, schema, _ = table_name.split(".")
authorization_details = [
{
"type": "unity_catalog_privileges",
"privileges": ["USE CATALOG"],
"object_type": "CATALOG",
"object_full_path": catalog,
},
{
"type": "unity_catalog_privileges",
"privileges": ["USE SCHEMA"],
"object_type": "SCHEMA",
"object_full_path": f"{catalog}.{schema}",
},
{
"type": "unity_catalog_privileges",
"privileges": ["SELECT", "MODIFY"],
"object_type": "TABLE",
"object_full_path": table_name,
},
]
response = requests.post(
f"{workspace_url}/oidc/v1/token",
auth=(client_id, client_secret),
data={
"grant_type": "client_credentials",
"scope": "all-apis",
"resource": f"api://databricks/workspaces/{workspace_id}/zerobusDirectWriteApi",
"authorization_details": json.dumps(authorization_details),
},
timeout=30,
)
response.raise_for_status()
return response.json()["access_token"]


connack_received = threading.Event()
puback_received = threading.Event()
connack_result = {}
puback_result = {}
disconnect_result = {}


def on_connect(client, userdata, flags, reason_code, properties):
connack_result["reason_code"] = reason_code
connack_received.set()


def on_publish(client, userdata, mid, reason_code, properties):
puback_result["mid"] = mid
puback_result["reason_code"] = reason_code
puback_received.set()


def on_disconnect(client, userdata, disconnect_flags, reason_code, properties):
disconnect_result["reason_code"] = reason_code
connack_received.set()
puback_received.set()


connect_properties = mqtt.Properties(mqtt.PacketTypes.CONNECT)
connect_properties.UserProperty = [
("Authorization", f"Bearer {fetch_zerobus_token()}"),
("x-databricks-zerobus-table-name", table_name),
]

client = mqtt.Client(
mqtt.CallbackAPIVersion.VERSION2,
client_id=f"zerobus-mqtt-{uuid.uuid4()}",
protocol=mqtt.MQTTv5,
reconnect_on_failure=False,
)
client.on_connect = on_connect
client.on_publish = on_publish
client.on_disconnect = on_disconnect
client.connect_timeout = CONNECT_TIMEOUT_SECONDS
client.tls_set(cert_reqs=ssl.CERT_REQUIRED, tls_version=ssl.PROTOCOL_TLS_CLIENT)

loop_started = False
try:
result = client.connect(
mqtt_host,
port=8883,
keepalive=60,
clean_start=True,
properties=connect_properties,
)
if result != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"MQTT connect failed: {mqtt.error_string(result)}")

result = client.loop_start()
if result != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"MQTT network loop failed: {mqtt.error_string(result)}")
loop_started = True

if not connack_received.wait(CONNACK_TIMEOUT_SECONDS):
raise TimeoutError("Timed out waiting for CONNACK")
if "reason_code" not in connack_result:
raise RuntimeError(
f"Disconnected before CONNACK: {disconnect_result['reason_code']}"
)
connack_reason = connack_result["reason_code"]
if connack_reason.is_failure:
raise RuntimeError(f"CONNECT rejected: {connack_reason}")

record = {"device_name": "sensor-1", "temp": 22, "humidity": 55}
message = client.publish(table_name, json.dumps(record), qos=1)
if message.rc != mqtt.MQTT_ERR_SUCCESS:
raise RuntimeError(f"PUBLISH failed: {mqtt.error_string(message.rc)}")

if not puback_received.wait(PUBACK_TIMEOUT_SECONDS):
raise TimeoutError("Timed out waiting for PUBACK")
if "reason_code" not in puback_result:
raise RuntimeError(
f"Disconnected before PUBACK: {disconnect_result['reason_code']}"
)
if puback_result["mid"] != message.mid:
raise RuntimeError("Received PUBACK for a different PUBLISH packet")
puback_reason = puback_result["reason_code"]
if puback_reason.value != 0:
raise RuntimeError(f"PUBLISH rejected: {puback_reason}")

print(f"Record {message.mid} is durable")
finally:
client.disconnect()
if loop_started:
client.loop_stop()

Limitations​

The MQTT interface has the following protocol limits:

Item

Support

Protocol version

MQTT v5 only.

Quality of service

QoS 0 and QoS 1. QoS 2 isn't supported.

In-flight QoS 1 messages

The CONNACK Receive Maximum is 4,096 messages per connection. Clients must honor the advertised limit.

CONNECT packet size

16 KiB for the complete initial CONNECT packet, including properties.

Packet size after CONNECT

64 KiB for the complete MQTT packet, including the fixed header, variable header, properties, topic, packet ID, and payload. CONNACK advertises this limit as the Maximum Packet Size. A larger packet closes the connection with reason code Packet too large.

Client ID

Required, non-empty, and no more than 128 bytes. Use a unique client ID for each active connection.

Session state

The server doesn't retain or resume session state and sets the Session Expiry Interval to 0.

Keep-alive interval

Between 5 and 300 seconds. If the CONNECT value is 0, the server uses 60 seconds; otherwise, the server limits the value to this range. Honor the Server Keep Alive value returned in CONNACK. The server closes the connection if it receives no complete packet for 1.5 times the effective interval. A packet counts as received only when its last byte arrives, so on slow links, choose a keep-alive interval long enough to send your largest record.

Front-end Private Link

Not supported. Connect over the public endpoint instead.

PUBLISH properties

Only the topic, QoS, and payload affect ingestion. The server doesn't persist or apply Payload Format Indicator, Content Type, Correlation Data, Message Expiry Interval, Response Topic, or User Properties. Subscription Identifier isn't supported.

Broker features

Subscriptions, retained messages, topic aliases, Last Will messages, and enhanced authentication aren't supported.

Item

Support

Protocol version

MQTT v5 only.

Quality of service

QoS 0 and QoS 1. QoS 2 isn't supported.

In-flight QoS 1 messages

The CONNACK Receive Maximum is 4,096 messages per connection. Clients must honor the advertised limit.

CONNECT packet size

16 KiB for the complete initial CONNECT packet, including properties.

Packet size after CONNECT

64 KiB for the complete MQTT packet, including the fixed header, variable header, properties, topic, packet ID, and payload. CONNACK advertises this limit as the Maximum Packet Size. A larger packet closes the connection with reason code Packet too large.

Client ID

Required, non-empty, and no more than 128 bytes. Use a unique client ID for each active connection.

Session state

The server doesn't retain or resume session state and sets the Session Expiry Interval to 0.

Keep-alive interval

Between 5 and 300 seconds. If the CONNECT value is 0, the server uses 60 seconds; otherwise, the server limits the value to this range. Honor the Server Keep Alive value returned in CONNACK. The server closes the connection if it receives no complete packet for 1.5 times the effective interval. A packet counts as received only when its last byte arrives, so on slow links, choose a keep-alive interval long enough to send your largest record.

Front-end Private Link

Not supported. Connect over the public endpoint instead.

PUBLISH properties

Only the topic, QoS, and payload affect ingestion. The server doesn't persist or apply Payload Format Indicator, Content Type, Correlation Data, Message Expiry Interval, Response Topic, or User Properties. Subscription Identifier isn't supported.

Broker features

Subscriptions, retained messages, topic aliases, Last Will messages, and enhanced authentication aren't supported.

Keep the connection open when publishing multiple records to the same table.

Troubleshooting​

Use the following checks to diagnose common failures:

  • Connection failures: Verify the cloud-specific hostname, port 8883, DNS resolution, outbound firewall rules, and TLS certificate verification. Also confirm that MQTT is enabled for the workspace.
  • Rejected CONNECT: Obtain a new table-scoped OAuth token. Verify both User Property names, the three-level table name, the target table, and the service principal grants.
  • Unsuccessful or missing PUBACK: Don't report the record as durable. Confirm that the topic equals the table name and that the payload is one schema-matching UTF-8 JSON object. A missing PUBACK doesn't prove that the record failed to become durable.
  • Disconnects: Check for a mismatched topic, a packet larger than 64 KiB, an unsupported MQTT feature, missed keep-alive traffic, or the bounded server-side connection lifetime. Obtain a fresh token before reconnecting, and apply your application's retry and deduplication policy to unacknowledged records.