第 11 章

Connectors 與 Sinks

Connector 是 EMQX 連到外部系統的底層通道,Sink 則把規則輸出的資料外送到 Webhook、HTTP、Kafka、PostgreSQL/MySQL 等目標;這章以 Webhook/HTTP 為主實作。

為什麼要學這個

規則引擎把資料轉好之後,「要怎麼送出去」是另一個難題。家裡的感測器事件若想送進手機通知、簡訊閘道、雲端 API 或關聯式資料庫,你得先建立一條 EMQX 通往外部系統的連線通道。

EMQX 5.x 把這過程拆成 Connector(連線)與 Sink(輸出):Connector 只建一次可被多個 Sink/Source 共用;Sink 把規則輸出的資料實際送到外部目標。Webhook/HTTP 是最容易先在智慧家庭裡試的起點。

核心概念

Connector(連線器)是 EMQX 對外部資料系統的底層連線通道,只負責「怎麼連」,不負責「處理哪筆資料」。建立 Sink/Source 時可以選用已有的 Connector,或在建立過程中順手生出一個。

  • 分離設定:連線資訊(伺服器位址、帳密)與資料設定(規則、topic 對應、payload 模板)分開,改連線不給影響規則。
  • 重用:要連多個 Sink/Source 的外部系統,只需一個 Connector。
  • 狀態可觀察:Dashboard 可看 Connector 連線狀態(Connecting、Connected、Disconnected、Inconsistent),有助排錯。

Sink(輸出端)決定「把規則輸出的資料送往哪個外部目標及格式」,例如 Webhook、HTTP Server、Kafka 或資料庫。它放在規則的「Action Outputs」裡,挑選一個 Connector並給出對應參數。

名詞對照

英文中文一句話
Connector連線器對外系統的底層連線通道
Sink輸出端把 EMQX 資料寫出去的外部目標
Source來源從外部系統把資料收進 EMQX 的入口
Webhook網路鉤子以 HTTP 將訊息送出的外部目標
HTTP Server IntegrationHTTP 資料離場透過 Connector+規則把資料送進 HTTP 服務

動手做

以官方教學的 HTTP Server Sink 為例,把 t/# 的 MQTT 訊息送到本機 Flask HTTP 服務。

  1. 建立 HTTP Server Connector

    Dashboard 左側選 Integration → Connectors,按 Create,選 HTTP Server 連線器,取名 my_httpserver,URL 填 http://localhost:5000;可先按 Test Connectivity 測試。

  2. 建立規則

    選 Integration → Rules → Create,SQL 輸入 SELECT * FROM "t/#"。

  3. 新增 HTTP Sink 動作

    按 + Add Action,Type 選 HTTP Server;Connector 下拉選 my_httpserver,Method 選 POST,填名稱說明。

  4. 建立規則

    回到 Create Rule 頁確認 Sink 在 Action Outputs,按 Create 完成。

  5. 用訊息觸發

    用 MQTTX 向 t/1 發布例如 {"msg":"hello HTTP Server"},再到 Flask server 端看收到的 POST request。

Connector 的目標有哪些

EMQX 5.8.9 提供多種外部目標做 Sink,但「開源版」能用的是有限的。官方在 Connector 文件標註:EMQX 開源版(Community)只支援 HTTP 與 MQTT 兩種 Connector;Kafka、PostgreSQL、MySQL、AWS、Azure 等目標屬 企業版(Enterprise)功能。

目標適合版本
Webhook不需規則處理、直接把 MQTT 資料送給 HTTP 服務開源版可用
HTTP Server先用規則過濾/加工再送 HTTP 外送開源版可用
MQTT Broker跨 broker 橋接(第 12 章)開源版可用
Kafka把事件流送進 Kafka 訊息佇列企業版
PostgreSQL/MySQL把裝置資料直接寫進關聯式資料庫企業版

Webhook 與 HTTP Server 的差別,主要是要不要先用規則(SQL)處理資料:Webhook 較簡單、不需 SQL;HTTP Server 則能把規則輸出再動態組出 request header、body 甚至 URL。

以 Webhook/HTTP 外送資料

最常先在智慧家庭情境試的是「把事件通知送到 HTTP 服務」,有兩條路:

  • Webhook(較簡):不需規則處理,適合「只要把收到的訊息送給外部 HTTP 端點」的簡單整合,目標清楚。
  • HTTP Server Integration(較進階):先把資料經規則引擎提取與轉換,再送進 HTTP 服務。做法是先建立 HTTP Server Connector(URL、方法),再加規則與 HTTP Sink 動作。

HTTP Sink 可選 POST/PUT 等方法,並可設定 request 內容與帳密(可使用 ${var} 佔位符)。EMQX 會顯示每個 Sink 的連線狀態與請求次數,方便確認資料是否送出。

只送一則狀態通知給某個 API,用 Webhook 最輕;需要過濾、加欄位或加工後再外送,就選 HTTP Sink 搭配規則。

故障排除

  • 外部 HTTP 服務沒收到資料:確認 Connector 狀態為 Connected、URL 與 Method 正確、Sink 已加進規則、規則已啟用。
  • 選不到 Kafka/PostgreSQL/MySQL 連接器:開源版只支援 HTTP 與 MQTT 連接器;這些目標屬企業版。
  • 改了 Connector 之後資料中斷:正在使用的 Connector 被更新會觸發 Sink/Source reload,可能暫時中斷,建議離峰更新。
  • 刪不掉 Connector:使用中的 Connector 不能刪;先刪依附它的 Sink/Source。Dashboard 會列出關聯的 Sink。

常見問題

Connector 與 Sink 差在哪

Connector 是底層連線通道(連到哪個外部系統);Sink 是「把規則輸出送到那個目標」的執行設定,放在規則的 Action 裡。一個 Connector 可被多個 Sink/Source 使用。

Webhook 與 HTTP Server 怎麼選

資料若不需規則處理,用 Webhook 最簡單;若要先去噪、再轉址,就用 HTTP Server Sink 搭配規則。官方建議不需規則處理時直接用 Webhook。

開源版能連哪些外部目標

開源版(Community edition)只支援 HTTP 與 MQTT 兩種 Connector;Kafka、PostgreSQL、MySQL 等目標屬企業版功能。

想送資料到資料庫一定要企業版嗎

官方原生的資料庫 Sink(PostgreSQL、MySQL 等)在企業版才可用。開源版可改用 HTTP Sink 搭配自己的資料庫 API,或先用 MQTT 橋接到另一套資料處理服務。

官方來源