TDengine 数据发布(Data Publisher)无代码投递全指南:MQTT、Kafka、Flink 与 Parquet 实时数据分发实战
2026/9/13 6:20:47 网站建设 项目流程

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(订阅取数)、mqttkafkaparquet等;+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_homelog_levellog_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数据库连接地址与端口
topicTMQ 订阅的 Topic 名称
with.meta是否同步建表、删表、改表、删数据等元数据,默认false(不同步)
with.meta.delete是否同步元数据中的删除数据事件,仅在启用with.meta时生效
with.meta.drop是否同步元数据中的删表事件,仅在启用with.meta时生效
group.idTMQ 订阅参数,必填,订阅所属消费组 ID
client.idTMQ 订阅参数,可选,订阅客户端 ID
auto.offset.reset订阅起始位置
experimental.snapshot.enable是否同步已落盘到 TSDB 存储文件(非 WAL)中的数据;关闭时仅同步仍在 WAL 中的数据

更多 TMQ 订阅参数参见 数据订阅文档。

MQTT DSN 参数

参数说明
ip/portMQTT broker 地址与端口
versionMQTT 协议版本,必填,可选3.1/3.1.1/5.0
qosMQTT QoS 级别,默认0
client_idMQTT 客户端 ID,必填,每个客户端必须唯一
topic数据发布的 MQTT 目标 Topic,必填
meta_topic元数据发布的 Topic,不指定时默认与数据 Topic 相同

MQTT 的topicmeta_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;数据消息与元数据消息可分别控制,保存前还可执行连通性检查与消息预览。

发布任务运行流程

任务创建成功后的运行顺序一般为:

  1. Topic DSN指定的 TMQ 地址订阅消息;
  2. 依据Enable Data SubscriptionEnable Meta Subscription决定读取哪些消息;
  3. 连接目标 Kafka 集群;
  4. 将数据消息发布到数据 Topic、元数据消息发布到元数据 Topic;
  5. 若启用了自动建 Topic,目标 Topic 不存在时尝试创建;
  6. 保存前执行连通性检查与预览,以验证链路与消息格式。

基础配置

任务名称(Task Name):必填,用于标识当前发布任务。建议在名称中包含源 Topic、目标用途与环境信息,例如prod-device-data-to-kafkatmq_order_meta_to_kafka

Topic DSN:必填,指定完整的 TMQ 连接地址,是任务的数据源。建议显式携带group.idauto.offset.reset等关键订阅参数,使任务行为可预测。常用格式:

tmq+ws://root:taosdata@localhost:6041/topic_meters

常见订阅参数包括group.id(消费组 ID)、auto.offset.reset(消费起始位置,在 taosExplorer 中由独立的必填字段Start From控制,映射为earliest/latest)、with.metawith.meta.deletewith.meta.dropexperimental.snapshot.enable。示例:

tmq+ws://root:taosdata@localhost:6041/topic_meters?group.id=pub-kafka-demo&auto.offset.reset=earliest&with.meta=true

Start 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:9093

SASL Authentication:默认关闭,配置 Kafka broker 的 SASL 认证机制与参数。支持PLAINSCRAM-SHA-256GSSAPI三种机制。规则如下:

  • 不选择机制即关闭 SASL;
  • 选择PLAINSCRAM-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.outdata.${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 数量,应与集群高可用策略一致。

高级选项

参数默认值说明
Parallelism1Kafka 生产者最大并发度,取值范围 1~128,高吞吐场景可从 1 或 2 起步逐步增大
Queue Timeout (ms)30000消息进入发送队列后的最长等待时间,网络不稳定时可增大,追求快速失败时可减小
Batch Size1000每个 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等发送策略。

保存前的校验规则

提交任务前通常会执行以下逻辑校验:

  1. Topic DSN必须填写,且必须显式选择Start From
  2. Enable Data SubscriptionEnable Meta Subscription不能同时关闭;
  3. 启用数据订阅时,Data Topic必须填写;
  4. 关闭数据订阅但启用元数据订阅时,Meta Topic必须填写;
  5. 提交前执行 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)。预览结果通常展示topickeyvalue三个字段;若等待时间内未收到数据,系统会提示当前条件下无可预览消息。

验证发布结果

任务保存并启动后,可用 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实现convertgetProducedType自定义反序列化;
  • 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+ 继承TDengineRecordDeserializationResultSourceDeserialization)的完整示例,见 source/Main.java 中的source_batch_testsource_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_sourcecdc_batch_sourcecdc_custom_type_test片段;其中自定义类型要求ResultBean的字段名与数据类型和列名及类型一一对应。

Table SQL 集成

通过 Table SQL 可从多个数据源(TDengine、MySQL、Oracle 等)抽取数据,执行清洗、格式转换、多表关联等算子操作后,将结果加载到目标数据源。Source 连接器参数(connector = 'tdengine-connector'):

参数名类型说明
connectorstring连接器标识,固定为tdengine-connector
td.jdbc.urlstring连接 URL
td.jdbc.modestring连接器类型:sourcesink
table.namestring源或目标表名
scan.querystring取数 SQL 语句
sink.db.namestring目标数据库名
sink.supertable.namestring目标超级表名
sink.batch.sizeinteger批量写入大小
sink.table.namestring子表或普通表名

Table CDC 连接器参数在此基础上增加userpasswordbootstrap.serverstopictd.jdbc.modecdcsink)、group.idauto.offset.resetearliest/latest,默认latest)、poll.interval_ms(默认 500ms)。完整的 Table Source 与 Table CDC 示例见 source/Main.java 的source_tablecdc_table片段,例如将powermeters超级表的子表数据写入power_sinksink_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:新的临时文件成功关闭后,是否允许替换已有的最终文件;
  • Compressionuncompressedzstdsnappygzipbrotlilz4_raw
  • Compression Levelzstdgzipbrotli可选;
  • 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查询;
  • 不允许SHOWDESCRIBEDESC及任何修改数据、结构、会话或权限的语句;
  • 不支持目录输出与多文件拆分;
  • 导出失败从头重新开始;
  • 不支持追加到已有 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询