TDengine 数据发布(Data Publisher)无代码投递全指南:MQTT、Kafka、Flink 与 Parquet 实时数据分发实战
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
TDengine 提供的数据发布(Data Publisher)能力,可以在不改写业务代码的前提下,把数据库中的实时数据推送到 MQTT、Kafka、Flink 等主流消息队列与流处理系统,或将 SQL 查询结果导出为 Parquet 文件。本指南以 docs/en/08-data-ingest-and-delivery/02-no-code-delivery 系列文档为主线,结合仓库中的 taosX 组件文档、Flink 连接器示例代码与配置说明,完整讲解每种投递目标的配置参数、操作步骤与验证方法,帮助你按需搭建 IoT 数据采集、实时监控与大数据分析场景下的数据分发链路。
数据发布能力概览
TDengine 的数据发布功能支持将实时数据推送到各类流处理与消息队列系统,显著增强数据的流动性(data mobility)与系统集成能力,可满足 IoT、大数据分析、实时监控等多样化场景的需求。目前支持向以下主流平台发布数据:
| 目标平台 | 能力说明 |
|---|---|
| MQTT | 配置后可将数据实时推送到 MQTT 服务器。MQTT 是 IoT 设备间通信广泛使用的轻量级消息协议,通过该能力可便捷地将传感器数据、设备状态等信息分发到各终端,实现高效的数据共享与交互。 |
| Kafka | 支持向 Kafka 集群发布数据。Kafka 是高吞吐的分布式消息队列系统,常用于大数据采集、日志聚合与流处理场景;与 Kafka 集成可将实时数据无缝接入大数据平台,支撑数据分析、监控、告警等应用。 |
| Flink | 可将数据流式投递到 Flink。Flink 是适合实时数据分析与复杂事件处理的高性能流处理引擎,借助该能力可构建端到端的实时数据处理管道,满足低延迟、高可靠需求。 |
| Parquet | 可将一条只读 SQL 查询的结果导出为 taosX 服务节点上的单个 Parquet 文件,适用于离线分析、基于文件的交换以及消费 Parquet 文件的下游系统。 |
版本限制:以上数据发布(Data Publisher)相关特性仅存在于TDengine TSDB-Enterprise(企业版)中,TSDB-OSS(社区版)不包含这些组件与功能(见 docs/en/08-data-ingest-and-delivery/resources/_resources.mdx)。
前置准备:企业服务与 taosX
数据发布依赖 taosX 组件(TDengine 企业版提供零代码数据接入/投递能力的核心组件)。在创建发布任务前,需要确认:
taosd服务正常运行;taosAdapter服务正常运行(Flink、Kafka 等场景需要);- 已安装 taosX,可用
taosx --version验证版本; - 若在 taosExplorer 图形界面中使用,需先完成 taosX 服务模式部署。
taosX 的两种运行模式
taosX 支持两种运行模式(详见 taosX 组件文档):
- 命令行模式:通过
taosx命令直接执行一次性的数据迁移/发布任务,命令格式为:
taosx -f <from-DSN> -t <to-DSN> <other parameters>其中-f指定数据源(Source DSN),-t指定写入目标(Sink DSN)。常用参数还包括:--jobs <number>指定并发任务数(仅支持 TMQ 任务)、-v/-vv/-vvv分别开启 info/debug/trace 级别日志。
- 服务模式:以系统服务方式运行(Linux 下
systemctl start taosx,Windows 下sc.exe start taosx),通过 taosExplorer 图形界面使用其功能。
DSN 采用类 URL 格式:<driver>[+<protocol>]://[[<username>:<password>@]<host>:<port>][/<object>][?<p1>=<v1>[&<p2>=<v2>]]。常用 driver 有taos(查询接口取数)、tmq(订阅取数)、mqtt、kafka、parquet等;+ws协议后缀表示通过 REST/WebSocket 取数(此时 taosx 可安装在非服务器节点),不带后缀则使用原生连接(taosx 必须与数据库同机部署)。
taosX 服务配置(taosx.toml)
服务模式的默认配置文件路径为:Linux/etc/taos/taosx.toml,WindowsC:\TDengine\cfg\taosx.toml。关键配置项包括:
data_dir:数据文件存储目录;serve.listen:REST API 监听地址,默认0.0.0.0:6050,支持 IPv6 与同端口多地址逗号分隔;serve.grpc:gRPC 监听地址,默认0.0.0.0:6055;monitor.fqdn / monitor.port / monitor.interval:taosKeeper 监控上报配置,默认 6043 端口、每 10 秒上报一次;log.path / log.level / log.rotationCount / log.rotationSize / log.keepDays:日志目录、级别与轮转保留策略(旧版logs_home、log_level、log_keep_days已弃用);jobs:tokio 运行时工作线程数,默认 0(表示核心数×2)。
准备订阅数据:数据库、超级表与 Topic
无论目标平台是 MQTT 还是 Kafka,发布的数据都来源于 TMQ(TDengine 消息队列)订阅的 Topic。可用taosCLI 或 taosExplorer 执行以下 SQL,创建数据库、超级表、Topic 并写入测试数据:
create database db vgroups 1; create table db.meters (ts timestamp, f1 int) tags(t1 int); create topic topic_meters as select ts, tbname, f1, t1 from db.meters; insert into db.tb using db.meters tags(1) values(now, 1);Topic 本质上是基于 SQL 查询定义的数据订阅集合,关于 Topic 定义、消费偏移与订阅参数的更多细节,参见 数据订阅文档。
MQTT 数据发布
MQTT(Message Queuing Telemetry Transport)是基于发布/订阅模型的轻量级消息协议,广泛用于 IoT 设备间通信。TDengine 支持将实时数据推送到 MQTT 服务器,从而便捷地将传感器数据、设备状态等分发到各类终端。
安装并配置 MQTT 服务器
使用该功能前需先部署 MQTT 服务器(如 Mosquitto、EMQX 等),可根据需要选择合适的实现。
创建 MQTT 数据发布任务
数据发布通过 taosx 命令行完成,将 TMQ Topic 数据发布到 MQTT:
taosx run -f "tmq+ws://username:password@ip:port/topic?param=value..." -t "mqtt://ip:port?param=value..."其中-f为 TMQ 订阅 DSN,-t为 MQTT broker DSN。taosx 与 DSN 的完整用法参见 taosX 组件文档。
TMQ DSN 参数:
| 参数 | 说明 |
|---|---|
username/password | 数据库用户名与密码 |
ip/port | 数据库连接地址与端口 |
topic | TMQ 订阅的 Topic 名称 |
with.meta | 是否同步建表、删表、改表、删数据等元数据,默认false(不同步) |
with.meta.delete | 是否同步元数据中的删除数据事件,仅在启用with.meta时生效 |
with.meta.drop | 是否同步元数据中的删表事件,仅在启用with.meta时生效 |
group.id | TMQ 订阅参数,必填,订阅所属消费组 ID |
client.id | TMQ 订阅参数,可选,订阅客户端 ID |
auto.offset.reset | 订阅起始位置 |
experimental.snapshot.enable | 是否同步已落盘到 TSDB 存储文件(非 WAL)中的数据;关闭时仅同步仍在 WAL 中的数据 |
更多 TMQ 订阅参数参见 数据订阅文档。
MQTT DSN 参数:
| 参数 | 说明 |
|---|---|
ip/port | MQTT broker 地址与端口 |
version | MQTT 协议版本,必填,可选3.1/3.1.1/5.0 |
qos | MQTT QoS 级别,默认0 |
client_id | MQTT 客户端 ID,必填,每个客户端必须唯一 |
topic | 数据发布的 MQTT 目标 Topic,必填 |
meta_topic | 元数据发布的 Topic,不指定时默认与数据 Topic 相同 |
MQTT 的topic与meta_topic中支持以下模板变量:
database:源数据库名(元数据与数据消息中均包含);tmq_topic:源 TMQ Topic 名(元数据与数据消息中均包含);vgroup_id:源 vgroup ID(元数据与数据消息中均包含);stable:源超级表名(仅包含在创建超级表、创建子表、删除超级表的元数据消息中);table:源表/子表名(仅包含在创建子表/普通表、修改表、删除子表/普通表及数据消息中)。
注意:如果 Topic 中包含消息里不存在的变量,该消息将不会被处理,也不会发布到 MQTT broker。
完整示例命令:
taosx run \ -f "tmq+ws://root:taosdata@localhost:6041/topic_meters?group.id=taosx-pub-test&auto.offset.reset=earliest" \ -t "mqtt://mqtt.tdengine.com:1883?topic=test/topic_meters&qos=1&version=5.0&client_id=taosx-pub-meters"验证数据发布
可使用 MQTTX 等工具订阅目标 Topic 验证发布结果。上述示例发布的消息体格式如下:
{"data":{"ts":1756957064991,"tbname":"tb","f1":1,"t1":1},"offset":{"database":"db","topic":"topic_meters","vgroupId":2,"offset":8}}其中data为实际数据行,offset携带源库、源 Topic、vgroup 与消息偏移信息,可用于下游做去重与位点追踪。
Kafka 数据发布
TDengine 可将 TMQ 数据消息与元数据消息发布到 Kafka,使时序数据得以转发到数据平台、实时计算引擎及下游业务系统。在 taosExplorer 中创建 Kafka 发布任务后,系统会从指定 TMQ Topic 读取消息,并按配置发布到一个或多个 Kafka Topic;数据消息与元数据消息可分别控制,保存前还可执行连通性检查与消息预览。
发布任务运行流程
任务创建成功后的运行顺序一般为:
- 从
Topic DSN指定的 TMQ 地址订阅消息; - 依据
Enable Data Subscription与Enable Meta Subscription决定读取哪些消息; - 连接目标 Kafka 集群;
- 将数据消息发布到数据 Topic、元数据消息发布到元数据 Topic;
- 若启用了自动建 Topic,目标 Topic 不存在时尝试创建;
- 保存前执行连通性检查与预览,以验证链路与消息格式。
基础配置
任务名称(Task Name):必填,用于标识当前发布任务。建议在名称中包含源 Topic、目标用途与环境信息,例如prod-device-data-to-kafka、tmq_order_meta_to_kafka。
Topic DSN:必填,指定完整的 TMQ 连接地址,是任务的数据源。建议显式携带group.id、auto.offset.reset等关键订阅参数,使任务行为可预测。常用格式:
tmq+ws://root:taosdata@localhost:6041/topic_meters常见订阅参数包括group.id(消费组 ID)、auto.offset.reset(消费起始位置,在 taosExplorer 中由独立的必填字段Start From控制,映射为earliest/latest)、with.meta、with.meta.delete、with.meta.drop、experimental.snapshot.enable。示例:
tmq+ws://root:taosdata@localhost:6041/topic_meters?group.id=pub-kafka-demo&auto.offset.reset=earliest&with.meta=trueStart From:必填,指定消费者首次启动或没有已提交偏移时的起始消费位置。可选earliest(从最早可读位置开始,适合初次全量验证或回放历史数据)与latest(只消费新到达的消息,适合生产环境长期运行的任务)。
Group ID:可选,标识 TMQ 消费组,默认由系统自动生成。生产环境建议使用固定值,使消费偏移与重启行为可预测。
Enable Data Subscription:默认开启,控制是否发布 TMQ 行数据消息。仅当任务只想发布元数据时才关闭;开启时提交前必须配置Data Topic。
Enable Meta Subscription:默认关闭,控制是否发布建表、删表、结构变更等元数据消息。仅当下游需要表生命周期或结构变更事件时开启;仅启用元数据订阅时必须配置Meta Topic。
TSDB Data:默认开启,控制是否同时订阅已落盘的 TSDB 数据,而不仅是仍在 WAL 中的数据。
Table Deletions / Data Deletions:默认开启,分别控制是否转发删表事件与删数据事件,仅在下游需要同步表生命周期/删除事件时保留开启。
Kafka 连接配置
Bootstrap Servers:必填,Kafka broker 地址列表,多个地址用逗号分隔,建议至少配置两个可达地址以提升可用性,例如:
127.0.0.1:9092,127.0.0.1:9093SASL Authentication:默认关闭,配置 Kafka broker 的 SASL 认证机制与参数。支持PLAIN、SCRAM-SHA-256、GSSAPI三种机制。规则如下:
- 不选择机制即关闭 SASL;
- 选择
PLAIN或SCRAM-SHA-256时必须配置用户名与密码; - 选择
GSSAPI时必须配置 Kerberos 相关参数:Kerberos 服务名、principal、初始化命令与 keytab 文件,初始化命令示例:
kinit -R -t '%{sasl.kerberos.keytab}' -k %{sasl.kerberos.principal}SSL Authentication:默认关闭,配置证书校验与双向认证参数。开启后需配置 CA、CA 密码、客户端证书、客户端私钥等 PEM 格式证书文件。生产环境启用 TLS 时应提前验证证书链、证书密码与私钥。
Kafka 发布配置
Data Topic:默认taosx.data.out,定义数据消息的 Kafka Topic(启用数据订阅时必填)。建议在 Topic 名中体现环境、业务域或源表信息以简化下游路由。支持模板变量:${database}、${table}、${stable}、${tmq_topic}、${vgroup_id}、${offset}。示例:taosx.data.out、data.${database}.${table}。
Meta Topic:默认空,定义元数据消息的 Topic。元数据与数据需要独立路由时使用独立 Topic;留空时回退到Data Topic。
Data Key Template / Meta Key Template:默认空,分别定义数据/元数据消息的 Kafka key。需要按库、表或设备分区时配置稳定的 key 模板,可用变量为${database}、${table}、${stable}、${tmq_topic}、${vgroup_id}。示例:${database}.${table}。Meta Key Template 留空时默认使用 Data Key Template。
Auto Create Topic:默认关闭,目标 Topic 在 Kafka 中不存在时尝试自动创建。仅当 broker 允许自动建 Topic 且当前账号具备足够权限时开启;创建是否成功仍取决于 broker 设置与账号权限。开启后可额外配置:
- Topic Partitions:自动建 Topic 时的分区数,默认使用 broker 默认配置,取值范围 1~1024,应根据下游消费并发度与预期吞吐规划;
- Replication Factor:自动建 Topic 时的副本因子,默认使用 broker 默认配置,取值范围 1~128,且不能超过 Kafka 集群可用 broker 数量,应与集群高可用策略一致。
高级选项
| 参数 | 默认值 | 说明 |
|---|---|---|
| Parallelism | 1 | Kafka 生产者最大并发度,取值范围 1~128,高吞吐场景可从 1 或 2 起步逐步增大 |
| Queue Timeout (ms) | 30000 | 消息进入发送队列后的最长等待时间,网络不稳定时可增大,追求快速失败时可减小 |
| Batch Size | 1000 | 每个 Kafka 批次的最大记录数,取值范围 1~100000,吞吐优先时增大、低延迟时调小 |
| Batch Timeout (ms) | 1000 | 批次发送前的最大等待时间,吞吐优先时增大、低延迟时减小 |
| Kafka Extra Parameters | 空 | 额外 Kafka 生产者参数,仅配置明确理解的原生参数并避免与标准字段冲突 |
Kafka Extra Parameters 示例:
compression.type=zstd acks=all linger.ms=100常见用途包括启用消息压缩、设置acks、调整linger.ms等发送策略。
保存前的校验规则
提交任务前通常会执行以下逻辑校验:
Topic DSN必须填写,且必须显式选择Start From;Enable Data Subscription与Enable Meta Subscription不能同时关闭;- 启用数据订阅时,
Data Topic必须填写; - 关闭数据订阅但启用元数据订阅时,
Meta Topic必须填写; - 提交前执行 TMQ 与 Kafka 连通性检查,检查失败则无法保存任务。
部分字段会根据其他设置动态显示:未选择 SASL 机制时隐藏 SASL 详情字段;关闭 SSL 认证时隐藏证书字段;关闭Auto Create Topic时隐藏分区与副本字段。
连通性检查与预览
保存前建议执行连通性检查,验证:TMQ 地址是否可达、Kafka broker 地址是否可达、SASL/SSL 参数是否正确、所需 Topic 与权限是否可用。检查失败时优先排查网络、认证、地址与权限配置。
预览功能可在保存前查看将生成的样例 Kafka 消息,常用参数为:Rows(默认1,范围 1~100)、Wait time seconds(默认30,范围 1~300)。预览结果通常展示topic、key、value三个字段;若等待时间内未收到数据,系统会提示当前条件下无可预览消息。
验证发布结果
任务保存并启动后,可用 Kafka 自带工具或第三方客户端验证消息发布是否正确。例如使用kcat消费目标 Topic:
kcat -b 127.0.0.1:9092 -t taosx.data.out -C配置了 key 模板时,消费输出通常同时包含消息 key 与消息体;可结合源数据库、表名、Topic 模板与偏移字段共同验证发布结果。
Flink 集成:TDengine Flink Connector
Apache Flink 是 Apache 软件基金会支持的开源分布式流批一体处理框架,可用于流处理、批处理、复杂事件处理、实时数仓构建,并为机器学习提供实时数据支撑。借助 TDengine 的 Flink Connector,Flink 可与 TDengine 无缝集成,高效稳定地从数据库读取海量数据并进行分析处理。
版本说明:Flink Connector 相关能力仅存在于 TDengine TSDB Enterprise;社区版集群可作为被读取的数据源,但 Connector 完整能力以企业版为准。Flink 需 v1.19.0 及以上,taosAdapter 需正常运行。
引入依赖
Maven 项目在pom.xml中添加:
<dependency> <groupId>com.taosdata.flink</groupId> <artifactId>flink-connector-tdengine</artifactId> <version>2.1.4</version> </dependency>Connector 版本历史(如 2.1.4 升级 JDBC 驱动至 3.7.3、2.0.0 起支持 Table SQL 写入等)参见 Flink 公共信息 与 Java 连接器版本历史。
连接参数
连接参数由 URL 与 Properties 组成。URL 规范格式:
jdbc: TAOS-WS://[host_name]:[port]/[database_name]?[user={user}|&password={password}|&timezone={timezone}]| 参数 | 说明 |
|---|---|
| user | 登录用户名,默认root |
| password | 登录密码,默认taosdata |
| database_name | 数据库名 |
| timezone | 时区 |
| HttpConnectTimeout | 连接超时时间(毫秒),默认 60000 |
| MessageWaitTimeout | 消息超时时间(毫秒),默认 60000 |
| UseSSL | 连接是否使用 SSL |
Source:并行读取与三种分片方式
Source 从 TDengine 读取数据并转换为 Flink 可处理的格式,通过设置数据源并行度可多线程并行读取,提升读取效率与吞吐。Source 关键属性(TDengineConfigParams)包括:
PROPERTY_KEY_USER/PROPERTY_KEY_PASSWORD:用户名/密码(默认root/taosdata);VALUE_DESERIALIZER:结果集反序列化方式,收到RowData类型时设为RowData;也可继承TDengineRecordDeserialization实现convert与getProducedType自定义反序列化;TD_BATCH_MODE:是否批量推送数据到下游算子,为 True 时需将数据类型指定为SourceRecords的模板形式;PROPERTY_KEY_MESSAGE_WAIT_TIMEOUT:消息超时(毫秒),默认 60000;PROPERTY_KEY_ENABLE_COMPRESSION:传输过程是否启用压缩,默认 false;PROPERTY_KEY_ENABLE_AUTO_RECONNECT:是否启用自动重连,默认 false;PROPERTY_KEY_RECONNECT_INTERVAL_MS:重连重试间隔(毫秒),默认 2000,仅自动重连开启时生效;PROPERTY_KEY_RECONNECT_RETRY_COUNT:自动重连重试次数,默认 3,仅自动重连开启时生效;PROPERTY_KEY_DISABLE_SSL_CERT_VALIDATION:是否关闭 SSL 证书校验,默认 false。
按时间分片:根据开始时间、结束时间、分片间隔与时间字段名,将 SQL 查询拆分为多个子任务并行取数(时间区间左闭右开)。示例代码见 docs/examples/flink/source/Main.java 中的time_interval片段:
SourceSplitSql splitSql = new SourceSplitSql(); splitSql.setSql("select ts, `current`, voltage, phase, groupid, location, tbname from meters") .setSplitType(SplitType.SPLIT_TYPE_TIMESTAMP) .setTimestampSplitInfo(new TimestampSplitInfo( "2024-12-19 16:12:48.000", "2024-12-19 19:12:48.000", "ts", Duration.ofHours(1), new SimpleDateFormat("yyyy-MM-dd HH:mm:ss.SSS"), ZoneId.of("Asia/Shanghai")));按超级表 TAG 分片:根据 TAG 字段将查询 SQL 拆分为多个查询条件,每个条件对应一个子任务并行取数:
SourceSplitSql splitSql = new SourceSplitSql(); splitSql.setSql("select ts, current, voltage, phase, groupid, location from meters where voltage > 100") .setTagList(Arrays.asList("groupid >100 and location = 'Shanghai'", "groupid >50 and groupid < 100 and location = 'Guangzhou'", "groupid >0 and groupid < 50 and location = 'Beijing'")) .setSplitType(SplitType.SPLIT_TYPE_TAG);按表分片:输入多个表结构相同的超级表或普通表,系统按"一表一任务"拆分后并行取数:
SourceSplitSql splitSql = new SourceSplitSql(); splitSql.setSelect("ts, current, voltage, phase, groupid, location") .setTableList(Arrays.asList("d1001", "d1002")) .setOther("order by ts limit 100") .setSplitType(SplitType.SPLIT_TYPE_TABLE);使用 Source 连接器:以RowData为例,通过TDengineSource<RowData>创建数据源并注册到 Flink 环境:
Properties connProps = new Properties(); connProps.setProperty(TDengineConfigParams.PROPERTY_KEY_ENABLE_AUTO_RECONNECT, "true"); connProps.setProperty(TDengineConfigParams.PROPERTY_KEY_TIME_ZONE, "UTC-8"); connProps.setProperty(TDengineConfigParams.VALUE_DESERIALIZER, "RowData"); connProps.setProperty(TDengineConfigParams.TD_JDBC_URL, "jdbc:TAOS-WS://localhost:6041/power?user=root&password=taosdata"); StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); TDengineSource<RowData> source = new TDengineSource<>(connProps, sql, RowData.class); DataStreamSource<RowData> input = env.fromSource(source, WatermarkStrategy.noWatermarks(), "tdengine-source");批模式(TD_BATCH_MODE=true,类型为SourceRecords模板形式)与自定义类型(ResultBean+ 继承TDengineRecordDeserialization的ResultSourceDeserialization)的完整示例,见 source/Main.java 中的source_batch_test与source_custom_type_test片段。
CDC 数据订阅
Flink CDC 提供数据订阅功能,可实时监听 TDengine 数据变化并以数据流形式传给 Flink 处理,同时保证数据一致性与完整性。CDC 关键参数(TDengineCdcParams)包括:
BOOTSTRAP_SERVERS:TDengine 服务器ip:port,WebSocket 连接时为 taosAdapter 所在地址;CONNECT_USER/CONNECT_PASS:用户名/密码(默认root/taosdata);POLL_INTERVAL_MS:拉取数据间隔,默认 500ms;VALUE_DESERIALIZER:结果集反序列化方式;可继承com.taosdata.jdbc.tmq.ReferenceDeserializer指定结果集 Bean,或继承com.taosdata.jdbc.tmq.Deserializer自定义;TMQ_BATCH_MODE:批量推送模式,为 True 时类型需指定为ConsumerRecords模板形式;GROUP_ID:消费组 ID,同组共享消费进度,最大长度 192;AUTO_OFFSET_RESET:消费组订阅起始位置(earliest从最早订阅、latest从最新订阅,默认latest);ENABLE_AUTO_COMMIT:是否自动提交消费点位,true 自动提交、false 按 checkpoint 提交,默认 false。注意:自动提交模式在获取数据后即提交,无论下游算子是否正确处理,存在数据丢失风险,主要用于无状态算子或一致性要求低的场景;AUTO_COMMIT_INTERVAL_MS:自动提交消费记录的时间间隔(毫秒),默认 5000,仅ENABLE_AUTO_COMMIT=true时生效;TMQ_SESSION_TIMEOUT_MS:消费者心跳丢失后的超时时间,触发 rebalance 后剔除该消费者(3.3.3.0 起支持),默认 12000,范围 [6000, 1800000];TMQ_MAX_POLL_INTERVAL_MS:消费者 poll 拉取的最长间隔,超时视为消费者离线并触发 rebalance(3.3.3.0 起支持),默认 300000,范围 [1000, INT32_MAX]。
CDC 连接器会根据用户设置的并行度创建消费者,因此应依据资源情况合理设置并行度。RowData、批量(ConsumerRecords)与自定义类型的 CDC 示例,分别见 source/Main.java 的cdc_source、cdc_batch_source、cdc_custom_type_test片段;其中自定义类型要求ResultBean的字段名与数据类型和列名及类型一一对应。
Table SQL 集成
通过 Table SQL 可从多个数据源(TDengine、MySQL、Oracle 等)抽取数据,执行清洗、格式转换、多表关联等算子操作后,将结果加载到目标数据源。Source 连接器参数(connector = 'tdengine-connector'):
| 参数名 | 类型 | 说明 |
|---|---|---|
| connector | string | 连接器标识,固定为tdengine-connector |
| td.jdbc.url | string | 连接 URL |
| td.jdbc.mode | string | 连接器类型:source、sink |
| table.name | string | 源或目标表名 |
| scan.query | string | 取数 SQL 语句 |
| sink.db.name | string | 目标数据库名 |
| sink.supertable.name | string | 目标超级表名 |
| sink.batch.size | integer | 批量写入大小 |
| sink.table.name | string | 子表或普通表名 |
Table CDC 连接器参数在此基础上增加user、password、bootstrap.servers、topic、td.jdbc.mode(cdc、sink)、group.id、auto.offset.reset(earliest/latest,默认latest)、poll.interval_ms(默认 500ms)。完整的 Table Source 与 Table CDC 示例见 source/Main.java 的source_table与cdc_table片段,例如将power库meters超级表的子表数据写入power_sink库sink_meters超级表对应的子表:
CREATE TABLE `meters` (...) WITH ( 'connector' = 'tdengine-connector', 'td.jdbc.url' = 'jdbc:TAOS-WS://localhost:6041/power?user=root&password=taosdata', 'td.jdbc.mode' = 'source', 'table-name' = 'meters', 'scan.query' = 'SELECT ts, `current`, voltage, phase, location, groupid, tbname FROM `meters`' ); -- CDC 模式则将 td.jdbc.mode 设为 'cdc',并配置 bootstrap.servers、group.id、topic 等处理语义与类型映射
由于 TDengine 不支持事务、无法频繁执行 checkpoint 与复杂事务协调,且使用时间戳作为主键(下游算子可对重复数据过滤去重),Connector 采用At-Least-Once语义以保障处理性能与低延迟:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);TDengine 与 Flink RowData 类型的映射关系:TIMESTAMP→TimestampData、INT→Integer、BIGINT→Long、FLOAT→Float、DOUBLE→Double、SMALLINT→Short、TINYINT→Byte、BOOL→Boolean、VARCHAR/BINARY/NCHAR/JSON→StringData、VARBINARY/GEOMETRY→byte[]。任务执行失败后,可依据 Flink 任务日志与错误码排查,常见错误码如0xa000(连接参数错误)、0xa013(未设置 value.deserializer)、0x231d(连接超时,可通过增加 httpConnectTimeout 或检查 taosAdapter 解决)、0x231e(任务超时,可增加 messageWaitTimeout)等,完整错误码表见 Flink 公共信息。
Parquet 数据导出
Parquet Data Out 将一条只读 TDengine SQL 查询的结果导出为 taosX 服务节点上的单个本地 Parquet 文件,输出路径由 taosX 服务端解释,而非浏览器端。
创建 Parquet 导出任务
在 taosExplorer 的 Data Publisher 页面选择 Parquet 作为目标类型,配置以下参数:
- TDengine DSN:TDengine 连接地址,例如
taos+ws://root:taosdata@localhost:6041/db; - SQL Query:一条只读
SELECT查询,支持WITH ... SELECT查询; - Output File:taosX 服务节点上的输出文件路径,必须以
.parquet结尾; - Overwrite Existing File:新的临时文件成功关闭后,是否允许替换已有的最终文件;
- Compression:
uncompressed、zstd、snappy、gzip、brotli或lz4_raw; - Compression Level:
zstd、gzip、brotli可选; - Row Group Size:每个 Parquet row group 的最大行数,默认 131072。Parquet 写入器仅在 row group 写满或文件关闭时才落盘;取值越小
.part文件增长越频繁,但会降低压缩率并增加文件元数据开销。有效范围 1024~10000000。
输出路径规则
- 输出文件创建在 taosX 服务节点上,不是浏览器本地路径;
- 相对路径写入
$DATA_DIR/tasks/<task_id>/<job_id>/目录下; - 绝对路径按 taosX 服务节点上的绝对路径处理;
- 写入器先在相同目录创建
<final_name>.part临时文件,关闭 Parquet 写入器并校验元数据后,再将临时文件重命名为最终路径; - 第一版不支持 agent 或 via 执行。
DSN 示例
FROM 'taos+ws://root:taosdata@localhost:6041/db?query=select%20*%20from%20meters' TO 'parquet:/tmp/meters.parquet?overwrite=false&compression=zstd&row_group_size=131072'对应命令行方式可参考 taosX 组件文档 中的 SQL 查询结果导出用法,例如:
taosx run -f "taos+ws://root:taosdata@localhost:6041/test?query=select * from test.meters" \ -t "parquet:./test.parquet"需注意 SQL 查询语句需进行 URL 编码(尤其是含特殊字符时),并控制查询结果集大小避免内存溢出。从源码结构看,taosx 的 Parquet 读取驱动也支持batch_size(默认 1000)、projection(列投影,可按列名或从 0 开始的索引)、unprocessed_batches(背压控制,默认 64)等参数,Parquet 类型与 TDengine 类型存在自动映射(BOOLEAN→BOOL、INT32→INT、INT64→BIGINT、FLOAT→FLOAT、DOUBLE→DOUBLE、BYTE_ARRAY(UTF8)→NCHAR、BYTE_ARRAY(Binary)→BINARY、INT96(Timestamp)→TIMESTAMP)。
限制
- 仅允许一条只读
SELECT查询; - 不允许
SHOW、DESCRIBE、DESC及任何修改数据、结构、会话或权限的语句; - 不支持目录输出与多文件拆分;
- 导出失败从头重新开始;
- 不支持追加到已有 Parquet 文件;
- 第一版任务页面不提供 Parquet 下载操作;
- 任务配置页面不估算行数、文件大小或剩余时间。
性能与可观测性
大规模导出可能长时间运行并持续消耗 taosX 节点的磁盘、CPU 与网络资源。TDengine 结果块以流式方式写入,并作为 Arrow record batches 传给 Parquet 写入器。任务指标包括 Parquet 输出行数、批次、块数、字节数、当前文件大小、耗时、查询时间、写入时间、关闭时间与失败批次数;活动日志会记录查询开始、首个结果块、进度、写入器关闭、完成、取消与失败等事件。
总结
TDengine 的数据发布(Data Publisher)以 TMQ 订阅与 taosX 组件为底座,将数据库实时数据零代码投递到 MQTT、Kafka、Flink 等主流消息与流处理系统,或将查询结果导出为 Parquet 文件:
- MQTT:通过
taosx run -f <tmq DSN> -t <mqtt DSN>命令行一键发布,支持 QoS、协议版本、元数据同步与主题模板变量; - Kafka:通过 taosExplorer 图形界面配置,数据/元数据消息独立控制,提供 SASL/SSL 认证、key 模板、自动建 Topic、批量与并发调优、保存前连通性检查与消息预览;
- Flink:通过 Flink Connector 实现并行 Source、时间/TAG/按表分片、CDC 数据订阅与 Table SQL 集成,采用 At-Least-Once 语义保障性能;
- Parquet:将只读 SQL 结果导出为单个 Parquet 文件,支持多种压缩算法与 row group 调优,适合离线分析与文件交换。
相关配置细节、参数默认值与注意事项可在本文引用的 MQTT、Flink、Kafka、Parquet、taosX 组件文档 与 数据订阅文档 中继续深入查阅。
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考