尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

TDengine 数据发布(Data Publisher)无代码投递全指南:MQTT、Kafka、Flink 与 Parquet 实时数据分发实战

TDengine 数据发布(Data Publisher)无代码投递全指南:MQTT、Kafka、Flink 与 Parquet 实时数据分发实战 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/TDengineTDengine 提供的数据发布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 taosxWindows 下sc.exe start taosx通过 taosExplorer 图形界面使用其功能。DSN 采用类 URL 格式driver[protocol]://[[username:password]host:port][/object][?p1v1[p2v2]]。常用 driver 有taos查询接口取数、tmq订阅取数、mqtt、kafka、parquet等ws协议后缀表示通过 REST/WebSocket 取数此时 taosx 可安装在非服务器节点不带后缀则使用原生连接taosx 必须与数据库同机部署。taosX 服务配置taosx.toml服务模式的默认配置文件路径为Linux/etc/taos/taosx.tomlWindowsC:\TDengine\cfg\taosx.toml。关键配置项包括data_dir数据文件存储目录serve.listenREST API 监听地址默认0.0.0.0:6050支持 IPv6 与同端口多地址逗号分隔serve.grpcgRPC 监听地址默认0.0.0.0:6055monitor.fqdn / monitor.port / monitor.intervaltaosKeeper 监控上报配置默认 6043 端口、每 10 秒上报一次log.path / log.level / log.rotationCount / log.rotationSize / log.keepDays日志目录、级别与轮转保留策略旧版logs_home、log_level、log_keep_days已弃用jobstokio 运行时工作线程数默认 0表示核心数×2。准备订阅数据数据库、超级表与 Topic无论目标平台是 MQTT 还是 Kafka发布的数据都来源于 TMQTDengine 消息队列订阅的 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 数据发布MQTTMessage Queuing Telemetry Transport是基于发布/订阅模型的轻量级消息协议广泛用于 IoT 设备间通信。TDengine 支持将实时数据推送到 MQTT 服务器从而便捷地将传感器数据、设备状态等分发到各类终端。安装并配置 MQTT 服务器使用该功能前需先部署 MQTT 服务器如 Mosquitto、EMQX 等可根据需要选择合适的实现。创建 MQTT 数据发布任务数据发布通过 taosx 命令行完成将 TMQ Topic 数据发布到 MQTTtaosx run -f tmqws://username:passwordip:port/topic?paramvalue... -t mqtt://ip:port?paramvalue...其中-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 订阅参数必填订阅所属消费组 IDclient.idTMQ 订阅参数可选订阅客户端 IDauto.offset.reset订阅起始位置experimental.snapshot.enable是否同步已落盘到 TSDB 存储文件非 WAL中的数据关闭时仅同步仍在 WAL 中的数据更多 TMQ 订阅参数参见 数据订阅文档。MQTT DSN 参数参数说明ip/portMQTT broker 地址与端口versionMQTT 协议版本必填可选3.1/3.1.1/5.0qosMQTT QoS 级别默认0client_idMQTT 客户端 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 tmqws://root:taosdatalocalhost:6041/topic_meters?group.idtaosx-pub-testauto.offset.resetearliest \ -t mqtt://mqtt.tdengine.com:1883?topictest/topic_metersqos1version5.0client_idtaosx-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等关键订阅参数使任务行为可预测。常用格式tmqws://root:taosdatalocalhost: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。示例tmqws://root:taosdatalocalhost:6041/topic_meters?group.idpub-kafka-demoauto.offset.resetearliestwith.metatrueStart 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 默认配置取值范围 11024应根据下游消费并发度与预期吞吐规划Replication Factor自动建 Topic 时的副本因子默认使用 broker 默认配置取值范围 1128且不能超过 Kafka 集群可用 broker 数量应与集群高可用策略一致。高级选项参数默认值说明Parallelism1Kafka 生产者最大并发度取值范围 1128高吞吐场景可从 1 或 2 起步逐步增大Queue Timeout (ms)30000消息进入发送队列后的最长等待时间网络不稳定时可增大追求快速失败时可减小Batch Size1000每个 Kafka 批次的最大记录数取值范围 1100000吞吐优先时增大、低延迟时调小Batch Timeout (ms)1000批次发送前的最大等待时间吞吐优先时增大、低延迟时减小Kafka Extra Parameters空额外 Kafka 生产者参数仅配置明确理解的原生参数并避免与标准字段冲突Kafka Extra Parameters 示例compression.typezstd acksall linger.ms100常见用途包括启用消息压缩、设置acks、调整linger.ms等发送策略。保存前的校验规则提交任务前通常会执行以下逻辑校验Topic DSN必须填写且必须显式选择Start FromEnable 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范围 1100、Wait time seconds默认30范围 1300。预览结果通常展示topic、key、value三个字段若等待时间内未收到数据系统会提示当前条件下无可预览消息。验证发布结果任务保存并启动后可用 Kafka 自带工具或第三方客户端验证消息发布是否正确。例如使用kcat消费目标 Topickcat -b 127.0.0.1:9092 -t taosx.data.out -C配置了 key 模板时消费输出通常同时包含消息 key 与消息体可结合源数据库、表名、Topic 模板与偏移字段共同验证发布结果。Flink 集成TDengine Flink ConnectorApache Flink 是 Apache 软件基金会支持的开源分布式流批一体处理框架可用于流处理、批处理、复杂事件处理、实时数仓构建并为机器学习提供实时数据支撑。借助 TDengine 的 Flink ConnectorFlink 可与 TDengine 无缝集成高效稳定地从数据库读取海量数据并进行分析处理。版本说明Flink Connector 相关能力仅存在于 TDengine TSDB Enterprise社区版集群可作为被读取的数据源但 Connector 完整能力以企业版为准。Flink 需 v1.19.0 及以上taosAdapter 需正常运行。引入依赖Maven 项目在pom.xml中添加dependency groupIdcom.taosdata.flink/groupId artifactIdflink-connector-tdengine/artifactId version2.1.4/version /dependencyConnector 版本历史如 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登录用户名默认rootpassword登录密码默认taosdatadatabase_name数据库名timezone时区HttpConnectTimeout连接超时时间毫秒默认 60000MessageWaitTimeout消息超时时间毫秒默认 60000UseSSL连接是否使用 SSLSource并行读取与三种分片方式Source 从 TDengine 读取数据并转换为 Flink 可处理的格式通过设置数据源并行度可多线程并行读取提升读取效率与吞吐。Source 关键属性TDengineConfigParams包括PROPERTY_KEY_USER/PROPERTY_KEY_PASSWORD用户名/密码默认root/taosdataVALUE_DESERIALIZER结果集反序列化方式收到RowData类型时设为RowData也可继承TDengineRecordDeserialization实现convert与getProducedType自定义反序列化TD_BATCH_MODE是否批量推送数据到下游算子为 True 时需将数据类型指定为SourceRecords的模板形式PROPERTY_KEY_MESSAGE_WAIT_TIMEOUT消息超时毫秒默认 60000PROPERTY_KEY_ENABLE_COMPRESSION传输过程是否启用压缩默认 falsePROPERTY_KEY_ENABLE_AUTO_RECONNECT是否启用自动重连默认 falsePROPERTY_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为例通过TDengineSourceRowData创建数据源并注册到 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?userrootpasswordtaosdata); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); TDengineSourceRowData source new TDengineSource(connProps, sql, RowData.class); DataStreamSourceRowData input env.fromSource(source, WatermarkStrategy.noWatermarks(), tdengine-source);批模式TD_BATCH_MODEtrue类型为SourceRecords模板形式与自定义类型ResultBean 继承TDengineRecordDeserialization的ResultSourceDeserialization的完整示例见 source/Main.java 中的source_batch_test与source_custom_type_test片段。CDC 数据订阅Flink CDC 提供数据订阅功能可实时监听 TDengine 数据变化并以数据流形式传给 Flink 处理同时保证数据一致性与完整性。CDC 关键参数TDengineCdcParams包括BOOTSTRAP_SERVERSTDengine 服务器ip:portWebSocket 连接时为 taosAdapter 所在地址CONNECT_USER/CONNECT_PASS用户名/密码默认root/taosdataPOLL_INTERVAL_MS拉取数据间隔默认 500msVALUE_DESERIALIZER结果集反序列化方式可继承com.taosdata.jdbc.tmq.ReferenceDeserializer指定结果集 Bean或继承com.taosdata.jdbc.tmq.Deserializer自定义TMQ_BATCH_MODE批量推送模式为 True 时类型需指定为ConsumerRecords模板形式GROUP_ID消费组 ID同组共享消费进度最大长度 192AUTO_OFFSET_RESET消费组订阅起始位置earliest从最早订阅、latest从最新订阅默认latestENABLE_AUTO_COMMIT是否自动提交消费点位true 自动提交、false 按 checkpoint 提交默认 false。注意自动提交模式在获取数据后即提交无论下游算子是否正确处理存在数据丢失风险主要用于无状态算子或一致性要求低的场景AUTO_COMMIT_INTERVAL_MS自动提交消费记录的时间间隔毫秒默认 5000仅ENABLE_AUTO_COMMITtrue时生效TMQ_SESSION_TIMEOUT_MS消费者心跳丢失后的超时时间触发 rebalance 后剔除该消费者3.3.3.0 起支持默认 12000范围 [6000, 1800000]TMQ_MAX_POLL_INTERVAL_MS消费者 poll 拉取的最长间隔超时视为消费者离线并触发 rebalance3.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参数名类型说明connectorstring连接器标识固定为tdengine-connectortd.jdbc.urlstring连接 URLtd.jdbc.modestring连接器类型source、sinktable.namestring源或目标表名scan.querystring取数 SQL 语句sink.db.namestring目标数据库名sink.supertable.namestring目标超级表名sink.batch.sizeinteger批量写入大小sink.table.namestring子表或普通表名Table CDC 连接器参数在此基础上增加user、password、bootstrap.servers、topic、td.jdbc.modecdc、sink、group.id、auto.offset.resetearliest/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?userrootpasswordtaosdata, 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 DSNTDengine 连接地址例如taosws://root:taosdatalocalhost:6041/dbSQL Query一条只读SELECT查询支持WITH ... SELECT查询Output FiletaosX 服务节点上的输出文件路径必须以.parquet结尾Overwrite Existing File新的临时文件成功关闭后是否允许替换已有的最终文件Compressionuncompressed、zstd、snappy、gzip、brotli或lz4_rawCompression Levelzstd、gzip、brotli可选Row Group Size每个 Parquet row group 的最大行数默认 131072。Parquet 写入器仅在 row group 写满或文件关闭时才落盘取值越小.part文件增长越频繁但会降低压缩率并增加文件元数据开销。有效范围 102410000000。输出路径规则输出文件创建在 taosX 服务节点上不是浏览器本地路径相对路径写入$DATA_DIR/tasks/task_id/job_id/目录下绝对路径按 taosX 服务节点上的绝对路径处理写入器先在相同目录创建final_name.part临时文件关闭 Parquet 写入器并校验元数据后再将临时文件重命名为最终路径第一版不支持 agent 或 via 执行。DSN 示例FROM taosws://root:taosdatalocalhost:6041/db?queryselect%20*%20from%20meters TO parquet:/tmp/meters.parquet?overwritefalsecompressionzstdrow_group_size131072对应命令行方式可参考 taosX 组件文档 中的 SQL 查询结果导出用法例如taosx run -f taosws://root:taosdatalocalhost:6041/test?queryselect * 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),仅供参考
返回列表