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