メインコンテンツまでスキップ

Zerobus Ingest で MQTT を使用する

備考

ベータ版

この機能はベータ版です。使用するには、ワークスペース管理者は、 プレビュー ページから Zerobus Ingest MQTT Endpoint を有効にする必要があります。Databricksのプレビューを管理するを参照してください。

Zerobus Ingest は、すでに MQTT メッセージをパブリッシュしているアプリケーション向けに、MQTT v5 Endpoint を提供します。各メッセージには1つの JSON オブジェクトが含まれており、既存の Unity Catalog の Delta テーブルに直接書き込まれます。このインターフェイスは、MQTT v5 クライアントを使用し、ブローカーのサブスクリプション機能を必要としない、モノのインターネット(IoT)、エッジ、およびテレメトリ アプリケーションに適しています。

前提条件​

接続する前に、次のリソースとアクセス権があることを確認してください:

  • ワークスペースでMQTTが有効になっています。
  • 既存の Unity Catalog Delta テーブル。Zerobus Ingest は、クライアントが接続したときにテーブルを作成しません。
  • ターゲットテーブルに対する USE CATALOG、USE SCHEMA、SELECT、および MODIFY の権限を持つ Service Principal。Service Principalの作成と権限の付与を参照してください。
  • The ワークスペース URL, workspace ID, region, and Zerobus Ingest Endpoint.See Get your ワークスペース URL and Zerobus Ingest Endpoint.
  • ポート 8883 での Zerobus Ingest ホスト名へのアウトバウンド TLS 接続。

ワークスペースのリージョンで Zerobus Ingest が利用可能であることを確認します。Zerobus Ingest のリージョン別の可用性を参照してください。MQTT の可用性は、通常の Zerobus Ingest の可用性とは異なる場合があります。

接続モデル​

各 MQTT 接続のターゲットは 1 つのテーブルです。クライアントが CONNECT を送信したときにテーブルを設定し、すべての PUBLISH パケットのトピックとして同じ 3 レベルのテーブル名を使用します。

Text
catalog.schema.table

各 PUBLISH ペイロードは、ターゲット テーブルのスキーマに一致する1つの UTF-8 でエンコードされた JSON オブジェクトである必要があります。接続でテーブルを切り替えることはできません。別のテーブルに公開するには、別の接続を開きます。

MQTT Endpoint は、ポート 8883 で TLS を使用し、次のホスト名形式を使用します。

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

これは Zerobus Ingest サーバーの Endpoint と同じホスト名です。このホスト名とポート 8883 で、https:// スキームを使用せずに MQTT クライアントを設定します。

認証​

CONNECT の実行中に次の MQTT v5 ユーザープロパティを使用して認証を行います。

ユーザープロパティ

Value

Authorization

Bearer <token>(ここで、<token> はターゲット テーブルをスコープとする Databricks OAuth トークンです)。

x-databricks-zerobus-table-name

main.default.air_quality などのターゲットテーブルのフルネーム。

ユーザープロパティ

Value

Authorization

Bearer <token>(ここで、<token> はターゲット テーブルをスコープとする Databricks OAuth トークンです)。

x-databricks-zerobus-table-name

main.default.air_quality などのターゲットテーブルのフルネーム。

Service Principal のクライアント資格情報フローを使用して OAuth トークンを取得します。カタログ、スキーマ、およびテーブルに authorization_details を含め、ワークスペースに zerobusDirectWriteApi リソースを使用します。Python の例は、この交換を示しています。

OAuth 認証情報は、接続が起動するときにのみ適用されます。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.

耐久性と確認応答​

