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 Integration | HTTP 資料離場 | 透過 Connector+規則把資料送進 HTTP 服務 |
動手做
以官方教學的 HTTP Server Sink 為例,把 t/# 的 MQTT 訊息送到本機 Flask HTTP 服務。
建立 HTTP Server Connector
Dashboard 左側選 Integration → Connectors,按 Create,選 HTTP Server 連線器,取名
my_httpserver,URL 填http://localhost:5000;可先按 Test Connectivity 測試。建立規則
選 Integration → Rules → Create,SQL 輸入
SELECT * FROM "t/#"。新增 HTTP Sink 動作
按 + Add Action,Type 選 HTTP Server;Connector 下拉選
my_httpserver,Method 選 POST,填名稱說明。建立規則
回到 Create Rule 頁確認 Sink 在 Action Outputs,按 Create 完成。
用訊息觸發
用 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 橋接到另一套資料處理服務。