Utilisez MQTT avec Zerobus Ingest
Bêta
Cette fonctionnalité est en version bêta. Pour l'utiliser, un administrateur de workspace doit activer Zerobus Ingest MQTT Endpoint depuis la page Previews . Consultez Gérer les aperçus Databricks.
Zerobus Ingest fournit un endpoint MQTT v5 pour les applications qui publient déjà des messages MQTT. Chaque message contient un objet JSON et écrit directement dans une table Delta du Unity Catalog existante. Cette interface convient parfaitement aux applications de l’Internet des objets (IoT), de périphérie (« edge ») et de télémétrie qui utilisent un client MQTT v5 et ne nécessitent pas de fonctionnalités d’abonnement de courtier.
Prérequis
Avant de vous connecter, assurez-vous de disposer des ressources et des accès suivants :
- MQTT activé pour votre Workspace.
- Une table Delta Unity Catalog existante. Zerobus Ingest ne crée pas de table lorsqu'un client se connecte.
- Un service principal doté des privilèges
USE CATALOG,USE SCHEMA,SELECTetMODIFYpour la table cible. Voir Créer un service principal et accorder des autorisations. - L'URL du Workspace, l'ID du Workspace, la région et l'endpoint Zerobus Ingest. Consultez Obtenir l'URL de votre workspace et l'endpoint Zerobus Ingest.
- Connectivité TLS sortante vers le hostname Zerobus Ingest sur le port
8883.
Confirmez que Zerobus Ingest est disponible dans la région de votre workspace. Voir la disponibilité régionale de Zerobus Ingest. La disponibilité de MQTT peut différer de la disponibilité générale de Zerobus Ingest.
Modèle de connexion
Chaque connexion MQTT cible une table. Définissez la table lorsque le client envoie CONNECT, et utilisez le même nom de table à trois niveaux comme rubrique pour chaque paquet PUBLISH :
catalog.schema.table
Chaque charge utile PUBLISH doit être un objet JSON encodé en UTF-8 qui correspond au schéma de la table cible. Une connexion ne peut pas changer de table. Ouvrez une autre connexion pour publier vers une autre table.
L’endpoint MQTT utilise TLS sur le port 8883 avec le format de hostname suivant :
<workspace-id>.zerobus.<region>.cloud.databricks.com
Il s'agit du même Hostname que l'Endpoint du serveur Zerobus Ingest. Configurez votre client MQTT avec ce Hostname et le port 8883, sans schéma https://.
Authentification
S’authentifier pendant CONNECT avec ces propriétés utilisateur MQTT v5 :
Propriété utilisateur | Valeur |
|---|---|
|
|
| Le nom complet de la table cible, tel que |
Obtenez le jeton OAuth avec le flux d’identifiants client du service principal. Incluez authorization_details pour le catalogue, le schéma et la table, et utilisez la ressource zerobusDirectWriteApi pour votre workspace. L’ exemple Python démontre cet échange.
Les identifiants OAuth s'appliquent uniquement lorsqu'une connexion start. Les connexions ont une durée de vie limitée côté serveur, après quoi Zerobus Ingest déconnecte le client. Lorsque vous vous reconnectez, obtenez un nouveau jeton et envoyez-le dans les propriétés du nouveau CONNECT. L'authentification renforcée MQTT et la réauthentification au sein de la connexion ne sont pas prises en charge.
Durabilité et accusés de réception
Zerobus Ingest prend en charge les niveaux de qualité de service (QoS) MQTT suivants :
- QoS 0 est fourni dans la mesure du possible. Zerobus Ingest tente d'envoyer des enregistrements valides via le pipeline d'ingestion normal, mais n'envoie ni accusé de réception de réussite ni réponse d'échec par enregistrement. Un enregistrement peut être rejeté ou perdu sans notification du client ; par conséquent, une connexion ouverte ne confirme pas la livraison. QoS 0 ne garantit pas une diffusion au moins unique. Utilisez-le uniquement lorsque votre application peut tolérer la perte d'enregistrements.
- QoS 1 renvoie un
PUBACKune fois que Zerobus Ingest a rendu l’enregistrement durable. Considérez l’enregistrement comme accusé de réception uniquement lorsque lePUBACKcorrespond à l’ID du paquet PUBLISH et possède un code de motif de réussite.
Un PUBACK réussi confirme la durabilité, et non une visibilité immédiate dans la table Delta. Zerobus Ingest matérialise les enregistrements durables dans la table de manière asynchrone. Voir latence.
La livraison QoS 1 peut contenir des doublons. Si une connexion se ferme avant que votre client ne reçoive PUBACK, l'enregistrement est peut-être déjà persistant. Le fait de relancer cet enregistrement peut créer un doublon. Les applications qui réessaient des enregistrements non acquittés doivent tolérer les doublons ou les dédupliquer dans le lakehouse.
Un enregistrement QoS 1 n'est pas acquitté si son PUBACK comporte un code de motif d'échec ou si la connexion se ferme avant l'arrivée de PUBACK. Appliquez la règle de nouvelles tentatives de votre application aux enregistrements non acquittés. Un paquet PUBLISH supérieur à la limite de taille des paquets ferme la connexion et est toujours rejeté ; par conséquent, ne réessayez pas.
Exemple : publier un enregistrement
L'exemple suivant utilise Paho MQTT 2.x et requests.
Cet exemple d'introduction désactive le comportement de reconnexion automatique de Paho. Dans une application longue durée, reconnectez-vous avec un jeton OAuth nouvellement récupéré. Suivez les enregistrements sans PUBACK réussi et appliquez votre propre règle de nouvelle tentative, en tenant compte du risque de duplication décrit dans la page Durability and acknowledgments.
L'exemple s'attend à ce qu'une table possède ce schéma :
CREATE TABLE main.default.air_quality (
device_name STRING,
temp INT,
humidity INT
);
Installez les dépendances :
pip install "paho-mqtt>=2,<3" requests
Définissez les variables d’environnement requises. Définissez ZEROBUS_MQTT_HOST sur le hostname MQTT décrit dans la section Modèle de connexion.
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>"
Exécutez le script suivant :
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
L’interface MQTT a les limites de protocole suivantes :
Élément | Aide |
|---|---|
Version du protocole | Réservé à MQTT v5. |
Qualité de service | QoS 0 et QoS 1. QoS 2 n'est pas pris en charge. |
Messages QoS 1 en cours de transmission | Le maximum de réception |
| 16 Kio pour le paquet initial complet |
Taille du paquet après | 64 Kio pour le paquet MQTT complet, y compris l'en-tête fixe, l'en-tête variable, les propriétés, la rubrique, l'ID de paquet et la charge utile. |
Identifiant client | Obligatoire, non vide et ne dépassant pas 128 octets. Utilisez un identifiant client unique pour chaque connexion active. |
État de session | Le serveur ne conserve pas et ne reprend pas l'état de la session et définit l'intervalle d'expiration de la session à |
Intervalle Keep-alive | Entre 5 et 300 secondes. Si la valeur |
Front-end Link | Non pris en charge. Connectez-vous plutôt via l’endpoint public. |
| Seuls le sujet, la QoS et la charge utile affectent l'importation. Le serveur ne conserve ni n'applique l'indicateur de format de charge utile (Payload Format Indicator), le type de contenu (Content Type), les données de corrélation (Correlation Data), l'intervalle d'expiration des messages (Message Expiry Interval), le sujet de réponse (Response Topic) ou les propriétés utilisateur (User Properties). L'identifiant d'abonnement n'est pas pris en charge. |
Fonctionnalités de courtier | Les abonnements, les messages conservés, les alias de rubriques, les messages Last Will et l'authentification renforcée ne sont pas pris en charge. |
Maintenez la connexion ouverte lors de la publication de plusieurs enregistrements dans la même table.
Dépannage
Utilisez les vérifications suivantes pour diagnostiquer les défaillances courantes :
- Échecs de connexion : vérifiez le hostname spécifique au cloud, le port
8883, la résolution DNS, les règles de pare-feu sortantes et la vérification du certificat TLS. Confirmez également que MQTT est activé pour le workspace. - Rejeté
CONNECT: obtenez un nouveau jeton OAuth limité à la table. Vérifiez à la fois les noms de propriété utilisateur, le nom de la table à trois niveaux, la table cible et les attributions du service principal. - Élément
PUBACKerroné ou manquant : ne signalez pas l'enregistrement comme durable. Confirmez que le sujet correspond au nom de la table et que la charge utile est un objet JSON en UTF-8 conforme au schéma. Un élémentPUBACKmanquant ne prouve pas que l'enregistrement n'a pas pu devenir durable. - Déconnexions : vérifiez la présence d'un sujet non concordant, d'un paquet supérieur à 64 KiB, d'une fonctionnalité MQTT non prise en charge, d'un trafic de conservation de connexion manqué ou de la durée de vie limitée de la connexion côté serveur. Obtenez un nouveau jeton avant de vous reconnecter, et appliquez la politique de nouvelles tentatives et de déduplication de votre application aux enregistrements non reconnus.