Zerobus Ingest は、次の MQTT サービス品質(QoS)レベルをサポートしています。

  • QoS 0 はベストエフォートです。Zerobus Ingest は、有効なレコードを通常のパイプラインを介して送信しようと試みますが、成功の通知もレコードごとの失敗の応答も送信しません。クライアントへの通知なしにレコードが拒棄されたり失われたりする可能性があるため、接続が開いているからといって配信が確認されるわけではありません。QoS 0 は最低 1 回の配信を提供しません。アプリケーションがレコードの損失を許容できる場合にのみ使用してください。
  • Zerobus Ingest がレコードを永続化すると、 QoS 1 は PUBACK を返します。PUBACK が PUBLISH パケット ID と一致し、成功の理由コードを持っている場合にのみ、レコードを承認済みとして扱います。

PUBACK が成功しても、Delta テーブルでの即座の可視性ではなく、耐久性が確認されます。Zerobus Ingest は、耐久性のあるレコードをテーブルに非同期でマテリアライズします。レイテンシー を参照してください。

QoS 1 の配信には重複が含まれる場合があります。クライアントが PUBACK を受信する前に接続が切断された場合、レコードはすでに持続性を持つようになっている可能性があります。そのレコードを再試行すると、重複が発生する可能性があります。未確認のレコードを再試行するアプリケーションは、重複を許容するか、レイクハウスで重複排除を行う必要があります。

QoS 1 レコードは、PUBACKに失敗の理由コードがある場合、またはPUBACKが届く前に接続が切断された場合、未確認となります。未確認レコードにアプリケーションの再試行ポリシーを適用します。パケットサイズ制限を超えるPUBLISHパケットは接続を切断し、常に拒否されるため、再試行しないでください。

例:レコードをパブリッシュする​

次の例では、Paho MQTT 2.x および requests を使用します。

注記

この入門用の例では、Paho の自動再接続動作を無効にします。長時間実行されるアプリケーションでは、新しく取得した OAuth トークンで再接続します。正常な PUBACK がないレコードを追跡し、「耐久性と確認応答」に記載されている重複のリスクを考慮したうえで、独自の再試行ポリシーを適用します。

この例では、次のスキーマを持つテーブルが想定されています。

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

依存関係をインストールします:

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

必要な環境変数を設定します。ZEROBUS_MQTT_HOST に、接続モデル で説明されている MQTT ホスト名を設定します。

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>"

次のスクリプトをラン:

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()

制限事項​

MQTT インターフェイスには次のプロトコル制限があります。

項目

サポート

プロトコル バージョン

MQTT v5 のみ。

サービス品質

QoS 0 および QoS 1。QoS 2 はサポートされていません。

転送中の QoS 1 メッセージ

CONNACK 受信最大数は接続あたり4,096メッセージです。クライアントは通知された制限を遵守する必要があります。

CONNECT パケットサイズ

プロパティを含む、完全な初期 CONNECT パケットに対して 16 KiB。

送信後のパケットサイズ CONNECT

固定ヘッダー、可変ヘッダー、プロパティ、トピック、パケット ID、ペイロードを含む、完全な MQTT パケットの 64 KiB。CONNACK では、この制限を最大パケットサイズとして通知します。より大きなパケットでは、理由コード「Packet too large」で接続が閉じられます。

クライアントID

必須であり、空ではなく、128バイト以下である必要があります。アクティブな接続ごとに一意のクライアント ID を使用してください。

セッションの状態

サーバーはセッション状態を保持または再開せず、セッション有効期限間隔(Session Expiry Interval)を 0 に設定します。

キープアライブの間隔

5~300 秒。CONNECT の値が 0 である場合、サーバーは 60 秒を使用します。そうでない場合、サーバーはこの範囲に値を制限します。CONNACK で返されたサーバーキープアライブ(Server Keep Alive)の値を尊重してください。有効な間隔の 1.5 倍の期間パケットがまったく受信されない場合、サーバーは接続を切断します。パケットは最後のバイトが到着したときにのみ受信されたとカウントされるため、低速なLinkでは、最大のレコードを送信するのに十分な長さのキープアライブ間隔を選択してください。

フロントエンドPrivateLink

サポートされていません。代わりにパブリック Endpoint 経由で接続してください。

