第 12 章

Sources 與資料匯入

Source 是把「外部 MQTT broker 的訊息接回你自己的 EMQX」的入口,配合 republish 動作就能把外部訊息寫回 EMQX,實現跨 broker 橋接。

為什麼要學這個

不是所有裝置都連到你的 EMQX。也許你有一台 Zigbee 整合器或雲端平台,已經把資料發布到 別的 MQTT broker。想讓這些外部消息也能進到你的 EMQX、再進規則或 HA,就需要一條「進來的通道」,這就是 Source(來源)。

11 章講的 Sink 是「往外送」,這章講相反方向的資料匯入:透過 MQTT Source 從遠端 broker 收資料,再用 republish(或規則)把外部訊息「回寫」成 EMQX 本地可用的 topic。

核心概念

MQTT broker 之間的資料橋接(許多人叫「MQTT 橋接」,EMQX 5.x 的文件叫 MQTT Broker Data Integration)讓 EMQX 以「client」身分連到另一個 MQTT 服務,做雙向訊息交換:

  • 出站(Sink):把本地 EMQX 的訊息發布到遠端 broker 的指定 topic。
  • 入站(Source):訂閱遠端 broker 的 topic,把收到的訊息送回本地 EMQX。

兩者共用同一個 Connector(第 11 章)。同一個連線上可以配置多條橋接規則,各自有不同的 topic 對應與轉換。Source 收進來的訊息要能讓本地 client 讀到,通常要加上一個 republish 動作寫回本地。

名詞對照

英文中文一句話
Source來源外部資料進入 EMQX 的入口(入站)
Sink輸出端EMQX 資料送往外部(出站)
MQTT Broker BridgeMQTT 橋接在兩個 MQTT broker 之間傳訊息
Ingress入站把遠端 broker 的訊息接進 EMQX
Egress出站把 EMQX 訊息發布到遠端 broker

動手做

以官方教學為例:把遠端 broker(broker.emqx.io)收到的 f/# 訊息橋接成本地 sub/#。

  1. 建立 MQTT Broker Connector

    Dashboard 選 Integration → Connectors → Create,選 MQTT Broker,填名稱 my_mqtt_bridge,伺服器設 broker.emqx.io:1883(若需要驗證則填帳密);按 Create。

  2. 建立規則並加入 MQTT Source

    選 Integration → Rules → Create;在「Data Inputs」分頁刪掉預設的 Message 輸入,再按 Add Input,選 MQTT Broker,選 Connector,設定訂閱 Topic(例如 $share/1/f/#)與 QoS。

  3. 查看規則 SQL

    加入 Source 後,規則 SQL 自動變成 SELECT * FROM "$bridges/mqtt:my_source",表示資料來源是這個 MQTT 橋接。

  4. 新增 Republish 動作

    切到「Action Outputs」分頁按 + Add Action,選 Republish;Topic 填 sub/${topic},QoS 選 ${qos},Payload 填 ${payload}。

  5. 建立並驗證

    按 Create 完成規則。訂閱本地 sub/#,在遠端 broker 的 f/1 發布訊息,本地應接到 sub/f/1。

MQTT Source 與 $bridges 語法

建立 MQTT Source 時留意兩件事:

  • 訂閱 topic 支援 + 與 # 萬用字元,可一次收一系列 remote topic。
  • 叢集或連線池時要用共享訂閱:當 EMQX 是叢集部署、或 Connector 開了連線池,多個 client 同時訂閱同一 topic 會收到重複訊息;官方建議用 $share/群組名/topic(例如 $share/1/f/#)把它分攤掉。

Source 收進來的訊息欄位(供規則 SQL 使用)包括:

欄位說明
topic原始訊息 topic
server來源 broker 位址
payload訊息內容
qos訊息 QoS
retain是否 retained
pub_propsMQTT 5.0 訊息屬性(user property 等)
message_received_at接收時間(毫秒)

規則用 $bridges/mqtt:<name> 當資料來源,其中 <name> 是 Connector/Source 的名稱。

跨 broker 橋接與回寫 EMQX

出站這側很單純:建立 MQTT Sink(egress),把本地 topic 發布到遠端。例如規則 SELECT * FROM "t/#" 對應 pub/${topic} 目標,本地 t/1 會被轉送到遠端 pub/t/1。發送時 QoS 與 retain 可用 ${qos}、${flags.retain} 佔位符沿用原訊息。

入站寫回這條就是 Source 的用途:Source 把遠端訂閱到的訊息引進規則,但不會就地發布到本地——你得再加一個 Republish 動作把它「回寫」成本地 topic。官方範例用 sub/${topic},所以遠端 f/1 進到你 sub/f/1,本地 client 訂 sub/# 即可收到。

想用設定檔宣告時,在 rule_engine 與 connector/bridges 區域建立對應的 ingress(Source)與 egress(Sink)bridge,與規則的 SQL 相呼應。若為叢集或多節點連線,MQTT 的固定 client ID 會彼此衝突,官方因此不支援在此處設定固定 client ID。

故障排除

  • Source 沒收到遠端訊息:確認 Connector 狀態為 Connected、訂閱 Topic 與 QoS 正確、萬用字元無誤、遠端 broker 確有資料發布。
  • 叢集/連線池資料重複:改用共享訂閱 $share/群組/topic,避免多個 bridge client 重複收到同一筆。
  • 本地收不到回寫的訊息:Source 只是把訊息引進;需另外用 Republish 動作才真的寫回本地 topic。確認動作已加並已創建。
  • SQL 欄位拿不到:確認規則 FROM 是正確的 $bridges/mqtt:<名稱>,且放進 的名稱與 Source 名稱一致。

常見問題

Source 與 Sink 的差別

Source 是入站(外部 → EMQX),Sink 是出站(EMQX → 外部)。同一個 MQTT Connector 可同時供 Source 與 Sink 使用,訊息的進出彼此獨立。

MQTT 橋接一定要叢集嗎

不必。但它支援叢集與連線池;當多節點或多 client 並連時,官方建議用共享訂閱避免重複訊息。

能用固定 client ID 連到遠端 broker 嗎

不建議。叢集或連線池時多個節點用同一 client ID 會衝突,斷後重連也不穩;MQTT 橋接不提供固定 client ID,EMQX 會自動產生唯一的 client ID。

入站訊息一定要用 Republish 嗎

要。Source 只把外部訊息帶進規則處理流程,不會自動發布到本地 topic;想讓本地 client(或 HA)能收到,就得加 Republish(或送往 Sink)動作。

官方來源