Pular para o conteúdo principal

Use MQTT com o Zerobus Ingest

info

Beta

Este recurso está em Beta. Para usá-lo, o administrador do workspace deve ativar Zerobus Ingest MQTT Endpoint na página Previews . Consulte Gerenciar prévias do Databricks.

Zerobus Ingest fornece um endpoint MQTT v5 para aplicativos que já publicam mensagens MQTT. Cada mensagem contém um objeto JSON e é gravada diretamente em uma tabela Delta do Unity Catalog existente. Esta interface é ideal para aplicativos de Internet das Coisas (IoT), edge e telemetria que usam um cliente MQTT v5 e não exigem recursos de inscrição de broker.

Prerequisites​

Antes de se conectar, certifique-se de ter os seguintes recursos e acesso:

  • MQTT habilitado para o seu workspace.
  • Uma tabela Delta do Unity Catalog existente. O Zerobus Ingest não cria uma tabela quando um cliente se conecta.
  • Um Service Principal com os privilégios USE CATALOG, USE SCHEMA, SELECT e MODIFY para a tabela de destino. Consulte Criar um Service Principal e conceder permissões.
  • The workspace URL, workspace ID, region, and Zerobus Ingest endpoint. See Get your workspace URL and Zerobus Ingest endpoint.
  • Conectividade TLS de saída para o hostname do Zerobus Ingest na porta 8883.

Confirme se o Zerobus Ingest está disponível na região do seu workspace. Consulte disponibilidade regional do Zerobus Ingest. A disponibilidade do MQTT pode diferir da disponibilidade geral do Zerobus Ingest.

Modelo de conexão​

Cada conexão MQTT tem como alvo uma tabela. Defina a tabela quando o cliente enviar CONNECT e use o mesmo nome de tabela de três níveis como o tópico para cada pacote PUBLISH:

Text
catalog.schema.table

Cada payload PUBLISH deve ser um objeto JSON codificado em UTF-8 que corresponda ao esquema da tabela de destino. Uma conexão não pode alternar tabelas. Abra outra conexão para publicar em outra tabela.

O endpoint do MQTT usa TLS na porta 8883 com o seguinte formato de hostname:

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

Este é o mesmo hostname do endpoint do servidor Zerobus Ingest. Configure seu cliente MQTT com este hostname e a porta 8883, sem um esquema https://.

Autenticação​

Autentique durante CONNECT com estas propriedades de usuário do MQTT v5:

Propriedade do usuário

Valor

Authorization

Bearer <token>, em que <token> é um token OAuth do Databricks com escopo para a tabela de destino.

x-databricks-zerobus-table-name

O nome completo da tabela de destino, como main.default.air_quality.

Propriedade do usuário

Valor

Authorization

Bearer <token>, em que <token> é um token OAuth do Databricks com escopo para a tabela de destino.

x-databricks-zerobus-table-name

O nome completo da tabela de destino, como main.default.air_quality.

Obtenha o token OAuth com o fluxo de credenciais de cliente do Service Principal. Inclua authorization_details para o catálogo, esquema e tabela, e use o recurso zerobusDirectWriteApi para o seu Workspace. O exemplo em Python demonstra essa troca.

OAuth credentials apply only when a connection começa. 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.

Durabilidade e confirmações​

O Zerobus Ingest oferece suporte aos seguintes níveis de qualidade de serviço (QoS) do MQTT:

  • O QoS 0 é de melhor esforço. O Zerobus Ingest tenta enviar registros válidos por meio do pipeline de ingestão normal, mas não envia uma confirmação de sucesso nem uma resposta de falha por registro. Um registro pode ser rejeitado ou perdido sem notificação ao cliente, portanto, uma conexão aberta não confirma a entrega. O QoS 0 não fornece entrega de pelo menos uma vez. Use-o apenas quando seu aplicativo puder tolerar a perda de registros.
  • O QoS 1 retorna um PUBACK após o Zerobus Ingest tornar o registro durável. Trate o registro como confirmado somente quando o PUBACK corresponder ao ID do pacote PUBLISH e tiver um código de motivo bem-sucedido.

Um PUBACK bem-sucedido confirma a durabilidade, não a visibilidade imediata na tabela Delta. O Zerobus Ingest materializa registros duráveis na tabela de forma assíncrona. Consulte Latência.

A entrega do QoS 1 pode conter duplicatas. Se uma conexão for fechada antes de seu cliente receber PUBACK, o registro poderá já ser durável. A nova tentativa desse registro pode criar uma duplicata. Os aplicativos que tentam novamente registrar registros não confirmados devem tolerar duplicatas ou deduplicá-las no lakehouse.

Um registro QoS 1 não tem confirmação se o código de motivo do seu PUBACK for malsucedido ou se a conexão for fechada antes da chegada do PUBACK. Aplique a política de novas tentativas do seu aplicativo aos registros sem confirmação. Um pacote PUBLISH maior que o limite de tamanho do pacote fecha a conexão e é sempre rejeitado; portanto, não tente enviá-lo novamente.