PUBLISH properties

トピック、QoS、およびペイロードのみが取り込みに影響します。サーバーは、ペイロード形式インディケータ、コンテンツタイプ、相関データ、メッセージ有効期限間隔、応答トピック、またはユーザープロパティを保持または適用しません。サブスクリプション識別子はサポートされていません。

ブローカーの機能

サブスクリプション、保持メッセージ、トピックエイリアス、Last Willメッセージ、および拡張認証はサポートされていません。

項目

サポート

プロトコル バージョン

MQTT v5 のみ。

サービス品質

QoS 0 および QoS 1。QoS 2 はサポートされていません。

転送中の QoS 1 メッセージ

CONNACK 受信最大数は接続あたり4,096メッセージです。クライアントは通知された制限を遵守する必要があります。

CONNECT パケットサイズ

プロパティを含む、完全な初期 CONNECT パケットに対して 16 KiB。

送信後のパケットサイズ CONNECT

固定ヘッダー、可変ヘッダー、プロパティ、トピック、パケット ID、ペイロードを含む、完全な MQTT パケットの 64 KiB。CONNACK では、この制限を最大パケットサイズとして通知します。より大きなパケットでは、理由コード「Packet too large」で接続が閉じられます。

クライアントID

必須であり、空ではなく、128バイト以下である必要があります。アクティブな接続ごとに一意のクライアント ID を使用してください。

セッションの状態

サーバーはセッション状態を保持または再開せず、セッション有効期限間隔(Session Expiry Interval)を 0 に設定します。

キープアライブの間隔

5~300 秒。CONNECT の値が 0 である場合、サーバーは 60 秒を使用します。そうでない場合、サーバーはこの範囲に値を制限します。CONNACK で返されたサーバーキープアライブ(Server Keep Alive)の値を尊重してください。有効な間隔の 1.5 倍の期間パケットがまったく受信されない場合、サーバーは接続を切断します。パケットは最後のバイトが到着したときにのみ受信されたとカウントされるため、低速なLinkでは、最大のレコードを送信するのに十分な長さのキープアライブ間隔を選択してください。

フロントエンドPrivateLink

サポートされていません。代わりにパブリック Endpoint 経由で接続してください。

PUBLISH properties

トピック、QoS、およびペイロードのみが取り込みに影響します。サーバーは、ペイロード形式インディケータ、コンテンツタイプ、相関データ、メッセージ有効期限間隔、応答トピック、またはユーザープロパティを保持または適用しません。サブスクリプション識別子はサポートされていません。

ブローカーの機能

サブスクリプション、保持メッセージ、トピックエイリアス、Last Willメッセージ、および拡張認証はサポートされていません。

同じテーブルに複数のレコードを公開する際、接続を開いたままにします。

トラブルシューティング​

一般的な障害を診断するには、次のチェックを使用してください。

  • 接続エラー: クラウド特有のホスト名、ポート 8883、DNS 解決、アウトバウンド ファイアウォール規則、および TLS 証明書の検証を確認してください。また、ワークスペースで MQTT が有効になっていることを確認してください。
  • 拒否された CONNECT:テーブルスコープの OAuth トークンを新しく取得してください。両方のユーザープロパティ名、3 レベルのテーブル名、ターゲットテーブル、および Service Principal の付与を確認します。
  • Unsuccessful or missing PUBACK: Don't report the record as durable.トピックがテーブル名と一致していること、およびペイロードがスキーマに一致する 1 つの UTF-8 JSON オブジェクトであることを確認します。PUBACK が欠落していても、レコードの永続化に失敗したとは証明されません。
  • 切断:トピックの不一致、64 KiB を超えるパケット、サポートされていない MQTT 機能、キープアライブトラフィックの欠落、または制限されたサーバー側の接続ライフタイムを確認してください。再接続する前に新しいトークンを取得し、未確認のレコードに対してアプリケーションの再試行および重複排除ポリシーを適用します。