
TDengine 零代碼接入 SparkplugB基于 taosExplorer 的 IIoT 數據同步任務配置指南【免費下載鏈接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios項目地址: https://gitcode.com/GitHub_Trending/tde/TDengine本文基于 TDengine 開源倉庫 docs/en/08-data-ingest-and-delivery/01-no-code-ingestion/18-sparkplugb.md 編寫。SparkplugB 是專為工業物聯網IIoT設計、構建于 MQTT 之上的開放消息規范。借助 TDengine 的零代碼數據接入平臺你可以在 taosExplorer 圖形界面中直接創建數據源任務讓 taosX 連接器從 MQTT Broker 訂閱 SparkplugB 消息并實時寫入 TDengine 集群——無需編寫任何代碼。讀完本文你將掌握從新增數據源、配置連接認證、訂閱過濾、Payload 轉換解析/拆分/過濾/表映射到高級選項與異常處理策略的完整任務創建流程。背景為什么需要 SparkplugB 數據接入SparkplugB 是一種開放的消息規范專為工業物聯網IIoT應用設計底層基于 MQTT 協議。它定義了 IIoT 場景下設備與 MQTT Broker 之間的標準消息格式如 NBIRTH、NDATA、DDATA 等被廣泛用于工廠、產線、設備監控等工業數據采集場景。TDengine 通過內置的 SparkplugB 連接器可以從 MQTT Broker 訂閱 SparkplugB 消息并將數據實時寫入 TDengine實現工業數據的實時入庫。整個流程在 taosExplorer 的數據寫入Data In頁面通過零代碼配置完成無需額外部署 ETL 工具。從 零代碼接入平臺總覽 可知taosExplorer 是 TDengine 的可視化數據管理工具支持在瀏覽器中通過簡單配置向 TDengine 提交任務實現多種數據源到 TDengine 的零代碼導入并在導入過程中自動完成數據的抽取、過濾與轉換。SparkplugB 正是該平臺支持的眾多數據源之一。說明SparkplugB 連接器在消息斷點續傳方面有明確限制——從 任務斷點續傳 一節可知與 MQTT、Kafka 等數據源不同SparkplugB 當前不支持消息持久化與恢復任務重啟后無法從上次斷點續傳。創建數據寫入任務新增數據源登錄 taosExplorer 后進入數據寫入頁面點擊 新增數據源Add Data Source按鈕進入任務創建頁面。配置基本信息在任務創建頁面中配置以下基本信息名稱輸入任務名稱例如test_spb類型在下拉列表中選擇SparkplugB代理可選選擇一個已創建的 Agent 代理或點擊右側 創建新的代理按鈕新建目標數據庫在下拉列表中選擇一個目標數據庫或點擊右側 創建數據庫按鈕新建。配置連接與認證信息在連接配置區域需要填寫 MQTT Broker 的接入參數配置項說明BrokersMQTT Broker 地址例如localhost:1883。可以填寫多個用逗號,分隔用于連接多個 BrokerMQTT 協議使用的 MQTT 協議版本默認5.0可選 3.1、3.1.1、5.0客戶端 ID連接到每個 Broker 時使用的客戶端標識符Keep Alive保持活動間隔。如果 Broker 在該間隔內沒有收到來自客戶端的任何消息會假定客戶端已斷開并關閉連接。該間隔是客戶端與 Broker 之間協商的、用于檢測客戶端活躍狀態的時長用戶 / 密碼MQTT Broker 認證所需的用戶名與密碼若 Broker 開啟了認證則必填實操提示同一 MQTT Broker 下如果創建多個同步任務各任務的客戶端 ID 必須互不相同否則會造成沖突導致任務無法正常運行。TLS 校驗模式TLS 校驗TLS Verification支持三種模式不開啟Disabled不進行 TLS 證書認證。連接 MQTT 時會先嘗試 TCP 連接若失敗則改為無證書認證模式的 TLS 連接。單向認證One-way開啟 TLS 連接并驗證服務端證書此時需要上傳CA 證書。雙向認證Mutual開啟 TLS 連接并與服務端進行雙向認證此時需要上傳CA 證書、客戶端證書以及客戶端私鑰。配置完成后點擊檢查連通性Check Connectivity按鈕驗證數據源是否可用若檢查失敗請根據頁面返回的具體錯誤提示修改配置。訂閱配置訂閱配置決定連接器從哪些主題、哪些設備、哪些消息類型中消費數據。Group ID填寫 SparkplugB 規范定義的 group id。通常一個 group id 代表一個集團/公司/工廠/流水線等概念。節點/設備列表Node/Device List填寫需要訂閱的節點和設備列表以逗號分隔。其中節點直接填寫 ID 即可設備需要按照節點ID/設備ID的格式填寫。消息類型Message Types填寫需要訂閱的 SparkplugB 消息類型以逗號分隔。支持的類型包括NBIRTH、NDEATH、NDATA、NCMD、DBIRTH、DDEATH、DDATA、DCMD、STATE。其中NBIRTH、NDEATH、NDATA、NCMD類型的消息只會匹配節點/設備列表中的節點而DBIRTH、DDEATH、DDATA、DCMD只會匹配節點/設備列表中的設備。下發 REBIRTH 命令Send REBIRTH Command開啟后taosX 會自動下發 NCMD 中的Node Control/Rebirth命令從而獲取節點和設備的所有 metric 信息包括 metric name 與 metric alias 的對應關系。如果節點/設備在上報數據時不使用 alias 別名機制可以不開啟此選項。配置 Payload 轉換Payload 轉換是 SparkplugB 數據接入的核心環節包含解析、字段拆分、數據過濾、表映射四步是 taosX 內置 ETL 能力的具體體現參見 數據抽取、過濾與轉換。解析 PayloadPayload 解析區域提供三種獲取示例數據的方式點擊從服務器檢索從已配置的 MQTT Broker 獲取示例數據點擊文件上傳上傳文件獲取示例數據在消息體中手動填寫 MQTT 消息體的示例數據。由于 SparkplugB 消息使用Protocol Buffersprotobuf編碼從服務器檢索到的數據會先被解碼為 JSON 格式。JSON 數據支持 JSONObject 或 JSONArray 兩種形態可以用于解析 SparkplugB 中的 metadata、properties 等 JSON 格式字段。點擊放大鏡圖標可預覽解析結果從列中提取或拆分字段解析后的數據可能仍不滿足目標表的要求此時可以在從列中提取或拆分Extract or Split區域填寫提取/拆分規則。典型的場景是將datatype_str字段的值轉換為 TDengine 數據類型。選擇映射mapping提取器在rule輸入框中填寫如下 JSON在name中填寫td_datatype{ Int8: TINYINT, UInt8: TINYINT UNSIGNED, Int16: SMALLINT, UInt16: SMALLINT UNSIGNED, Int32: INT, UInt32: INT UNSIGNED, Int64: BIGINT, UInt64: BIGINT UNSIGNED, Float: FLOAT, DOUBLE: DOUBLE, Boolean: BOOL, String: VARCHAR(128), DateTime: TIMESTAMP }該規則會將datatype_str列的值例如字符串Int8轉換為對應的 TDengine 類型例如TINYINT并生成新的列td_datatype。你可以點擊新增添加更多提取規則點擊刪除移除當前規則點擊放大鏡圖標預覽提取/拆分結果。數據過濾在過濾Filter區域填寫過濾表達式只有滿足條件的數據行才會被寫入 TDengine。例如填寫datatype_str ! Int8則只有datatype_str值不為Int8的數據才會被寫入。過濾表達式的結果必須為布爾類型支持基于字段類型的判斷函數與比較運算符、、、、、!多個條件可通過邏輯運算符、||、!組合。例如location.starts_with(beijing) voltage 200表示只同步北京地區電壓大于 200 的智能電表數據。相關過濾語法細節可參考 零代碼接入平臺的過濾章節。點擊刪除可移除當前過濾規則點擊放大鏡圖標可預覽過濾結果。表映射表映射將解析、提取、拆分后的源字段映射到 TDengine 目標表。目標超級表在下拉列表中選擇一個目標超級表或點擊右側創建超級表按鈕新建。創建模板當超級表需要根據消息動態生成時選擇創建模板。此時超級表名稱、列名、列類型等均可以使用模板變量。接收到數據后程序會自動計算模板變量并生成對應的超級表模板當數據庫中該超級表不存在時使用模板創建超級表對于已創建的超級表如果缺少通過模板變量計算得到的列也會自動創建對應列。映射填寫目標超級表中的子表名稱例如t_{id}根據需求填寫映射規則其中 mapping 支持設置缺省值默認值。點擊預覽可查看映射結果確認子表名稱、列與標簽的映射是否符合預期。配置高級選項高級選項Advanced Options區域默認折疊點擊展開。MQTT 與 SparkplugB 數據源常用的選項如下字段名可能因連接器而異參見 高級選項詳解選項說明Message Queue Size消息隊列大小接收緩沖區大小。隊列滿且未開啟緩存實時數據時新到達的數據會被丟棄設為0表示禁用緩沖Maximum In-Process Batches最大進行中批次可并發處理的批次數上限。達到上限后連接器停止從接收隊列取消息消息會在隊列中累積最小值為1Batch Size批量大小每次送入處理管道的消息條數。與批量延遲配合使用即使延遲未到批量已滿也會立即發送最小值為1Batch Delay批量延遲每批次的超時時間毫秒從該批次第一條消息到達開始計時。超時后即使未達到 Batch Size 也會發送該批次最小值為1Write Concurrency寫入并發并發寫入 TDengine 的任務數Cache Realtime Data緩存實時數據開啟后消費到的數據先寫入本地文件由后臺任務轉發下游當下游處理跟不上時起到流量整形作用積壓消費完畢后緩存文件會被清除。默認關閉。詳見 Store and ForwardCache Storage Directory緩存存儲目錄覆蓋緩存文件的存儲目錄僅在開啟緩存實時數據時生效否則默認使用 taosX 啟動時配置的數據目錄Save Raw Data保存原始數據開啟后可進一步配置最大保留天數與原始數據存儲目錄此外高級選項中還包含健康監控設置Health Check Duration、Busy State Threshold、Max Write Queue Length、Write Error Threshold用于任務列表頁的健康狀態展示具體說明參見 Health Status。配置異常處理策略異常處理策略Exception Handling Strategy區域默認折疊點擊展開。taosX 為各類寫入異常提供了統一的分流策略參見 異常處理策略詳解歸檔Archive將無效數據寫入歸檔文件默認位于${data_dir}/tasks/id/datetime下不寫入目標數據庫丟棄Discard忽略無效數據報錯Error報告錯誤緩存Cache目標連接失敗或資源不足時將數據寫入緩存文件待目標恢復后再行入庫。可針對以下場景分別配置策略目標連接超時歸檔 / 丟棄 / 報錯 / 緩存目標數據庫不存在歸檔 / 丟棄 / 報錯表不存在歸檔 / 丟棄 / 報錯 / 自動建表并重試主時間戳超出范圍now - keep1至now 100y歸檔 / 丟棄 / 報錯主時間戳為空歸檔 / 丟棄 / 報錯 / 使用當前時間復合主鍵為空歸檔 / 丟棄 / 報錯表名超過 192 字符歸檔 / 丟棄 / 報錯 / 截斷 / 截斷并歸檔表名含非法字符如.歸檔 / 丟棄 / 報錯 / 用配置的字符串替換非法字符表名模板變量為空丟棄 / 變量留空 / 用配置的字符串替換列不存在歸檔 / 丟棄 / 報錯 / 自動補列并重試列名超過 64 字符歸檔 / 丟棄 / 報錯列值超出定義長度歸檔 / 丟棄 / 報錯 / 截斷 / 截斷并歸檔也可通過自動擴列修改表結構后重試其他數據錯誤歸檔 / 丟棄 / 報錯。附加設置項連接超時Connection Timeout目標連接超時時間秒取值范圍1~600臨時存儲位置相對${data_dir}/tasks/id/的路徑歸檔保留天數Archive Retention Days非負整數0表示不限歸檔可用空間Archive Available Space取值范圍0~655350表示不限歸檔位置Archive Location相對${data_dir}/tasks/id/的路徑歸檔寫入失敗策略刪除舊文件 / 丟棄數據 / 報錯并停止任務。提交任務完成上述所有配置后點擊提交Submit按鈕即完成 SparkplugB 到 TDengine 的數據同步任務創建自動回到數據源列表Data Source List頁面。提交成功后可在任務列表頁查看任務執行情況包括寫入記錄數、流量等運行指標任務狀態會切換為 Running。你也可以在任務列表頁對任務進行啟動、停止、查看、刪除、復制等管理操作并查看每個任務的健康狀態Ready、Idle、Active、Pending、Busy、Bounce、SourceError、SinkError、Fatal 等詳見 任務管理。小結SparkplugB 數據接入任務的核心鏈路可概括為taosX 連接器訂閱 MQTT Broker → 解碼 protobuf 為 JSON → 解析/拆分/過濾 → 映射到超級表與子表 → 實時寫入 TDengine。整個過程完全通過 taosExplorer 的零代碼界面完成涵蓋連接認證含 TLS 單向/雙向認證、訂閱配置Group ID、節點/設備、消息類型、REBIRTH、Payload 轉換四種 ETL 步驟以及高級選項與異常處理兜底策略。配置時需特別注意SparkplugB 當前不支持消息持久化與斷點續傳對于需要高可靠連續采集的工業場景建議結合網絡穩定性保障與異常歸檔策略共同使用。相關參考文檔SparkplugB 接入指南英文原檔SparkplugB 接入指南中文原檔零代碼數據接入平臺總覽taosX Agent 存儲轉發Store and Forward【免費下載鏈接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios項目地址: https://gitcode.com/GitHub_Trending/tde/TDengine創作聲明:本文部分內容由AI輔助生成(AIGC),僅供參考