Exemplo: publicar um registro​

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

nota

Este exemplo introdutório desativa o comportamento de reconexão automática do Paho. Em um aplicativo de execução longa, reconecte-se com um token OAuth recém-obtido. Acompanhe os registros sem um PUBACK bem-sucedido e aplique sua própria política de nova tentativa, contabilizando o risco de duplicata descrito em Durabilidade e confirmações.

O exemplo espera uma tabela com este esquema:

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

Instale as dependências:

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

Defina as variáveis de ambiente necessárias. Defina ZEROBUS_MQTT_HOST como o hostname MQTT descrito em 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>"

Execução o seguinte 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()

Limitações​

A interface MQTT tem os seguintes limites de protocolo:

Item

Apoiar

Versão de protocolo

Somente MQTT v5.

Qualidade do serviço

QoS 0 e QoS 1. O QoS 2 não é compatível.

Mensagens de QoS 1 em trânsito

O CONNACK Receive Maximum é de 4.096 mensagens por conexão. Os clientes devem respeitar o limite anunciado.

CONNECT tamanho do pacote

16 KiB para o pacote CONNECT inicial completo, incluindo propriedades.

Tamanho do pacote após CONNECT

64 KiB para o pacote MQTT completo, incluindo o cabeçalho fixo, o cabeçalho variável, as propriedades, o tópico, o ID do pacote e o payload. CONNACK anuncia esse limite como o tamanho máximo de pacote. Um pacote maior fecha a conexão com o código de motivo Packet too large.

ID do cliente

Obrigatório, não vazio e com no máximo 128 bytes. Use uma ID de cliente exclusiva para cada conexão ativa.

Estado da sessão

O servidor não retém nem resume o estado da sessão e define o Session Expiry Interval como 0.

Intervalo de keep-alive

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.

Link privado front-end

Não suportado. Conecte-se pelo endpoint público.

PUBLISH Propriedades

Somente o tópico, o QoS e o Payload afetam a ingestão. O servidor não persiste nem aplica Payload Format Indicator, Content Type, Correlation Data, Message Expiry Interval, Response Topic ou User Properties. O identificador da inscrição não é compatível.

Recursos do broker

Inscrições, retained messages, topic aliases, Last Will messages, and enhanced authentication aren't supported.

Item

Apoiar

Versão de protocolo

Somente MQTT v5.

Qualidade do serviço

QoS 0 e QoS 1. O QoS 2 não é compatível.

Mensagens de QoS 1 em trânsito

O CONNACK Receive Maximum é de 4.096 mensagens por conexão. Os clientes devem respeitar o limite anunciado.

CONNECT tamanho do pacote

16 KiB para o pacote CONNECT inicial completo, incluindo propriedades.

Tamanho do pacote após CONNECT

64 KiB para o pacote MQTT completo, incluindo o cabeçalho fixo, o cabeçalho variável, as propriedades, o tópico, o ID do pacote e o payload. CONNACK anuncia esse limite como o tamanho máximo de pacote. Um pacote maior fecha a conexão com o código de motivo Packet too large.

ID do cliente

Obrigatório, não vazio e com no máximo 128 bytes. Use uma ID de cliente exclusiva para cada conexão ativa.

Estado da sessão

O servidor não retém nem resume o estado da sessão e define o Session Expiry Interval como 0.

Intervalo de keep-alive

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.

Link privado front-end

Não suportado. Conecte-se pelo endpoint público.

PUBLISH Propriedades

Somente o tópico, o QoS e o Payload afetam a ingestão. O servidor não persiste nem aplica Payload Format Indicator, Content Type, Correlation Data, Message Expiry Interval, Response Topic ou User Properties. O identificador da inscrição não é compatível.

Recursos do broker

Inscrições, retained messages, topic aliases, Last Will messages, and enhanced authentication aren't supported.

Mantenha a conexão aberta ao publicar vários registros na mesma tabela.

Solução de problemas​

Use as seguintes verificações para diagnosticar falhas comuns:

  • Falhas de conexão: verifique o hostname específico da cloud, a porta 8883, a resolução de DNS, as regras de firewall de saída e a verificação do certificado TLS. Confirme também se o MQTT está habilitado para o Workspace.
  • CONNECT rejeitado: obtenha um novo token OAuth com escopo de tabela. Verifique os dois nomes de propriedade do usuário, o nome da tabela de três níveis, a tabela de destino e as concessões (grants) da Service Principal.
  • 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.
  • Desconexões: verifique se há um tópico incompatível, um pacote maior que 64 KiB, um recurso MQTT sem suporte, tráfego de keep-alive perdido ou o tempo de vida limitado da conexão no lado do servidor. Obtenha um token novo antes de reconectar e aplique a política de repetição e deduplicação do seu aplicativo aos registros não reconhecidos.