资讯动态

TDengine 企业版 Flink 连接器实战:Source/CDC 数据订阅与 Sink 写入完整指南

发布时间:2026/9/13 14:54:56 来源:尧图企业网站定制
TDengine 企业版 Flink 连接器实战Source/CDC 数据订阅与 Sink 写入完整指南【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineApache Flink 是 Apache 软件基金会支持的开源分布式流批一体化处理框架广泛用于流处理、批处理、复杂事件处理、实时数据仓库构建及为机器学习提供实时数据支持等大数据场景。本文基于 TDengine 企业版 Flink 连接器系统讲解如何让 Flink 从 TDengine 读取数据Source/CDC并将处理结果写回 TDengineSink覆盖连接参数、数据分片、CDC 订阅、Table SQL 集成、错误码排障等完整链路读者可据此在工业物联网IIoT等实时数据场景中落地流批一体分析与实时链路集成。注意本功能仅适用于 TDengine 企业版TSDB Enterprise。环境准备安装与前置检查安装和配置 Flink 服务器安装Apache Flink v1.19.0 或以上版本详细安装方式请参考 Apache Flink 官方文档。Flink Connector 支持所有可运行 Flink 1.19 及以上版本的平台。确认企业版服务正常在启动 Flink 任务前需确认以下服务处于正常运行状态taosd 服务正常TDengine 数据库主服务负责数据存储与查询taosAdapter 服务正常负责将 REST/WebSocket 请求转发给 taosdFlink 连接器通过jdbc:TAOS-WS://协议经 taosAdapter 访问 TDengine。引入 Maven 依赖如果使用 Maven 管理项目只需在pom.xml中加入以下依赖以 2.1.4 版本为例dependency groupIdcom.taosdata.flink/groupId artifactIdflink-connector-tdengine/artifactId version2.1.4/version /dependency连接器版本与 TDengine 企业版版本对应关系、各版本主要变化参见 版本历史其中 2.0.0 起支持自定义数据结构序列化写入与 Table SQL 写入要求 TSDB-Enterprisev3.3.5.1及以上1.0.0 起支持 Sink 功能要求v3.3.2.0及以上。Flink 语义选择说明连接器采用At-Least-Once至少一次语义原因如下TDengine 目前不支持事务不能进行频繁的检查点操作和复杂的事务协调TDengine 采用时间戳作为主键重复数据下游算子可以进行过滤操作避免重复计算采用 At-Least-Once 可确保较高的数据处理性能和较低的数据延迟。设置方式StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);建立连接URL 与 Properties 参数建立连接的参数有URL和Properties两类。URL 规范格式为jdbc:TAOS-WS://[host_name]:[port]/[database_name]?[user{user}|password{password}|timezone{timezone}]URL 参数说明参数说明默认值user登录 TDengine 用户名rootpassword用户登录密码taosdatadatabase_name数据库名称无timezone时区设置无httpConnectTimeout连接超时时间单位 ms60000messageWaitTimeout消息超时时间单位 ms60000useSSL连接中是否使用 SSL无数据准备通过命令行工具taos或管理界面 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);更完整的示例可参考仓库中的 Flink Source 示例其中使用CREATE DATABASE IF NOT EXISTS power vgroups 5创建数据库、CREATE STABLE IF NOT EXISTS meters (ts timestamp, current float, voltage int, phase float) TAGS (location binary(64), groupId int)创建超级表、CREATE TOPIC topic_meters as SELECT ts,current, voltage, phase, location, groupid, tbname FROM meters创建订阅主题并通过一条多值 INSERT 写入power.d1001、power.d1002两个子表的数据。注意prepare()中先执行DROP TOPIC IF EXISTS与DROP DATABASE IF EXISTS以保证示例可重复执行。Source从 TDengine 并行读取数据Source 拉取 TDengine 数据库中的数据并将其转换为 Flink 内部可处理的格式和类型以并行的方式进行读取和分发为后续数据处理提供高效输入。通过设置数据源的并行度env.setParallelism(n)可实现多个线程并行读取提高读取效率和吞吐量充分利用集群资源进行大规模数据处理。Source Properties 配置参数参数TDengineConfigParams 常量说明默认值PROPERTY_KEY_USER登录 TDengine 用户名rootPROPERTY_KEY_PASSWORD用户登录密码taosdataVALUE_DESERIALIZER下游算子接收结果集反序列化方法。若接收结果集类型是 Flink 的RowData仅需设置为RowData也可继承TDengineRecordDeserialization并实现convert和getProducedType方法根据 SQL 的ResultSet自定义反序列化方式无TD_BATCH_MODE是否批量将数据推送给下游算子。若为True创建TDengineSource对象时需指定数据类型为SourceRecords类型的泛型形式falsePROPERTY_KEY_MESSAGE_WAIT_TIMEOUT消息超时时间单位 ms60000PROPERTY_KEY_ENABLE_COMPRESSION传输过程是否启用压缩true启用false不启用falsePROPERTY_KEY_ENABLE_AUTO_RECONNECT是否启用自动重连true启用false不启用falsePROPERTY_KEY_RECONNECT_INTERVAL_MS自动重连重试间隔单位 ms。仅在PROPERTY_KEY_ENABLE_AUTO_RECONNECT为true时生效2000PROPERTY_KEY_RECONNECT_RETRY_COUNT自动重连重试次数。仅在PROPERTY_KEY_ENABLE_AUTO_RECONNECT为true时生效3PROPERTY_KEY_DISABLE_SSL_CERT_VALIDATION关闭 SSL 证书验证true启用false不启用false按时间分片用户可对查询的 SQL 按照时间拆分为多个子任务输入开始时间、结束时间、拆分间隔、时间字段名称系统会按照设置的间隔时间左闭右开进行拆分并行获取数据。使用SourceSplitSql、SplitType.SPLIT_TYPE_TIMESTAMP与TimestampSplitInfo完成配置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)));上述配置表示将2024-12-19 16:12:48.000至19:12:48.000之间的查询按 1 小时间隔拆分为多个子任务并行执行。完整上下文见 Source 示例 time_interval 片段。按超级表 TAG 分片用户可按照超级表的 TAG 字段将查询的 SQL 拆分为多个查询条件系统会以一个查询条件对应一个子任务的方式进行拆分进而并行获取数据。使用SplitType.SPLIT_TYPE_TAG并传入条件列表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);该示例将原查询按groupid与location的不同组合拆成 3 个并行子查询适合按地域、设备分组等维度并行拉取。完整上下文见 Source 示例 tag_split 片段。按表名分片支持输入多个相同表结构的超级表或普通表进行分片系统会按照一个表一个任务的方式进行拆分进而并行获取数据。使用SplitType.SPLIT_TYPE_TABLE并传入表名列表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 示例 table_split 片段。使用 Source 连接器查询结果为 RowData 数据类型示例RowData SourceProperties 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, splitSql, RowData.class); DataStreamSourceRowData input env.fromSource(source, WatermarkStrategy.noWatermarks(), tdengine-source); DataStreamString resultStream input.map((MapFunctionRowData, String) rowData - { StringBuilder sb new StringBuilder(); sb.append(ts: rowData.getTimestamp(0, 0) , current: rowData.getFloat(1) , voltage: rowData.getInt(2) , phase: rowData.getFloat(3) , location: rowData.getString(4).toString()); sb.append(\n); return sb.toString(); }); resultStream.print(); env.execute(tdengine flink source);完整示例见 Source 示例 source_test 片段。批量查询结果示例Batch Source设置TD_BATCH_MODEtrue并将TDengineSource的泛型指定为SourceRecordsRowData下游通过迭代器逐条消费connProps.setProperty(TDengineConfigParams.TD_BATCH_MODE, true); ClassSourceRecordsRowData typeClass (ClassSourceRecordsRowData) (Class?) SourceRecords.class; TDengineSourceSourceRecordsRowData source new TDengineSource(connProps, sql, typeClass);完整示例见 Source 示例 source_batch_test 片段。查询结果为自定义数据类型示例Custom Type Source将VALUE_DESERIALIZER设置为自定义反序列化类的全限定名泛型指定为自定义 BeanconnProps.setProperty(TDengineConfigParams.VALUE_DESERIALIZER, com.taosdata.flink.entity.ResultSourceDeserialization); ... TDengineSourceResultBean source new TDengineSource(connProps, splitSql, ResultBean.class);完整示例见 Source 示例 source_custom_type_test 片段。其中ResultBean是自定义的一个内部类用于定义 Source 查询结果的数据类型ResultSourceDeserialization是自定义的一个内部类通过继承TDengineRecordDeserialization并实现convert和getProducedType方法实现按需的反序列化逻辑。CDC 数据订阅实时监控数据变更Flink CDC 主要用于提供数据订阅功能能实时监控 TDengine 数据库的数据变化并将这些变更以数据流形式传输到 Flink 中进行处理同时确保数据的一致性和完整性。CDC 连接器会根据用户设置的并行度创建 consumer因此请根据资源情况合理设置并行度。CDC Properties 配置参数参数说明默认值TDengineCdcParams.BOOTSTRAP_SERVERSTDengine 服务端所在的ip:port若使用 WebSocket 连接则为 taosAdapter 所在的ip:port无TDengineCdcParams.CONNECT_USER登录 TDengine 用户名rootTDengineCdcParams.CONNECT_PASS用户登录密码taosdataTDengineCdcParams.POLL_INTERVAL_MS拉取数据间隔500msTDengineCdcParams.VALUE_DESERIALIZER结果集反序列化方法。若接收结果集类型是 Flink 的RowData仅需设置为RowData也可继承com.taosdata.jdbc.tmq.ReferenceDeserializer并指定结果集 bean 实现反序列化无TDengineCdcParams.TMQ_BATCH_MODE是否批量将数据推送给下游算子。若为True创建TDengineCdcSource对象时需指定数据类型为ConsumerRecords类型的泛型形式falseTDengineCdcParams.GROUP_ID消费组 ID同一消费组共享消费进度。最大长度192无TDengineCdcParams.AUTO_OFFSET_RESET消费组订阅的初始位置earliest从头开始订阅latest仅从最新数据开始订阅latestTDengineCdcParams.ENABLE_AUTO_COMMIT是否启用消费位点自动提交true自动提交false依赖 checkpoint 时间来提交falseTDengineCdcParams.AUTO_COMMIT_INTERVAL_MS消费记录自动提交消费位点的时间间隔单位 ms。仅在ENABLE_AUTO_COMMIT为true时生效5000TDengineConfigParams.PROPERTY_KEY_ENABLE_COMPRESSION传输过程是否启用压缩true启用false不启用falseTDengineConfigParams.PROPERTY_KEY_ENABLE_AUTO_RECONNECT是否启用自动重连falseTDengineConfigParams.PROPERTY_KEY_RECONNECT_INTERVAL_MS自动重连重试间隔单位 ms仅重连开启时生效2000TDengineConfigParams.PROPERTY_KEY_RECONNECT_RETRY_COUNT自动重连重试次数仅重连开启时生效3TDengineCdcParams.TMQ_SESSION_TIMEOUT_MSconsumer 心跳丢失后的超时时间超时后触发 rebalance 逻辑成功后该 consumer 会被删除从v3.3.3.0开始支持12000取值范围[6000, 1800000]TDengineCdcParams.TMQ_MAX_POLL_INTERVAL_MSconsumer poll 拉取数据间隔的最长时间超过该时间认为该 consumer 离线触发 rebalance 逻辑成功后该 consumer 会被删除300000取值范围[1000, INT32_MAX]注意自动提交模式下reader 获取完成数据后自动提交不管下游算子是否正确处理了数据存在数据丢失的风险主要用于追求高效的无状态算子场景或是数据一致性要求不高的场景。生产环境建议保持ENABLE_AUTO_COMMITfalse配合 Flink Checkpoint 机制提交消费位点。使用 CDC 连接器订阅结果为 RowData 数据类型示例CDC SourceStreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); env.enableCheckpointing(100, AT_LEAST_ONCE); env.getConfig().setRestartStrategy(RestartStrategies.noRestart()); Properties config new Properties(); config.setProperty(TDengineCdcParams.CONNECT_TYPE, ws); config.setProperty(TDengineCdcParams.BOOTSTRAP_SERVERS, localhost:6041); config.setProperty(TDengineCdcParams.AUTO_OFFSET_RESET, earliest); config.setProperty(TDengineCdcParams.MSG_WITH_TABLE_NAME, true); config.setProperty(TDengineCdcParams.AUTO_COMMIT_INTERVAL_MS, 1000); config.setProperty(TDengineCdcParams.GROUP_ID, group_1); config.setProperty(TDengineCdcParams.ENABLE_AUTO_COMMIT, true); config.setProperty(TDengineCdcParams.CONNECT_USER, root); config.setProperty(TDengineCdcParams.CONNECT_PASS, taosdata); config.setProperty(TDengineCdcParams.VALUE_DESERIALIZER, RowData); config.setProperty(TDengineCdcParams.VALUE_DESERIALIZER_ENCODING, UTF-8); TDengineCdcSourceRowData tdengineSource new TDengineCdcSource(topic_meters, config, RowData.class); DataStreamSourceRowData input env.fromSource(tdengineSource, WatermarkStrategy.noWatermarks(), tdengine-source);完整示例见 Source 示例 cdc_source 片段示例通过env.executeAsync(...)异步提交任务并在 5 秒后取消便于本地验证通过 Flink UI 提交的任务无法直接 cancel需要在 UI 页面停止。将订阅结果批量下发到算子的示例CDC Batch Source设置TMQ_BATCH_MODEtrue并将TDengineCdcSource的泛型指定为ConsumerRecordsRowData通过records.iterator()批量消费config.setProperty(TDengineCdcParams.TMQ_BATCH_MODE, true); ClassConsumerRecordsRowData typeClass (ClassConsumerRecordsRowData) (Class?) ConsumerRecords.class; TDengineCdcSourceConsumerRecordsRowData tdengineSource new TDengineCdcSource(topic_meters, config, typeClass);完整示例见 Source 示例 cdc_batch_source 片段。订阅结果为自定义数据类型示例CDC Custom Type将VALUE_DESERIALIZER设置为自定义反序列化类全限定名例如com.taosdata.flink.entity.ResultDeserializer泛型指定为ResultBeanconfig.setProperty(TDengineCdcParams.VALUE_DESERIALIZER, com.taosdata.flink.entity.ResultDeserializer); config.setProperty(TDengineCdcParams.VALUE_DESERIALIZER_ENCODING, UTF-8); TDengineCdcSourceResultBean tdengineSource new TDengineCdcSource(topic_meters, config, ResultBean.class);完整示例见 Source 示例 cdc_custom_type_test 片段。ResultBean是自定义的一个内部类其字段名和数据类型与列的名称和数据类型一一对应这样根据TDengineCdcParams.VALUE_DESERIALIZER属性对应的反序列化类即可反序列化出ResultBean类型的对象。Sink将处理结果写回 TDengineFlink 任务处理完的数据可通过TDengineSink批量写回 TDengine 超级表或普通表。以下示例构造 10 条GenericRowData字段顺序ts、current、voltage、phase、location、groupid、tbname并写入power_sink库的sink_meters超级表StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); RowData[] rows new GenericRowData[10]; Random random new Random(System.currentTimeMillis()); for (int i 0; i 10; i) { GenericRowData row new GenericRowData(7); long current System.currentTimeMillis() i * 1000; row.setField(0, TimestampData.fromEpochMillis(current)); // ts row.setField(1, random.nextFloat() * 30); // current row.setField(2, 300 (i 1)); // voltage row.setField(3, random.nextFloat()); // phase row.setField(4, StringData.fromString(location_ i)); // location row.setField(5, i); // groupid row.setField(6, StringData.fromString(d0 i)); // tbname rows[i] row; } DataStreamRowData dataStream env.fromElements(RowData.class, rows); Properties sinkProps new Properties(); sinkProps.setProperty(TSDBDriver.PROPERTY_KEY_ENABLE_AUTO_RECONNECT, true); sinkProps.setProperty(TSDBDriver.PROPERTY_KEY_CHARSET, UTF-8); sinkProps.setProperty(TSDBDriver.PROPERTY_KEY_TIME_ZONE, UTC-8); sinkProps.setProperty(TDengineConfigParams.VALUE_DESERIALIZER, RowData); sinkProps.setProperty(TDengineConfigParams.PROPERTY_KEY_DBNAME, power_sink); sinkProps.setProperty(TDengineConfigParams.TD_SUPERTABLE_NAME, sink_meters); sinkProps.setProperty(TDengineConfigParams.TD_JDBC_URL, jdbc:TAOS-WS://localhost:6041/power_sink?userrootpasswordtaosdata); sinkProps.setProperty(TDengineConfigParams.TD_BATCH_SIZE, 2000); TDengineSinkRowData sink new TDengineSink(sinkProps, Arrays.asList(ts, current, voltage, phase, location, groupid, tbname)); dataStream.sinkTo(sink); env.execute(flink tdengine sink);完整示例见 Sink 示例 RowDataToSuperTable 片段。关键点说明TDengineSink构造函数的第二个参数为目标表字段名列表Arrays.asList(...)需与数据顺序保持一致写入超级表时配置PROPERTY_KEY_DBNAME或TD_DATABASE_NAME与TD_SUPERTABLE_NAME数据行中需包含tbname字段用于指定子表名写入普通表时配置TD_TABLE_NAME见 Sink 示例 RowDataToNormalTable 片段Sink 亦支持自定义数据结构序列化VALUE_DESERIALIZER指向自定义序列化类见 Sink 示例 CustomTypeToNormalTable 片段。Table SQL以声明式 SQL 完成跨源集成使用 Table SQL 的方式可以从多个不同的数据源数据库如 TDengine、MySQL、Oracle 等中提取数据再进行自定义的算子操作如数据清洗、格式转换、关联不同表的数据等然后将处理后的结果加载到目标数据源如 TDengine、MySQL 等中。Table Source 连接器参数配置说明参数名称类型参数说明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写入的普通表或子表名称使用示例将power库的meters表的子表数据写入power_sink库的sink_meters超级表对应的子表中StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(3); env.enableCheckpointing(1000, CheckpointingMode.AT_LEAST_ONCE); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env, fsSettings); String tdengineSourceTableDDL CREATE TABLE meters ( ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT, location VARCHAR(255), groupid INT, tbname VARCHAR(255) ) 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 ); String tdengineSinkTableDDL CREATE TABLE sink_meters ( ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT, location VARCHAR(255), groupid INT, tbname VARCHAR(255) ) WITH ( connector tdengine-connector, td.jdbc.mode sink, td.jdbc.url jdbc:TAOS-WS://localhost:6041/power_sink?userrootpasswordtaosdata, sink.db.name power_sink, sink.supertable.name sink_meters ); tableEnv.executeSql(tdengineSourceTableDDL); tableEnv.executeSql(tdengineSinkTableDDL); tableEnv.executeSql(INSERT INTO sink_meters SELECT ts, current, voltage, phase, location, groupid, tbname FROM meters);完整示例见 Source 示例 source_table 片段。示例中通过CREATE TABLE ... WITH (...)声明 source 表td.jdbc.modesourcescan.query与 sink 表td.jdbc.modesinksink.db.name/sink.supertable.name再以一条INSERT INTO ... SELECT ...完成数据搬移。sink 表 DDL 的更多用法普通表场景使用sink.table.name、批模式inBatchMode()下用Row构造数据后executeInsert可参考 Sink 示例。Table CDC 连接器参数配置说明参数名称类型参数说明connectorstring连接器标识设置tdengine-connectoruserstring用户名默认rootpasswordstring密码默认taosdatabootstrap.serversstring服务器地址topicstring订阅主题td.jdbc.modestring连接器类型cdc、sinkgroup.idstring消费组 ID同一消费组共享消费进度auto.offset.resetstring消费组订阅的初始位置earliest从头开始订阅latest仅从最新数据开始订阅默认latestpoll.interval_msinteger拉取数据间隔默认500mssink.db.namestring目标数据库名称sink.supertable.namestring写入的超级表名称sink.batch.sizeinteger写入的批大小sink.table.namestring写入的普通表或子表名称使用示例订阅power库的meters超级表的子表数据写入power_sink库的sink_meters超级表对应的子表中StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(5); env.enableCheckpointing(1000, CheckpointingMode.AT_LEAST_ONCE); StreamTableEnvironment tableEnv StreamTableEnvironment.create(env, fsSettings); String tdengineSourceTableDDL CREATE TABLE meters ( ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT, location VARCHAR(255), groupid INT, tbname VARCHAR(255) ) WITH ( connector tdengine-connector, bootstrap.servers localhost:6041, td.jdbc.mode cdc, group.id group_22, auto.offset.reset earliest, enable.auto.commit false, topic topic_meters ); // sink 表 DDL 与 Table Source 示例中的 sink_meters 声明一致 tableEnv.executeSql(tdengineSourceTableDDL); tableEnv.executeSql(tdengineSinkTableDDL); TableResult tableResult tableEnv.executeSql(INSERT INTO sink_meters SELECT ts, current, voltage, phase, location, groupid, tbname FROM meters);完整示例见 Source 示例 cdc_table 片段。CDC 表连接器通过bootstrap.servers指向 taosAdapter 地址以td.jdbc.modecdc订阅topic指定的主题数据变更实时进入 Flink 后写回 sink 表实现实时数据管道。数据类型映射TDengine 目前支持时间戳、数字、字符、布尔类型与 Flink RowData Type 对应类型转换关系如下TDengine DataTypeFlink RowDataTypeTIMESTAMPTimestampDataINTIntegerBIGINTLongFLOATFloatDOUBLEDoubleSMALLINTShortTINYINTByteBOOLBooleanVARCHARStringDataBINARYStringDataNCHARStringDataJSONStringDataVARBINARYbyte[]GEOMETRYbyte[]异常与错误码排查任务执行失败后首先查看 Flink 任务执行日志确认失败原因。连接器常见错误码及处理建议如下Error CodeDescriptionSuggested Actions0xa000connection param error连接器参数错误。0xa010database name configuration error数据库名配置错误。0xa011table name configuration error表名配置错误。0xa013value.deserializer parameter not set未设置序列化方式。0xa014list of column names for target table not set未设置目标表的列名列表。0x2301connection already closed连接已经关闭检查连接情况或重新创建连接去执行相关指令。0x2302this operation is NOT supported currently!当前使用接口不支持可以更换其他连接方式。0x2303invalid variables参数不合法请检查相应接口规范调整参数类型及大小。0x2304statement is closedstatement 已经关闭请检查 statement 是否关闭后再次使用或是连接是否正常。0x2305resultSet is closedresultSet 结果集已经释放请检查 resultSet 是否释放后再次使用。0x230dparameter index out of range参数越界请检查参数的合理范围。0x230econnection already closed连接已经关闭请检查 Connection 是否关闭后再次使用或是连接是否正常。0x230funknown sql type in TDengine请检查 TDengine 支持的 Data Type 类型。0x2315unknown taos type in TDengine在 TDengine 数据类型与 JDBC 数据类型转换时是否指定了正确的 TDengine 数据类型。0x2319user is required创建连接时缺少用户名信息。0x231apassword is required创建连接时缺少密码信息。0x231dcant create connection with server within通过增加参数httpConnectTimeout增加连接耗时或检查与 taosAdapter 之间的连接情况。0x231efailed to complete the task within the specified time通过增加参数messageWaitTimeout增加执行耗时或检查与 taosAdapter 之间的连接情况。0x2352Unsupported encoding本地连接下指定了不支持的字符编码集。0x2353internal error of database, please see taoslog for more details本地连接执行 prepareStatement 时出现错误请检查 taos log 进行问题定位。0x2354connection is NULL本地连接执行命令时 Connection 已经关闭请检查与 TDengine 的连接情况。0x2355result set is NULL本地连接获取结果集时结果集异常请检查连接情况并重试。0x2356invalid num of fields本地连接获取结果集的 meta 信息不匹配。小结TDengine 企业版 Flink 连接器为 Flink 生态提供了完整的数据通路Source支持按时间、按超级表 TAG、按表名三种分片策略并行读取CDC支持基于主题订阅的实时数据变更捕获配合消费组与位点提交机制保证一致性Sink支持 RowData、批量及自定义类型的写入Table SQL则让跨源数据集成以声明式 SQL 完成。任务失败时可按错误码快速定位参数、连接或类型问题。相关可运行示例均位于仓库 docs/examples/flink/source/Main.java 与 docs/examples/flink/sink/Main.java可结合本文逐步验证。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价