分类:14.生态 | 篇章:04 第三方工具
免费详情
TDengine 通过 InfluxDB 兼容协议、JDBC、连接器等方式与主流数据生态对接。本文汇总 Telegraf、Kafka Connect、Flink、Spark、Logstash 等工具的集成方式。
集成方式速查
| 工具 | 集成方式 | 用途 |
|---|---|---|
| Telegraf | InfluxDB output | 系统/IoT 采集 |
| collectd | collectd protocol | 服务器监控 |
| StatsD | StatsD protocol | 应用指标 |
| Prometheus | remote_write | 长期存储 |
| Kafka Connect | JDBC Sink | Kafka → TD |
| Apache Flink | JDBC Sink | 流处理结果存储 |
| Apache Spark | JDBC | 大数据分析 |
| Logstash | JDBC output | 日志数据 |
| DBeaver | JDBC | SQL IDE |
| Hive | JDBC | 数据仓库 |
详细解析
1. Telegraf 集成
1# /etc/telegraf/telegraf.conf 2 3# 输入插件(按需) 4[[inputs.cpu]] 5 percpu = true 6 totalcpu = true 7 8[[inputs.mem]] 9 10[[inputs.system]] 11 12 13# 输出到 TDengine(通过 InfluxDB Line 协议) 14[[outputs.http]] 15 url = "http://taosadapter:6041/influxdb/v1/write?db=telegraf" 16 method = "POST" 17 username = "root" 18 password = "taosdata" 19 data_format = "influx" 20 21 22# 启动 23systemctl start telegraf 24
2. collectd 集成
1# /etc/collectd/collectd.conf 2 3LoadPlugin network 4 5<Plugin "network"> 6 Server "taosadapter" "6045" 7</Plugin> 8 9LoadPlugin "cpu" 10LoadPlugin "memory" 11LoadPlugin "disk" 12LoadPlugin "interface" 13 14 15# taosAdapter 配置开启 collectd 16# /etc/taos/taosadapter.toml 17[collectd] 18enable = true 19port = 6045 20db = "collectd" 21user = "root" 22password = "taosdata" 23
3. StatsD 集成
1# StatsD 协议(UDP) 2echo "myapp.requests:1|c|@1.0" | nc -u -w0 taosadapter 6044 3 4 5# taosAdapter 配置 6# /etc/taos/taosadapter.toml 7[statsd] 8enable = true 9port = 6044 10db = "statsd" 11
4. Prometheus remote_write
1# prometheus.yml 2remote_write: 3 - url: "http://taosadapter:6041/prometheus/v1/remote_write/prometheus_data" 4 basic_auth: 5 username: root 6 password: taosdata 7 8 9# 长期存储所有 Prometheus 指标 10# 利用 TDengine 高压缩比 11 12 13# 反向:从 TDengine 读 Prometheus 指标 14remote_read: 15 - url: "http://taosadapter:6041/prometheus/v1/remote_read/prometheus_data" 16
5. Kafka Connect JDBC Sink
1// kafka-connect-jdbc-sink.json 2{ 3 "name": "tdengine-sink", 4 "config": { 5 "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", 6 "tasks.max": "4", 7 "topics": "sensor_data", 8 "connection.url": "jdbc:TAOS-RS://taosadapter:6041/iot", 9 "connection.user": "root", 10 "connection.password": "taosdata", 11 "insert.mode": "insert", 12 "auto.create": "true", 13 "auto.evolve": "true", 14 "pk.mode": "record_value", 15 "pk.fields": "ts" 16 } 17} 18 19 20# 部署到 Kafka Connect 21curl -X POST http://kafka-connect:8083/connectors \ 22 -H "Content-Type: application/json" \ 23 -d @kafka-connect-jdbc-sink.json 24
6. Apache Flink 集成
1// 用 JDBC Sink 写 TDengine 2import org.apache.flink.connector.jdbc.*; 3 4DataStream<MeterReading> stream = ...; // 你的数据流 5 6stream.addSink(JdbcSink.sink( 7 "INSERT INTO meters VALUES (?, ?, ?)", 8 (ps, reading) -> { 9 ps.setTimestamp(1, reading.ts); 10 ps.setFloat(2, reading.current); 11 ps.setInt(3, reading.voltage); 12 }, 13 JdbcExecutionOptions.builder() 14 .withBatchSize(1000) 15 .withBatchIntervalMs(200) 16 .build(), 17 new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() 18 .withUrl("jdbc:TAOS-WS://taosadapter:6041/iot") 19 .withDriverName("com.taosdata.jdbc.ws.WebSocketDriver") 20 .withUsername("root") 21 .withPassword("taosdata") 22 .build() 23)); 24
7. Apache Spark 集成
1// 读 TDengine 2val df = spark.read 3 .format("jdbc") 4 .option("url", "jdbc:TAOS-WS://taosadapter:6041/iot") 5 .option("driver", "com.taosdata.jdbc.ws.WebSocketDriver") 6 .option("user", "root") 7 .option("password", "taosdata") 8 .option("dbtable", "(SELECT * FROM meters WHERE ts > NOW - 1d) tmp") 9 .load() 10 11df.show() 12 13 14// 写 TDengine 15processedDf.write 16 .format("jdbc") 17 .mode("append") 18 .option("url", "jdbc:TAOS-WS://taosadapter:6041/iot") 19 .option("driver", "com.taosdata.jdbc.ws.WebSocketDriver") 20 .option("dbtable", "processed_meters") 21 .save() 22
8. Logstash 集成
1# logstash.conf 2input { 3 file { 4 path => "/var/log/sensors/*.log" 5 codec => json 6 } 7} 8 9filter { 10 date { 11 match => ["ts", "ISO8601"] 12 } 13} 14 15output { 16 http { 17 url => "http://taosadapter:6041/influxdb/v1/write?db=logstash" 18 http_method => "post" 19 user => "root" 20 password => "taosdata" 21 format => "message" 22 message => "sensors,device=%{device} value=%{value} %{ts}" 23 } 24} 25
DBeaver SQL IDE
1连接配置: 2 3Driver: 4 下载 taos-jdbcdriver-3.x.x.jar 5 6连接: 7 URL: jdbc:TAOS-WS://localhost:6041/test 8 User: root 9 Password: taosdata 10 11 12功能: 13 - 表浏览 14 - SQL 编辑 15 - 数据导出 16 - 可视化查询 17
Hive 集成
1-- 用 Hive 查 TDengine(通过 JDBC StorageHandler) 2CREATE EXTERNAL TABLE hive_meters ( 3 ts TIMESTAMP, 4 current FLOAT, 5 voltage INT, 6 location STRING 7) 8STORED BY 'org.apache.hive.storage.jdbc.JdbcStorageHandler' 9TBLPROPERTIES ( 10 "hive.sql.database.type" = "MYSQL", -- 借用 MySQL Driver 11 "hive.sql.jdbc.driver" = "com.taosdata.jdbc.TSDBDriver", 12 "hive.sql.jdbc.url" = "jdbc:TAOS://taosd:6030/iot", 13 "hive.sql.dbcp.username" = "root", 14 "hive.sql.dbcp.password" = "taosdata", 15 "hive.sql.table" = "meters" 16); 17 18SELECT * FROM hive_meters LIMIT 10; 19
代码示例
Telegraf 监控 + Grafana 完整链路
1# docker-compose.yml 2services: 3 taosd: 4 image: tdengine/tdengine:3.x.x 5 6 adapter: 7 image: tdengine/tdengine:3.x.x 8 command: taosadapter 9 ports: ["6041:6041", "6045:6045"] 10 11 telegraf: 12 image: telegraf:latest 13 volumes: 14 - ./telegraf.conf:/etc/telegraf/telegraf.conf 15 depends_on: 16 - adapter 17 18 grafana: 19 image: grafana/grafana:latest 20 ports: ["3000:3000"] 21 environment: 22 - GF_INSTALL_PLUGINS=tdengine-datasource 23
一站式系统监控部署
1#!/bin/bash 2# install_monitoring.sh 3 4# 1. Telegraf 安装 5apt install -y telegraf 6 7# 2. 配置 Telegraf → TDengine 8cat > /etc/telegraf/telegraf.conf <<EOF 9[[inputs.cpu]] 10[[inputs.mem]] 11[[inputs.disk]] 12[[inputs.net]] 13 14[[outputs.http]] 15 url = "http://taosadapter:6041/influxdb/v1/write?db=monitoring" 16 username = "root" 17 password = "taosdata" 18 data_format = "influx" 19EOF 20 21# 3. 启动 22systemctl start telegraf 23 24# 4. 验证数据 25sleep 30 26taos -s "USE monitoring; SHOW TABLES;" 27
性能考量
集成性能对比
| 集成 | 单实例吞吐 |
|---|---|
| Telegraf | 几万指标/秒 |
| collectd | 几万指标/秒 |
| Prometheus remote_write | 几十万样本/秒 |
| Kafka Connect | 几十万行/秒 |
| Flink JDBC Sink | 几十万行/秒 |
| Spark JDBC | 视分区 |
选型建议
| 场景 | 推荐 |
|---|---|
| 系统监控 | Telegraf |
| 应用指标 | StatsD / Prometheus |
| Kafka 消息 | Kafka Connect / taosX |
| 流处理 | Flink |
| 批分析 | Spark |
| 日志 | Logstash + Schemaless |
FAQ
Q1: 用 InfluxDB 兼容协议有何限制?
- 通用写入功能完整
- 不支持 InfluxQL(用 TDengine SQL)
- 部分 Flux 函数无对应
Q2: Kafka Connect 推荐 taosX 还是 JDBC Sink?
- 简单场景:JDBC Sink 即可
- 复杂转换/高吞吐:taosX
Q3: Flink CDC 接 TDengine?
通过 JDBC Sink 写 TDengine。或用 Flink CDC 源 + 自定义 Sink。
Q4: Spark Streaming 写 TDengine?
可用 foreachBatch 调用 JDBC:
1stream.foreachBatch { (df, _) => 2 df.write.format("jdbc")... 3} 4
Q5: 国产生态怎么对接?
- DolphinScheduler:JDBC
- Apache SeaTunnel:连接器支持
- DataX:插件
- 海豚调度等:JDBC 标准接口
参考
系统构架篇
- 01-《TDengine 整体架构全景》
- 02-《集群拓扑深度解析》
- 03-《MNode 内部机制深度解析》
- 04-《RPC 通信层深度解析》
- 05-《VNode 生命周期》
- 06-《RAFT 共识协议》
- 07-《端到端的消息流》
数据模型
- 01-《数据库创建与参数详解》
- 02-《超级表/子表/普通表》
- 03-《支持数据类型深度解析》
- 04-《TDengine Tag 设计哲学与 Schema 变更机制》
- 05-《TDengine 虚拟表实现原理》
存储引擎
- 01-《TDengine 存储引擎概览》
- 02-《TDengine MemTable 深度解析》
- 03-《TDengine WAL 预写日志机制》
- 04-《TDengine 数据文件格式》
- 05-《TDengine Commit 与 Flush 机制 》
- 06-《TDengine Compaction 合并策略 》
- 07-《TDengine 数据保留与 TTL》
- 08-《TDengine 压缩编码机制》
- 09-《TDengine Cache 与 Last 查询加速》
- 10-《TDengine 逻辑计划生成》
查询引擎
- 01-《TDengine 查询引擎概览》
- 02-《TDengine SQL 解析与词法分析》
- 03-《TDengine 语义分析与 AST 重写》
- 04-《TDengine 逻辑计划生成》
- 05-《TDengine 物理计划生成》
- 06-《TDengine 扫描算子》
- 07-《TDengine 聚合算子》
- 08-《TDengine 连接算子》
- 09-《TDengine 排序、填充与投影》
- 10-《TDengine 分布式查询执行》
- 11-《TDengine EXPLAIN 与查询优化》
数据写入
- 01-《TDengine SQL INSERT》
- 02-《TDengine 无模式写入》
- 03-《TDengine STMT 写入》
- 04-《TDengine 写入内部流程》
- 05-《TDengine 数据更新删除》
数据订阅
- 01-《TDengine 数据订阅》
- 02-《TDengine 订阅 vs Kafka》
- 03-《TDengine TMQ 消费流程》
- 04-《TDengine 内部机制》
- 05-《TDengine TMQ 最佳实践》
预聚合
索引
SQL 语句
- 01-《TDengine DDL》
- 02-《TDengine DML SELECT》
- 03-《TDengine DML 函数完整参考》
- 04-《TDengine JOIN 完整语法》
- 05-《TDengine 窗口完整语法》
- 06-《TDengine 操作符与表达式》
- 07-《TDengine 系统表》
- 08-《TDengine SQL 与标准 SQL 差异》
客户端与连接器
- 01-《TDengine 的连接方式》
- 02-《TDengine C/C++ 连接器》
- 03-《TDengine java 连接器》
- 04-《TDengine Python 连接器》
- 05-《TDengine Go 与 Rust 连接器》
- 06-《TDengine Node.js 与 C# 连接器》
运维
- 01-《TDengine 部署指南》
- 02-《TDengine 配置详解》
- 03-《TDengine 监控系统》
- 04-《TDengine 备份与恢复》
- 05-《TDengine 版本升级》
- 06-《TDengine 加密使用指南》
安全
生态
关于 TDengine
TDengine 专为物联网IoT平台、工业大数据平台设计。其中,TDengine TSDB 是一款高性能、分布式的时序数据库(Time Series Database),同时它还带有内建的缓存、流式计算、数据订阅等系统功能;TDengine IDMP 是一款AI原生工业数据管理平台,它通过树状层次结构建立数据目录,对数据进行标准化、情景化,并通过 AI 提供实时分析、可视化、事件管理与报警等功能。
《TDengine 第三方工具 — Telegraf、Kafka Connect、Flink、Spark》 是转载文章,点击查看原文。

