Use MQTT com o Zerobus Ingest
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,SELECTeMODIFYpara 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:
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 |
|---|---|
|
|
| O nome completo da tabela de destino, como |
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
PUBACKapós o Zerobus Ingest tornar o registro durável. Trate o registro como confirmado somente quando oPUBACKcorresponder 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.
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:
CREATE TABLE main.default.air_quality (
device_name STRING,
temp INT,
humidity INT
);
Instale as dependências:
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.
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:
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 |
| 16 KiB para o pacote |
Tamanho do pacote após | 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. |
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 |
Intervalo de keep-alive | Between 5 and 300 seconds. If the |
Link privado front-end | Não suportado. Conecte-se pelo endpoint público. |
| 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. CONNECTrejeitado: 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 missingPUBACKdoesn'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.