TDengine Flink Connector 使用指南:从 Sink 写入到 Table Sink 的流批集成实战
【免费下载链接】TDengineHigh-performance, scalable time-series database designed for Industrial IoT (IIoT) scenarios项目地址: https://gitcode.com/GitHub_Trending/tde/TDengine
Apache Flink 是 Apache 软件基金会支持的开源分布式流批一体处理框架,广泛用于流处理、批处理、复杂事件处理与实时数仓构建。本文基于 TDengine 官方文档与仓库中的可运行示例,系统讲解如何借助flink-connector-tdengine连接器,将 Flink 作业中的处理结果写入 TDengine(Sink),并通过 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.3):
| Flink 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.2 | Table Sink 支持 RowKind.UPDATE_BEFORE、RowKind.UPDATE_AFTER、RowKind.DELETE 等类型 | - |
| 2.0.1 | Sink 支持写入 RowData 实现类型 | - |
| 2.0.0 | 1. Sink 支持自定义数据结构序列化后写入 TDengine;2. 支持使用 Table SQL 写入 TDengine | 3.3.5.1 及以上 |
| 1.0.0 | 支持 Sink 功能,将其他数据源的数据写入 TDengine | 3.3.2.0 及以上 |
连接参数
建立连接的参数由 URL 和 Properties 两部分组成。URL 的规范格式为:
jdbc:TAOS-WS://[host_name]:[port]/[database_name]?[user={user}|&password={password}|&timezone={timezone}]参数说明:
| 参数 | 说明 | 默认值 |
|---|---|---|
| user | 登录 TDengine 的用户名 | root |
| password | 用户登录密码 | taosdata |
| database_name | 数据库名称 | - |
| timezone | 时区设置 | - |
| httpConnectTimeout | 连接超时时间,单位毫秒 | 60000 |
| messageWaitTimeout | 消息超时时间,单位毫秒 | 60000 |
| useSSL | 连接中是否使用 SSL | - |
需要注意的是,TDengine 的 Java 原生连接与 REST 连接已被标记为弃用(将于 2027-01-01 停止),连接器统一推荐使用 WebSocket 连接方式,即jdbc:TAOS-WS://协议前缀,并使用com.taosdata.jdbc.ws.WebSocketDriver驱动类。仓库中的示例代码即遵循该方式,例如:
static String jdbcUrl = "jdbc:TAOS-WS://localhost:6041?user=root&password=taosdata";Sink:将 Flink 处理结果写入 TDengine
Sink 的核心功能是将 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.java:
static 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; } DataStream<RowData> 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?user=root&password=taosdata"); sinkProps.setProperty(TDengineConfigParams.TD_BATCH_SIZE, "2000"); TDengineSink<RowData> 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; } DataStream<RowData> 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?user=root&password=taosdata"); sinkProps.setProperty(TDengineConfigParams.TD_BATCH_SIZE, "2000"); TDengineSink<RowData> 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; } DataStream<ResultBean> 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?user=root&password=taosdata"); sinkProps.setProperty(TDengineConfigParams.TD_BATCH_SIZE, "2000"); TDengineSink<ResultBean> 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.java):
CREATE 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子句声明连接器参数:
| 参数名称 | 类型 | 参数说明 |
|---|---|---|
| connector | string | 连接器标识,设置为tdengine-connector |
| td.jdbc.url | string | 连接的 URL |
| td.jdbc.mode | string | 连接器类型,设置为sink |
| sink.db.name | string | 目标数据库名称 |
| sink.batch.size | integer | 写入的批大小 |
| sink.supertable.name | string | 写入的超级表名称 |
| sink.table.name | string | 写入的普通表或子表名称 |
示例 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?user=root&password=taosdata'," + " '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 字段中无需tbname:
static 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?user=root&password=taosdata'," + " '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?user=root&password=taosdata'," + " '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"; List<Row> 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 类型 |
|---|---|
| TIMESTAMP | TimestampData |
| INT | Integer |
| BIGINT | Long |
| FLOAT | Float |
| DOUBLE | Double |
| SMALLINT | Short |
| TINYINT | Byte |
| BOOL | Boolean |
| VARCHAR | StringData |
| BINARY | StringData |
| NCHAR | StringData |
| JSON | StringData |
| VARBINARY | byte[] |
| GEOMETRY | byte[] |
异常与错误码
任务执行失败后,请检查 Flink 任务的执行日志确认失败原因。常见的错误码及其处理建议如下:
| 错误码 | 说明 | 处理建议 |
|---|---|---|
| 0xa000 | 连接参数错误 | 检查连接器参数配置 |
| 0xa010 | 数据库名称配置错误 | 检查数据库名称配置 |
| 0xa011 | 表名称配置错误 | 检查表名称配置 |
| 0xa013 | value.deserializer 参数未设置 | 设置序列化方法 |
| 0xa014 | 目标表的列名列表设置错误 | 检查目标表的列名列表 |
| 0x2301 | 连接已关闭 | 检查连接状态或新建连接后执行相关指令 |
| 0x2302 | 当前不支持该操作 | 当前接口不支持,可切换其他连接方式 |
| 0x2303 | 参数无效 | 检查对应接口规范,调整参数类型和大小 |
| 0x2304 | statement 已关闭 | 检查 statement 是否被关闭后复用,或连接是否正常 |
| 0x2305 | resultSet 已释放 | 检查 ResultSet 是否被释放后再次使用 |
| 0x230d | 参数索引超出范围 | 检查参数的合理范围 |
| 0x230e | 连接已关闭 | 检查连接是否关闭后被再次使用 |
| 0x230f | TDengine 中存在未知 SQL 类型 | 检查 TDengine 支持的数据类型 |
| 0x2315 | TDengine 中存在未知类型 | 检查将 TDengine 类型转换为 JDBC 类型时是否指定了正确的类型 |
| 0x2319 | 缺少用户名 | 创建连接时补充用户名信息 |
| 0x231a | 缺少密码 | 创建连接时补充密码信息 |
| 0x231d | 无法在指定时间内建立连接 | 增加httpConnectTimeout参数,或检查与 taosAdapter 的连接状态 |
| 0x231e | 未在指定时间内完成任务 | 增加messageWaitTimeout参数,或检查与 taosAdapter 的连接 |
| 0x2352 | 不支持的编码 | 本地连接指定了不支持的字符编码集 |
| 0x2353 | 数据库内部错误,详见 taoslog | 本地连接执行 prepareStatement 时出错,请检查 taoslog 定位问题 |
| 0x2354 | 连接为空 | 本地连接执行命令时连接已关闭,请检查与 TDengine 的连接 |
| 0x2355 | 结果集为空 | 本地连接获取结果集异常,请检查连接状态后重试 |
| 0x2356 | 字段数量无效 | 本地连接结果集获取的元信息不匹配 |
Maven 依赖
如果使用 Maven 管理项目,只需在pom.xml中添加如下依赖:
<dependency> <groupId>com.taosdata.flink</groupId> <artifactId>flink-connector-tdengine</artifactId> <version>2.1.4</version> </dependency>小结
通过flink-connector-tdengine,Apache 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),仅供参考