Note
Access to this page requires authorization. You can try signing in or changing directories.
Access to this page requires authorization. You can try changing directories.
Important
This feature is in Beta. To use it, a workspace admin must turn on Zerobus Ingest MQTT Endpoint from the Previews page. See Manage Azure 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:
- MQTT enabled for your workspace.
- An existing Unity Catalog Delta table. Zerobus Ingest doesn't create a table when a client connects.
- A service principal with
USE CATALOG,USE SCHEMA,SELECT, andMODIFYprivileges for the target table. See Create a service principal and grant permissions. - The workspace URL, workspace ID, region, and Zerobus Ingest endpoint. See Get your workspace URL and Zerobus Ingest endpoint.
- Outbound TLS connectivity to the Zerobus Ingest hostname on port
8883.
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:
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>.azuredatabricks.net
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 Azure 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
PUBACKafter Zerobus Ingest makes the record durable. Treat the record as acknowledged only when thePUBACKmatches 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:
CREATE TABLE main.default.air_quality (
device_name STRING,
temp INT,
humidity INT
);
Install the dependencies:
pip install "paho-mqtt>=2,<3" requests
Set the required environment variables. Set ZEROBUS_MQTT_HOST to the MQTT hostname described in Connection model.
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:
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. |
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 missingPUBACKdoesn'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.