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

资讯详情

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

SeaTunnel TiDB-CDC 源连接器实战指南:通过 tikv-client-java 读取 TiDB 快照与增量变更

SeaTunnel TiDB-CDC 源连接器实战指南:通过 tikv-client-java 读取 TiDB 快照与增量变更 SeaTunnel TiDB-CDC 源连接器实战指南通过 tikv-client-java 读取 TiDB 快照与增量变更【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文以 SeaTunnel 仓库中的 TiDB CDC 官方文档 为核心系统讲解 TiDB-CDC 源连接器的工作原理、全部配置参数、类型映射与作业配置示例并结合connector-cdc-tidb模块的源码工厂类、Split 枚举器、SourceReader 与 Split 定义深入剖析其快照 增量两阶段读取机制、并行切分策略与 exactly-once 状态恢复流程帮助你既会配置这个连接器又理解它底层是如何工作的。连接器概述TiDB CDC 源连接器通过tikv-client-javaJava 客户端直接与 TiDB 的 TiKV 集群交互——它连接 TiKV 的放置驱动器Placement Driver即 PD和 TiKV 节点从而读取快照数据与增量变更事件。它支持并行快照读取和 exactly-once 流式语义是官方推荐的将 TiDB 表引入 SeaTunnel 数据管道的方式。支持的执行引擎文档 Support Those Engines 一节SeaTunnel ZetaFlink依赖准备使用 TiDB-CDC 连接器需要额外提供两个依赖 jarMySQL JDBC 驱动com.mysql.cj.jdbc.Driver和tikv-client-java3.2.0。放置位置取决于所用引擎Flink 引擎两个 jar 都放入${SEATUNNEL_HOME}/plugins/目录Zeta 引擎两个 jar 都放入${SEATUNNEL_HOME}/lib/目录。从源码看JDBC 驱动并非仅用于发现元数据这么简单TiDBSourceFactory 在createSource时、TiDBSource 在createReader/createEnumerator/restoreEnumerator时都会显式执行Class.forName(com.mysql.cj.jdbc.Driver)把驱动注册进DriverManager因此缺少该 jar 会直接导致作业启动失败。核心特性Key Features按官方文档的特性勾选表特性支持情况batch批处理不支持stream流式支持exactly-once精确一次支持column projection列裁剪不支持parallelism并行支持user-defined split自定义切分不支持这一特性表与源码一致TiDBSource 的getBoundedness()固定返回Boundedness.UNBOUNDED即该连接器只作为无界流式源运行这正是它必须配合job.mode STREAMING和 checkpoint 间隔使用的原因。数据源版本与驱动信息数据源支持版本驱动连接串示例MySQL 协议TiDB 兼容MySQL 5.5 / 5.6 / 5.7 / 8.0.xRDS MySQL 5.6 / 5.7 / 8.0.xcom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testtikv-client-java3.2.0--注意这里的 MySQL 版本约束指的是 TiDB 对外暴露的 MySQL 兼容协议层TiDB 通过该协议提供表结构元数据供 Catalog 读取而真正读取数据走的是 TiKV 的 KV 接口。数据类型映射TiDB/MySQL 类型到 SeaTunnel 类型的映射关系如下完整继承自官方文档MySQL 数据类型SeaTunnel 数据类型BIT(1)、TINYINT(1)BOOLEANTINYINTTINYINTTINYINT UNSIGNED、SMALLINTSMALLINTSMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINTINT UNSIGNED、INTEGER UNSIGNED、BIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20, 0)DECIMAL(p, s)、DECIMAL(p, s) UNSIGNED、NUMERIC(p, s)、NUMERIC(p, s) UNSIGNEDDECIMAL(p, s)FLOAT、FLOAT UNSIGNEDFLOATDOUBLE、DOUBLE UNSIGNED、REAL、REAL UNSIGNEDDOUBLECHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、ENUM、JSONSTRINGDATEDATETIME(s)TIME(s)DATETIME、TIMESTAMP(s)TIMESTAMP(s)BINARY、VARBINARY、BIT(p)、TINYBLOB、MEDIUMBLOB、BLOB、LONGBLOB、GEOMETRYBYTES配置参数详解Source Options下表完整覆盖官方文档列出的全部参数并补充了源码中的默认值与校验逻辑参数名类型必填默认值说明urlString是-用于发现表元数据的 MySQL 兼容 JDBC URL例如jdbc:mysql://tidb0:4000/inventoryusernameString是-连接 TiDB 服务器的用户名passwordString是-连接 TiDB 服务器的密码pd-addressesString是-TiKV 集群的 PD 地址逗号分隔例如pd0:2379,pd1:2379database-nameString是-要监控的 TiDB 数据库名table-nameString是-database-name下要监控的表名不要包含库名startup.modeEnum否initial启动模式合法取值initial/earliest/latest。initial先做历史快照再持续读取增量earliest从最早可用 offset 开始latest跳过快照只消费之后的新变更batch-size-per-scanInt否1000每次 TiKV scan 请求读取的行数tikv.grpc.timeout_in_msLong否-TiKV gRPC 客户端超时毫秒。TiKV 在高负载下响应慢时调大tikv.grpc.scan_timeout_in_msLong否-TiKV gRPC 扫描超时毫秒。大表扫描超时时调大tikv.batch_get_concurrencyInteger否-TiKVBatchGet请求并发度读请求受 TiKV CPU 瓶颈时调大tikv.batch_scan_concurrencyInteger否-TiKVBatchScan请求并发度快照读受 TiKV CPU 瓶颈时调大参数在源码中的定义与校验所有参数集中定义在 TiDBSourceOptions 中STARTUP_MODE通过singleChoice(StartupMode.class, [INITIAL, EARLIEST, LATEST])声明因此startup.mode配置了其他取值如specific会直接触发选项校验失败BATCH_SIZE_PER_SCAN的默认值 1000 与文档一致。参数合法性规则由 TiDBSourceFactory 的optionRule()声明必填项为database-name、table-name、pd-addresses可选项为四个tikv.*调优参数与startup.mode。url/username/password则由内部的TiDBCatalogFactory在创建 Catalog 时使用用于读取表结构。getTiConfiguration(ReadonlyConfig)方法把用户配置注入TiConfigurationtikv-client-java 的核心配置对象pd-addresses通过TiConfiguration.createDefault(pdAddrsStr)初始化其余四项通过setTimeout/setScanTimeout/setBatchGetConcurrency/setBatchScanConcurrency覆盖默认值——未配置时完全沿用 tikv-client-java 的内置默认。完整作业示例示例一快照 增量写入 JDBC Sink流式消费一张 TiDB 表的 CDC 事件并写入 JDBC Sink。需要job.mode STREAMING以及 checkpoint 间隔增量事件才能持续流动该示例与仓库中 TiDB CDC 端到端测试使用的作业形态一致env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { TiDB-CDC { plugin_output products_tidb_cdc url jdbc:mysql://tidb0:4000/tidb_cdc driver com.mysql.cj.jdbc.Driver tikv.grpc.timeout_in_ms 20000 pd-addresses pd0:2379 username root password database-name tidb_cdc table-name tidb_cdc_e2e_source_table } } sink { Jdbc { plugin_input products_tidb_cdc url jdbc:mysql://tidb0:4000/tidb_cdc driver com.mysql.cj.jdbc.Driver username root password database tidb_cdc table tidb_cdc_e2e_sink_table generate_sink_sql true primary_keys [id] } }示例二只消费最新变更latest 模式当只需要新变更、需要跳过历史快照时使用startup.mode latestsource { TiDB-CDC { url jdbc:mysql://tidb0:4000/tidb_cdc driver com.mysql.cj.jdbc.Driver pd-addresses pd0:2379 username root password database-name tidb_cdc table-name tidb_cdc_e2e_source_table startup.mode latest } }工作原理源码级的三阶段剖析1. Split 切分按 KeyRange 并行切表TiDBSourceSplitEnumerator 在open()中创建TiSession并通过 Catalog 拿到表 IDtableId在run()中调用TableKeyRangeUtils.getTableKeyRanges(tableId, currentParallelism)按并行度把表的数据 Key 空间切分成若干Coprocessor.KeyRange每个 KeyRange 生成一个 TiDBSourceSplit。几个源码细节值得注意Split 的初始位点由启动模式决定枚举器构造 Split 时传入resolvedTs (startupMode INITIAL) ? -1 : 0。-1表示先做快照Reader 收到后以当前 TSO 时间戳为快照读位点0表示直接进入增量阶段。轮询分配枚举器用assignCount % parallelism把 Split 轮询分配给各个 Reader并在snapshotState()中把shouldEnumerate、pendingSplit与assignCount持久化到TiDBSourceCheckpointState——这意味着作业从 checkpoint 恢复后Split 归属保持原有序列不会因重分配产生倾斜。失败回补addSplitsBack(splits, subtaskId)支持把失败 Reader 尚未被 checkpoint 确认的 Split 退回重新分配这是 exactly-once 语义中至少投递一次、由下游去重环节的关键。2. 快照阶段RegionStoreClient 逐 Region 扫描TiDBSourceReader 的pollNext()遵循固定顺序当startup.mode为INITIAL时先遍历所有未完成的 Split 执行snapshotEvents()再做增量captureStreamingEvents()。snapshotEvents()的实现逻辑是通过session.getTimestamp().getVersion()取一个 TSO 时间戳作为快照读位点保证快照的一致性视图以RegionStoreClient从 Split 的snapshotStart位置开始按 Region 逐段 scan扫描到空结果时前进到下一个 Region 的起始 Key直到越过 KeyRange 的结束位置每读一批就更新split.setSnapshotStart(start)——这个游标会随snapshotState()进入 checkpoint因此快照中断后恢复可以从断点继续而不必重扫整个 Split只反序列化记录键TableKeyRangeUtils.isRecordKey索引键被跳过。batch-size-per-scan控制每轮 poll 从 CDC 客户端拉取的最大行数配合maxPollAttemptsbatchSize 的 10 倍上限避免长时间空转等待元数据事件。3. 增量阶段CDCClient prewrite/commit 配对增量读取是 TiDB CDC 最有技术含量的部分每 Split 一个 CDCClientgetCdcClient()为每个 Split 懒创建org.tikv.cdc.CDCClienttikv-client-java 提供的 watch 客户端并从 Split 的resolvedTs开始订阅只监听该 Split 的 KeyRange。事务两阶段配对TiKV 的 CDC 事件是 MVCC 层的PREWRITE/COMMIT/ROLLBACK事件。Reader 内部维护两个TreeMap——preWrites按 startTs 排序与commits按 commitTs 排序收到 COMMIT 事件后到preWrites中查找对应的预写记录配对成功才把该行事件放入committedEvents阻塞队列找不到对应 prewrite 时safeResolvedTs会回退Math.min(safeResolvedTs, commitTs - 1)并暂缓推进防止乱序丢失。ROLLBACK事件则直接从preWrites中删除从而不产生任何输出。resolvedTs 推进每轮 flush 后取cdcClient.getMinResolvedTs()所有 TiKV region 都到达的时间戳作为新的安全位点写回 Split。Reader 还会周期性输出[TiDB-CDC-DIAG]前缀的统计日志polledRows、emittedRows、resolvedLagMs、pull/flush/emit 耗时等排障时可直接在运行日志中观察消费延迟。读写字段解耦committedEvents队列把拉取 CDC 事件与反序列化输出分离源码注释明确说明这是为了防止下游 sink如 JDBC写慢时阻塞 CDC 事件拉取——Region 分裂期间若 pull 阻塞会丢事件这是该队列存在的直接原因。4. 状态恢复exactly-once 的落点TiDBSource实现了restoreEnumerator(context, checkpointState)作业恢复时枚举器用 checkpoint 中的TiDBSourceCheckpointState重建pendingSplit与轮询游标Reader 侧则通过snapshotState()返回当前 Split 列表含snapshotStart快照游标与resolvedTs增量游标。快照游标 增量游标两个恢复点共同保证作业失败重启后已 checkpoint 确认的数据不重复未确认的数据重放后由下游主键去重如 JDBC Sink 的primary_keys配置保证最终 exactly-once 效果。使用注意事项Notes官方文档明确列出三条注意事项均与源码行为对应一个 source 块只读一张表TiDBSourceFactory每次只按单个database-nametable-name构建一个TablePath并取出一个CatalogTable。一个作业要采集多张表就写多个TiDB-CDCsource 块配合不同的plugin_output。startup.mode specific不是合法选项TiDBSourceOptions中singleChoice只注册了initial/earliest/latest三个枚举值配置specific会在选项校验阶段报错。tikv.grpc.*与tikv.batch_*_concurrency仅在默认值不满足集群需求时调整这些参数直接透传给 tikv-client-java 的TiConfiguration默认值通常足够盲目调大并发可能加剧 TiKV 侧 CPU 压力。变更记录该连接器的演进记录维护在 connector-cdc-tidb 变更记录 中关键节点包括2.3.8TiDB CDC 源连接器首次落地Feature: Support tidb cdc connector source #7199 / #74772.3.11新增 source/sink 状态类serialVersionUID缺失检查脚本#9118对本连接器的 checkpoint 兼容性有直接意义2.3.12修正batch-size-per-scan选项 key 的拼写问题#9434、优化 CDC JAR 体积#9546、优化枚举器 API 语义以降低连接器层锁调用#9671、增加每个连接器独立插件目录支持#9650。这也解释了为什么当前文档中的选项名是batch-size-per-scan——早期版本的拼写错误已在 2.3.12 修正旧拼写不再可用。小结TiDB-CDC 源连接器是一个以 TiKV 为一等公民的 CDC 实现元数据走 MySQL 协议 Catalog快照与增量全部通过 tikv-client-java 的 Region 扫描和 CDC watch 完成。掌握本文要点后你可以完成三件事按引擎类型正确放置 JDBC 驱动与 tikv-client-java 依赖根据业务诉求选择initial/earliest/latest启动模式并编写完整可运行的 HOCON 作业当作业出现快照慢、增量延迟或恢复异常时结合[TiDB-CDC-DIAG]诊断日志与 Split 状态快照游标snapshotStart、增量位点resolvedTs进行定位必要时通过tikv.grpc.*超时与tikv.batch_*_concurrency参数针对性调优。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表