拓冰建站拓冰建站
首页 / 资讯中心 / 正文

TDengine Flink Connector 使用指南:从 Sink 写入到 Table Sink 的流批集成实战

TDengine Flink Connector 使用指南从 Sink 写入到 Table Sink 的流批集成实战【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengineApache Flink 是 Apache 软件基金会支持的开源分布式流批一体处理框架广泛用于流处理、批处理、复杂事件处理与实时数仓构建。本文基于 TDengine 官方文档与仓库中的可运行示例系统讲解如何借助flink-connector-tdengine连接器将 Flink 作业中的处理结果写入 TDengineSink并通过 Flink Table API 以 SQL 声明式方式完成写入Table Sink。读完本文你将掌握连接参数的配置方法、三种 Sink 写入模式、数据类型的映射规则、异常排查手段以及 At-Least-Once 语义的正确选择。前置条件在开始集成前需要准备以下运行环境TDengine 服务已部署并正常运行企业版与社区版均可。taosAdapter 能够正常运行连接器通过 WebSocket 方式经由 taosAdapter 访问数据库。Apache Flink 1.19.0 或以上版本已安装安装方式请参考 Apache Flink 官方文档。支持的平台Flink Connector 支持所有能够运行 Flink 1.19 及以上版本的平台。由于连接器基于jdbc:TAOS-WS://WebSocket协议工作WebSocket 连接方式不依赖 TDengine 原生客户端驱动天然具备跨平台能力。连接器版本演进连接器持续迭代各版本的主要变更如下完整版本历史可参考 Java 连接器版本历史其中 2.1.4 版本将 JDBC 驱动升级至 3.7.3Flink Connector 版本主要变更对应 TDengine TSDB 企业版2.1.4将 JDBC 驱动升级至 3.7.3-2.1.3增加数据转换时的异常信息输出-2.1.2增加对写入字段的反引号过滤-2.1.1修复 Stmt 中同一张表数据绑定失败的问题-2.1.0修复来自不同数据源的 varchar 类型写入问题-2.0.2Table Sink 支持 RowKind.UPDATE_BEFORE、RowKind.UPDATE_AFTER、RowKind.DELETE 等类型-2.0.1Sink 支持写入 RowData 实现类型-2.0.01. Sink 支持自定义数据结构序列化后写入 TDengine2. 支持使用 Table SQL 写入 TDengine3.3.5.1 及以上1.0.0支持 Sink 功能将其他数据源的数据写入 TDengine3.3.2.0 及以上连接参数建立连接的参数由 URL 和 Properties 两部分组成。URL 的规范格式为jdbc:TAOS-WS://[host_name]:[port]/[database_name]?[user{user}|password{password}|timezone{timezone}]参数说明参数说明默认值user登录 TDengine 的用户名rootpassword用户登录密码taosdatadatabase_name数据库名称-timezone时区设置-httpConnectTimeout连接超时时间单位毫秒60000messageWaitTimeout消息超时时间单位毫秒60000useSSL连接中是否使用 SSL-需要注意的是TDengine 的 Java 原生连接与 REST 连接已被标记为弃用将于 2027-01-01 停止连接器统一推荐使用 WebSocket 连接方式即jdbc:TAOS-WS://协议前缀并使用com.taosdata.jdbc.ws.WebSocketDriver驱动类。仓库中的示例代码即遵循该方式例如static String jdbcUrl jdbc:TAOS-WS://localhost:6041?userrootpasswordtaosdata;Sink将 Flink 处理结果写入 TDengineSink 的核心功能是将 Flink 作业中来自不同数据源或算子处理后的数据高效、准确地写入 TDengine其高效写入机制保证了数据的快速稳定落库。:::note写入的目标数据库必须已经创建。写入的超级表/普通表必须已经创建。 :::Sink Properties 配置说明Sink 通过Properties传递配置常用参数如下TDengineConfigParams.PROPERTY_KEY_USER登录 TDengine 用户名默认值root。TDengineConfigParams.PROPERTY_KEY_PASSWORD用户登录密码默认值taosdata。TDengineConfigParams.PROPERTY_KEY_DBNAME写入的数据库名称。TDengineConfigParams.TD_SUPERTABLE_NAME写入的超级表名称。写入的数据必须带有tbname字段用于确定写入哪张子表。TDengineConfigParams.TD_TABLE_NAME写入子表或普通表的表名此参数与TD_SUPERTABLE_NAME仅需设置一个。TDengineConfigParams.VALUE_DESERIALIZER接收结果集的反序列化方法。若接收的结果集类型是 Flink 的RowData设置为RowData即可也可以继承TDengineSinkRecordSerializer并实现serialize方法根据接收的数据类型自定义序列化方式。TDengineConfigParams.TD_BATCH_SIZE设置一次写入 TDengine 数据库的批大小。当达到批数量后触发写入或在一个 checkpoint 到达时也会触发写入。TDengineConfigParams.PROPERTY_KEY_MESSAGE_WAIT_TIMEOUT消息超时时间单位毫秒默认值 60000。TDengineConfigParams.PROPERTY_KEY_ENABLE_COMPRESSION传输过程是否启用压缩。true启用false不启用默认false。TDengineConfigParams.PROPERTY_KEY_ENABLE_AUTO_RECONNECT是否启用自动重连。true启用false不启用默认false。TDengineConfigParams.PROPERTY_KEY_RECONNECT_INTERVAL_MS自动重连重试间隔单位毫秒默认值 2000仅在PROPERTY_KEY_ENABLE_AUTO_RECONNECT为true时生效。TDengineConfigParams.PROPERTY_KEY_RECONNECT_RETRY_COUNT自动重连重试次数默认值 3仅在PROPERTY_KEY_ENABLE_AUTO_RECONNECT为true时生效。TDengineConfigParams.PROPERTY_KEY_DISABLE_SSL_CERT_VALIDATION关闭 SSL 证书验证。true启用false不启用默认false。示例 1将 RowData 类型数据写入超级表对应的子表以下示例将RowData类型的数据写入power_sink库中sink_meters超级表对应的子表。完整代码见 docs/examples/flink/sink/Main.javastatic void testRowDataToSuperTable() throws Exception { 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); }写入超级表时构造TDengineSink的字段列表必须包含tbname列即Arrays.asList(...)中最后一个元素连接器根据该字段的值决定将数据写入哪张子表。示例 2将 RowData 类型数据写入普通表写入普通表时不再需要tbname字段只需配置TD_TABLE_NAME指定目标表名static void testRowDataToNormalTable() throws Exception { 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(4); 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 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_TABLE_NAME, sink_normal); 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)); dataStream.sinkTo(sink); env.execute(flink tdengine sink); }示例 3将自定义类型数据写入超级表对应的子表当上游数据不是 Flink 的RowData而是自定义 POJO 时可通过继承TDengineSinkRecordSerializer并实现serialize方法来自定义序列化逻辑然后在VALUE_DESERIALIZER参数中指定该序列化类全名static void testCustomTypeToSink() throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); ResultBean[] rows new ResultBean[10]; Random random new Random(System.currentTimeMillis()); for (int i 0; i 10; i) { ResultBean rowData new ResultBean(); long current System.currentTimeMillis() i * 1000; rowData.setTs(new Timestamp(current)); rowData.setCurrent(random.nextFloat() * 30); rowData.setVoltage(300 (i 1)); rowData.setPhase(random.nextFloat()); rowData.setLocation(location_ i); rowData.setGroupid(i); rowData.setTbname(d0 i); rows[i] rowData; } DataStreamResultBean dataStream env.fromElements(ResultBean.class, rows); Properties sinkProps new Properties(); sinkProps.setProperty(TDengineConfigParams.PROPERTY_KEY_ENABLE_AUTO_RECONNECT, true); sinkProps.setProperty(TDengineConfigParams.PROPERTY_KEY_CHARSET, UTF-8); sinkProps.setProperty(TSDBDriver.PROPERTY_KEY_TIME_ZONE, UTC-8); sinkProps.setProperty(TDengineConfigParams.VALUE_DESERIALIZER, com.taosdata.flink.entity.ResultBeanSinkSerializer); 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); TDengineSinkResultBean sink new TDengineSink(sinkProps, Arrays.asList(ts, current, voltage, phase, location, groupid, tbname)); dataStream.sinkTo(sink); env.execute(flink tdengine sink); }其中ResultBean是自定义类用于定义写入字段的数据类型ResultBeanSinkSerializer是自定义类通过继承 TDengine 的序列化基类并实现serialize方法完成自定义序列化。从源码理解 Sink 的写入机制从示例代码可以看到TDengineSink的构造参数除了Properties外还需要传入目标表字段名的有序列表该列表必须与数据的排列顺序保持一致示例注释中明确说明The list of target table field names needs to be consistent with the data order。Sink 内部会将 RowData/自定义对象按字段顺序转换为 TDengine 的写入语句并通过TD_BATCH_SIZE控制批量大小达到批大小或 checkpoint 触发时执行一次写入。写入前请在prepare()阶段完成库表准备参考 docs/examples/flink/sink/Main.javaCREATE DATABASE IF NOT EXISTS power_sink vgroups 5; CREATE STABLE IF NOT EXISTS sink_meters (ts timestamp, current float, voltage int, phase float) TAGS (location binary(64), groupId int); CREATE TABLE IF NOT EXISTS sink_normal (ts timestamp, current float, voltage int, phase float);Table Sink通过 SQL 声明式写入Table Sink 允许使用 Flink Table API 从多个不同的数据源如 MySQL、Oracle、Kafka 等中提取数据进行自定义算子操作数据清洗、格式转换、关联不同表的数据等后将处理结果写入 TDengine。参数配置说明通过 DDL 的WITH子句声明连接器参数参数名称类型参数说明connectorstring连接器标识设置为tdengine-connectortd.jdbc.urlstring连接的 URLtd.jdbc.modestring连接器类型设置为sinksink.db.namestring目标数据库名称sink.batch.sizeinteger写入的批大小sink.supertable.namestring写入的超级表名称sink.table.namestring写入的普通表或子表名称示例 1通过 Table SQL 写入超级表对应的子表static void testTableSqlToSink() throws Exception { EnvironmentSettings settings EnvironmentSettings.newInstance().inStreamingMode().build(); TableEnvironment tEnv TableEnvironment.create(settings); 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 ); tEnv.executeSql(tdengineSinkTableDDL); String insertQuery INSERT INTO sink_meters VALUES (CAST(2024-12-19 19:12:45 AS TIMESTAMP(6)), 50.30000, 201, 3.31003, California.SanFrancisco, 1, d1001), (CAST(2024-12-19 19:12:46 AS TIMESTAMP(6)), 82.60000, 202, 0.33000, California.SanFrancisco, 1, d1001), (CAST(2024-12-19 19:12:47 AS TIMESTAMP(6)), 92.30000, 203, 0.31000, California.SanFrancisco, 1, d1001), (CAST(2024-12-19 19:12:45 AS TIMESTAMP(6)), 50.30000, 204, 3.25003, Alabama.Montgomery, 2, d1002), (CAST(2024-12-19 19:12:46 AS TIMESTAMP(6)), 62.60000, 205, 0.33000, Alabama.Montgomery, 2, d1002), (CAST(2024-12-19 19:12:47 AS TIMESTAMP(6)), 72.30000, 206, 0.31000, Alabama.Montgomery, 2, d1002);; TableResult tableResult tEnv.executeSql(insertQuery); tableResult.await(); }示例 2通过 Table SQL 写入普通表写入普通表时使用sink.table.name替代sink.supertable.name且 DDL 字段中无需tbnamestatic void testNormalTableSqlToSink() throws Exception { EnvironmentSettings settings EnvironmentSettings.newInstance().inStreamingMode().build(); TableEnvironment tEnv TableEnvironment.create(settings); String tdengineSinkTableDDL CREATE TABLE sink_normal ( ts TIMESTAMP, current FLOAT, voltage INT, phase FLOAT ) 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.table.name sink_normal ); tEnv.executeSql(tdengineSinkTableDDL); String insertQuery INSERT INTO sink_normal VALUES (CAST(2024-12-19 19:12:45 AS TIMESTAMP(6)), 50.30000, 201, 3.31003), (CAST(2024-12-19 19:12:46 AS TIMESTAMP(6)), 82.60000, 202, 0.33000), (CAST(2024-12-19 19:12:47 AS TIMESTAMP(6)), 92.30000, 203, 0.31000), (CAST(2024-12-19 19:12:45 AS TIMESTAMP(6)), 50.30000, 204, 3.25003), (CAST(2024-12-19 19:12:46 AS TIMESTAMP(6)), 62.60000, 205, 0.33000), (CAST(2024-12-19 19:12:47 AS TIMESTAMP(6)), 72.30000, 206, 0.31000);; TableResult tableResult tEnv.executeSql(insertQuery); tableResult.await(); }示例 3将 Row 类型数据写入超级表对应的子表也可以通过tEnv.fromValues构造行数据Batch 模式再通过executeInsert写入static void testTableRowToSink() throws Exception { EnvironmentSettings settings EnvironmentSettings.newInstance().inBatchMode().build(); TableEnvironment tEnv TableEnvironment.create(settings); 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 ); tEnv.executeSql(tdengineSinkTableDDL); int sum 0; String tbname d001; int groupId 1; String location California.SanFrancisco; ListRow rows new ArrayList(); Random random new Random(System.currentTimeMillis()); for (int i 0; i 50; i) { sum 300 (i 1); long timestampInMillis System.currentTimeMillis() i * 1000; Row row Row.of( new Timestamp(timestampInMillis), // ts random.nextFloat() * 30, // current 300 (i 1), // voltage random.nextFloat(), // phase location, groupId, tbname ); rows.add(row); } Table inputTable tEnv.fromValues( DataTypes.ROW( DataTypes.FIELD(ts, DataTypes.TIMESTAMP(6)), DataTypes.FIELD(current, DataTypes.FLOAT()), DataTypes.FIELD(voltage, DataTypes.INT()), DataTypes.FIELD(phase, DataTypes.FLOAT()), DataTypes.FIELD(location, DataTypes.STRING()), DataTypes.FIELD(groupid, DataTypes.INT()), DataTypes.FIELD(tbname, DataTypes.STRING()) ), rows ); TableResult result inputTable.executeInsert(sink_meters); result.await(); // waiting for task completion }以上示例的完整可运行版本位于 docs/examples/flink/sink/Main.java其中main方法通过参数sink或table分别触发 DataStream Sink 与 Table Sink 相关测试。Flink 语义选择At-Least-Once连接器建议使用At-Least-Once语义原因如下TDengine 当前不支持事务无法进行频繁的 checkpoint 操作和复杂的事务协调。TDengine 使用时间戳作为主键下游算子可以通过对重复数据的过滤操作避免重复计算。使用At-Least-Once可以保证较高的数据处理性能和较低的数据延迟。设置方式StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(5000); env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.AT_LEAST_ONCE);数据类型映射TDengine 目前支持时间戳、数值、字符和布尔类型与 Flink RowData 类型的转换关系如下TDengine 数据类型Flink RowData 类型TIMESTAMPTimestampDataINTIntegerBIGINTLongFLOATFloatDOUBLEDoubleSMALLINTShortTINYINTByteBOOLBooleanVARCHARStringDataBINARYStringDataNCHARStringDataJSONStringDataVARBINARYbyte[]GEOMETRYbyte[]异常与错误码任务执行失败后请检查 Flink 任务的执行日志确认失败原因。常见的错误码及其处理建议如下错误码说明处理建议0xa000连接参数错误检查连接器参数配置0xa010数据库名称配置错误检查数据库名称配置0xa011表名称配置错误检查表名称配置0xa013value.deserializer 参数未设置设置序列化方法0xa014目标表的列名列表设置错误检查目标表的列名列表0x2301连接已关闭检查连接状态或新建连接后执行相关指令0x2302当前不支持该操作当前接口不支持可切换其他连接方式0x2303参数无效检查对应接口规范调整参数类型和大小0x2304statement 已关闭检查 statement 是否被关闭后复用或连接是否正常0x2305resultSet 已释放检查 ResultSet 是否被释放后再次使用0x230d参数索引超出范围检查参数的合理范围0x230e连接已关闭检查连接是否关闭后被再次使用0x230fTDengine 中存在未知 SQL 类型检查 TDengine 支持的数据类型0x2315TDengine 中存在未知类型检查将 TDengine 类型转换为 JDBC 类型时是否指定了正确的类型0x2319缺少用户名创建连接时补充用户名信息0x231a缺少密码创建连接时补充密码信息0x231d无法在指定时间内建立连接增加httpConnectTimeout参数或检查与 taosAdapter 的连接状态0x231e未在指定时间内完成任务增加messageWaitTimeout参数或检查与 taosAdapter 的连接0x2352不支持的编码本地连接指定了不支持的字符编码集0x2353数据库内部错误详见 taoslog本地连接执行 prepareStatement 时出错请检查 taoslog 定位问题0x2354连接为空本地连接执行命令时连接已关闭请检查与 TDengine 的连接0x2355结果集为空本地连接获取结果集异常请检查连接状态后重试0x2356字段数量无效本地连接结果集获取的元信息不匹配Maven 依赖如果使用 Maven 管理项目只需在pom.xml中添加如下依赖dependency groupIdcom.taosdata.flink/groupId artifactIdflink-connector-tdengine/artifactId version2.1.4/version /dependency小结通过flink-connector-tdengineApache Flink 可以与 TDengine 无缝集成一方面将复杂计算和深度分析得到的结果准确写入 TDengine 实现高效存储与管理另一方面企业版也可以快速稳定地读取 TDengine 中的海量数据做进一步分析。本文覆盖的 Sink 与 Table Sink 均基于jdbc:TAOS-WS://WebSocket 连接社区版即可使用是流批一体数据管道中连接 Flink 与 TDengine 的推荐方式。【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门