尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
SeaTunnel 多表 CDC 实战:一个作业同步 MySQL 多张表并按表自动路由到 PostgreSQL
SeaTunnel 多表 CDC 实战一个作业同步 MySQL 多张表并按表自动路由到 PostgreSQL【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文是一份基于 Apache SeaTunnel 的「多表 CDC」实战配方Recipe通过单个流式作业用 MySQL-CDC Source 的正则table-pattern捕获多张上游表的变化再用 JDBC Sink 的${table_name}占位符将每个源表自动路由到 PostgreSQL 中各自独立的目标表如st_orders、st_customers、st_products并借助generate_sink_sql true自动建表、自动生成 UPSERT 语句。读完本文你将掌握CDC 前置环境binlog、用户权限、驱动的搭建方法、最小可运行的多表 CDC 配置、端到端验证步骤以及多表路由在 SeaTunnel 内部的源码级实现原理。适用场景当你有以下需求时本配方是标准解法单作业多表不希望为每张表各起一个同步作业而是用一个作业批量捕获数十上百张表一致的快照起点多张表需要从同一时间点开始全量快照之后无缝切换到 binlog 增量按表自动路由每张源表的数据自动落到对应的目标表而不是被混写进同一张表流式持续同步作业常驻运行MySQL 侧的新增、更新持续流入目标库。多表同步的架构目标与设计取舍可进一步阅读 多表同步架构文档本文先聚焦如何把它跑起来。前置条件1. 已完成第一个作业先确保你能跑通 SeaTunnel 的基本作业流程参见 运行你的第一个作业。2. 安装本配方所需的插件按照 部署文档 下载连接器插件 的方式安装插件并在config/plugin_config中只保留以下两个插件--seatunnel-connectors-- connector-cdc-mysql connector-jdbc --end--执行安装并确认插件就位cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-(cdc-mysql|jdbc)3. 放置 JDBC 驱动SeaTunnel Zeta 引擎如果你使用 SeaTunnel Zeta 引擎需要同时把 MySQL 与 PostgreSQL 的 JDBC 驱动放入${SEATUNNEL_HOME}/lib然后确认可见ls ${SEATUNNEL_HOME}/lib | rg mysql-connector|postgresqlSeaTunnel 不捆绑所有 JDBC 驱动驱动 JAR 需自行下载Zeta 引擎统一从${SEATUNNEL_HOME}/lib加载驱动放置后需重启 SeaTunnel 进程。如果使用 Spark/Flink 引擎则需放到${SEATUNNEL_HOME}/plugins/Jdbc/lib/详见 JDBC Sink 文档。4. 准备 MySQL 源表每张上游表都应具备稳定的主键因为本配方会把 CDC 变更自动路由到下游的 upsert 目标表主键同时是路由与去重的依据CREATE DATABASE IF NOT EXISTS inventory; CREATE TABLE IF NOT EXISTS inventory.orders ( id BIGINT PRIMARY KEY, order_status VARCHAR(32), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS inventory.customers ( id BIGINT PRIMARY KEY, customer_name VARCHAR(64), city VARCHAR(64), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); CREATE TABLE IF NOT EXISTS inventory.products ( id BIGINT PRIMARY KEY, product_name VARCHAR(64), unit_price DECIMAL(10, 2), updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP ); INSERT INTO inventory.orders (id, order_status, updated_at) VALUES (2001, CREATED, NOW()); INSERT INTO inventory.customers (id, customer_name, city, updated_at) VALUES (3001, Alice, Shanghai, NOW()); INSERT INTO inventory.products (id, product_name, unit_price, updated_at) VALUES (4001, Keyboard, 99.00, NOW());关于无主键表MySQL CDC 默认期望源表有主键。若某张表没有主键但存在唯一列可通过table-names-config.primaryKeys指定自定义主键如果完全没有稳定的唯一键UPDATE/DELETE 事件将无法安全地在下游应用详见 MySQL CDC 文档。5. 创建 MySQL CDC 用户并授权CDC 用户需要SELECT读快照、RELOAD、SHOW DATABASES以及REPLICATION SLAVE, REPLICATION CLIENT读 binlogCREATE USER IF NOT EXISTS st_user_source% IDENTIFIED BY mysqlpw; GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO st_user_source%; FLUSH PRIVILEGES;6. 确认 MySQL binlog 配置SHOW VARIABLES WHERE variable_name IN (log_bin, binlog_format, binlog_row_image);期望值为log_bin ON、binlog_format ROW、binlog_row_image FULL。若未开启需在 MySQL 配置文件中启用[mysqld] server-id 223344 log_bin mysql-bin expire_logs_days 10 binlog_format row binlog_row_image FULL并重启 MySQL。注意当初始快照较大时还需适当调大interactive_timeout与wait_timeout防止快照期间连接超时MySQL CDC 文档 中给出了更完整的 binlog 与会话超时说明。7. 准备 PostgreSQL 目标库与写入用户CREATE USER st_user_sink WITH PASSWORD pgpw; CREATE DATABASE sync_demo; GRANT ALL PRIVILEGES ON DATABASE sync_demo TO st_user_sink;重新连接sync_demo后授予publicschema 上的建表权限GRANT USAGE, CREATE ON SCHEMA public TO st_user_sink;本配方使用了generate_sink_sql true因此 SeaTunnel 会在首次运行时自动创建public.st_orders、public.st_customers等目标表前提是 sink 用户对publicschema 拥有CREATE权限。最小配置下面这份配置通过一个table-pattern读取 MySQL 的多张表并写入 PostgreSQL 中名为st_上游表名的目标表env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { MySQL-CDC { plugin_output mysql_multi startup.mode initial server-id 5652 username st_user_source password mysqlpw database-pattern inventory table-pattern inventory\\.(orders|customers|products) url jdbc:mysql://mysql:3306/inventory } } sink { Jdbc { plugin_input mysql_multi driver org.postgresql.Driver url jdbc:postgresql://postgresql:5432/sync_demo username st_user_sink password pgpw generate_sink_sql true database sync_demo table public.st_${table_name} primary_keys [${primary_key}] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }关键参数逐项解读env 段参数值说明parallelism1本配方用单并行度演示多表场景可调大并行度以提升吞吐job.modeSTREAMING流式作业快照完成后持续消费 binlogcheckpoint.interval5000每 5 秒做一次 checkpoint决定下游可见性与故障恢复粒度source 段MySQL-CDC参数值说明plugin_outputmysql_multi为数据流命名下游plugin_input须与之对应startup.modeinitial启动时先做全量快照再无缝切换增量其他可选值earliest、latest、specific、timestampserver-id5652该 CDC 读取器在 MySQL 集群中的唯一 ID也可配区间如5652-5660多读取器/多表并行时区间要足够大重复会导致 MySQL 踢掉其中一个客户端database-patterninventory要捕获的库名正则例如database_prefix.*table-patterninventory\\.(orders\|customers\|products)要捕获的表名正则匹配名包含库名table-pattern与table-names互斥二选一urljdbc:mysql://mysql:3306/inventoryJDBC 连接地址其中table-pattern是「一个作业捕获多表」的关键SeaTunnel 会把正则命中的每张表各自解析成独立的数据分片与元数据表结构、主键而不是把多张表混成一个 schema。对应的选项定义与互斥校验可查看源码 MySqlIncrementalSourceFactory.java 中对TABLE_NAMES/TABLE_PATTERN的解析逻辑。sink 段Jdbc参数值说明plugin_inputmysql_multi对应上游plugin_outputdriver/urlPostgreSQL 驱动与地址本配方向 PostgreSQL 写入generate_sink_sqltrue由 SeaTunnel 依据上游 schema 与行类型INSERT/UPDATE/DELETE自动生成 SQL若为false则必须手写querydatabasesync_demo目标库名tablepublic.st_${table_name}路由核心${table_name}占位符会被替换为上游记录携带的真实表名于是orders→public.st_orders、customers→public.st_customersprimary_keys[${primary_key}]同样支持占位符从上游元数据继承每张表的主键用于生成数据库原生的 UPSERT / UPDATE / DELETEschema_save_modeCREATE_SCHEMA_WHEN_NOT_EXIST目标表不存在时自动创建RECREATE_SCHEMA会先删后建、ERROR_WHEN_SCHEMA_NOT_EXIST缺失即报错、IGNORE跳过建表逻辑data_save_modeAPPEND_DATA保留已有数据继续追加DROP_DATA会清空、ERROR_WHEN_DATA_EXISTS有数据即报错关于${table_name}占位符JDBC Sink 的table参数支持${table_name}与${schema_name}变量${schema_name}会被替换为目标侧 schema 名${table_name}会被替换为目标侧表名。示例写法如test_${schema_name}_${table_name}_test、public.${table_name}_test等详见 JDBC Sink 文档。对带 schema 概念的数据库如 PostgreSQL、Oracle、SQL Servertable需写成xxx.xxx形式。运行作业将配置保存为config/multi-table-cdc.conf然后以本地模式启动 SeaTunnelcd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/multi-table-cdc.conf -m local因为这是一个流式 CDC 管道请保持作业运行同时到 MySQL 侧制造新的变更来验证增量同步。验证结果第一步确认建表等初始快照阶段结束后在 PostgreSQL 中查询 SeaTunnel 自动创建的目标表SELECT table_name FROM information_schema.tables WHERE table_schema public AND table_name LIKE st_% ORDER BY table_name;预期能看到st_orders、st_customers、st_products三张表。第二步在 MySQL 各源表上制造变更INSERT INTO inventory.orders (id, order_status, updated_at) VALUES (2002, PAID, NOW()); UPDATE inventory.customers SET city Hangzhou, updated_at NOW() WHERE id 3001; INSERT INTO inventory.products (id, product_name, unit_price, updated_at) VALUES (4002, Mouse, 59.00, NOW());第三步确认每张下游表只收到自己的数据SELECT id, order_status FROM public.st_orders ORDER BY id; SELECT id, customer_name, city FROM public.st_customers ORDER BY id; SELECT id, product_name, unit_price FROM public.st_products ORDER BY id;st_orders应包含 2001/2002 两条订单st_customers中 Alice 的城市应已变为 HangzhouUPSERT 生效st_products应包含 4001/4002 两条产品。如果每个上游表都被路由到各自的目标表且变更持续流动说明多表 CDC 管道工作正常。底层原理多表路由是怎么实现的理解了「怎么配」再看「为什么这样就能路由」。整条链路由三层机制协作完成1. Source 侧正则展开为多张表的独立元数据table-pattern命中多张表后MySQL CDC Source 会为每张表生成独立的 catalog 元数据CatalogTable与读取分片split。每张表拥有自己的 schema 与主键信息为下游按表建表和按表生成 SQL 提供了基础。全量快照阶段完成后Source 会自动切换到从快照开始时记录的 binlog 位点继续消费增量保证切换期间不丢事件MySQL CDC 文档 的 FAQ 有明确说明。server-id在源码层面被解析为ServerIdRange每个并行子任务会从区间中分得一个唯一 ID 并写入 Debezium 配置database.server.id避免多读取器在同一 MySQL 集群上发生 server-id 冲突见 MySqlSourceConfigFactory.java。2. 传输侧每行数据都携带 tableId多表管道中每一行SeaTunnelRow都会携带一个tableId即TablePath的序列化形式databaseName.schemaName.tableName。TablePath是记录归属的唯一标识定义于 TablePath.java。正是这个 tableId 让下游无需猜测「这行数据属于哪张表」。3. Sink 侧MultiTableSink 按表分桶写入JDBC Sink 在接收到多表输入时会被包装为框架层的MultiTableSink。它会为每张表创建独立的子 Sinkwriter并在运行时根据行的tableId把记录路由到对应子 writer。核心实现见 MultiTableSink.java 与 MultiTableSinkWriter.java。路由策略分两种情况有主键时哈希路由对主键字段值取哈希并映射到队列下标保证同一主键的记录永远进入同一个队列从而保持该键内的写入顺序int index (object.hashCode() Integer.MAX_VALUE) % blockingQueues.size();源码中刻意用 Integer.MAX_VALUE清掉符号位而非Math.abs因为Math.abs(Integer.MIN_VALUE)仍返回负数会导致下标越界。无主键时随机路由在各队列间随机分布以均衡负载但不保证同一键的顺序。每个队列对应一个 replica副本 writer可通过multi_table_sink_replica配置每个表的并发 writer 数量以提升吞吐。schema 变更如新增列则以 barrier 形式广播到所有队列确保每张表的行流在同一个位置应用变更、不乱序。4.${table_name}占位符如何生效子 Sink 写入时table参数中的${table_name}、${schema_name}会被替换为该行所属TablePath中的真实表名、schema 名primary_keys中的${primary_key}则替换为该表元数据中的主键列。配合generate_sink_sql trueSeaTunnel 即可为每张表生成专属的CREATE TABLE与INSERT ... ON CONFLICT ... DO UPDATEPostgreSQL 原生 UPSERT语句。这就是「每张源表自动路由到各自目标表」的完整闭环。常见陷阱以下是本配方最容易踩的坑逐一对照排查table-pattern中的正则转义错误在 HOCON 里.表示任意字符匹配字面点号时必须写成\\.。例如inventory\\.(orders|customers|products)写错会导致表匹配不到或误匹配。MySQL binlog 或 CDC 用户权限不完整表现是作业能读完快照但无法继续读增量。核对binlog_format ROW、binlog_row_image FULL并确认用户具备REPLICATION SLAVE, REPLICATION CLIENT。未配置占位符路由如果 sink 的table写死成单张表名没有${table_name}多张源表会被意外混写进同一张目标表。PostgreSQL sink 用户缺少CREATE权限能连接数据库但无法在publicschema 上建表。需执行GRANT USAGE, CREATE ON SCHEMA public TO st_user_sink;。上游表没有主键却按 upsert 语义配置没有稳定主键时UPSERT/UPDATE/DELETE 无法安全应用。要么给表加主键要么用table-names-config.primaryKeys指定唯一列。schema/table 占位符用错了位置不同数据库对 schema 的语义不同如 MySQL 无独立 schema、PostgreSQL 有public${schema_name}与${table_name}的用法需与目标库的命名规则匹配。相关文档MySQL CDC 连接器文档table-pattern/table-names互斥规则、server-id区间、startup.mode全量增量切换、binlog 与用户权限的完整说明JDBC Sink 连接器文档generate_sink_sql两种写入模式、${table_name}/${schema_name}占位符规则、schema_save_mode/data_save_mode取值、UPSERT 行为多表同步架构文档TablePath、MultiTableSink、replica 副本机制、schema 变更路由与故障隔离策略的源码级剖析【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED

相关推荐

摘要插入无关指令案例,TaoToken Key 跑 27 组样本

摘要插入无关指令案例,TaoToken Key 跑 27 组样本

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📅 2026/9/18 4:19:25
神经网络数据集规模与模型性能:学习曲线实战指南

神经网络数据集规模与模型性能:学习曲线实战指南

1. 神经网络和数据集的思考:先把"数据越多越好"这句话拆开看做神经网络这几年,我被问得最多的问题之一就是:"我这个模型效果不行,是不是训练数据太少了?再弄几万条是不是就好了?"问这话…

📅 2026/9/18 4:19:25
舆情实体识别 GLiFormer,情感归因 Agent Key 来自 TaoToken

舆情实体识别 GLiFormer,情感归因 Agent Key 来自 TaoToken

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📅 2026/9/18 4:19:25
MORE NEWS

更多资讯

📰

泛微E9 API接口调用全流程详解:从Token获取到签名校验的实战指南

泛微E9的API接口调用,说难不难,说简单也不简单。很多第一次接触泛微E9二开的同学,最容易卡住的地方不是Java语法,也不是HTTP请求怎么写,而是根本摸不清整个调用过程的全貌:token怎么拿、请求地址拼到哪、签…

📰

MiroFish:轻量级Miro白板本地化部署方案

1. 项目概述:MiroFish不是鱼,而是一套面向协作白板场景的轻量级镜像部署方案MiroFish这个名称乍一听容易让人联想到某种生物实验或海洋科技项目,但实际在当前协作工具生态中,它指的是一套专为Miro白板平台设计的、可本地化快速部署…

📰

大语言模型的技术潜力与局限分析

1. 大语言模型的技术潜力边界2023年ChatGPT的爆发让LLM(大语言模型)成为技术焦点,但从业界讨论来看,对其潜力评估呈现两极分化。我参与过多个NLP项目开发,发现LLM在特定场景表现惊人,但在某些基础能力上仍存…

📰

SpringBoot+Vue构建智能农业疾病防治系统

1. 项目概述果蔬作物疾病防治系统是一个面向现代农业的智能化管理平台,旨在解决传统农业中疾病防治效率低下、专业知识获取困难等问题。作为一名长期从事农业信息化系统开发的工程师,我在实际项目中发现,许多农户在面对作物疾病时往往缺乏有效…

📰

hermes智能体运行环境:从部署到配置DeepSeek的完整实践

这几个月我一直在折腾一个叫 hermes 的智能体,最开始只是出于好奇,后来发现它几乎把我桌面上那些零散的 AI 脚本全收编了。hermes 本身是一个开源的智能体运行环境,你可以把它理解为 AI 助手的“运行时”——它负责接收任务、调度模型、调用工…

📰

LiveTalking:5 步让数字人在浏览器里开口对话,本地跑、免费开源

LiveTalking:5 步让数字人在浏览器里开口对话,本地跑、免费开源 【免费下载链接】metahuman-stream Real time interactive streaming digital human 项目地址: https://gitcode.com/GitHub_Trending/me/metahuman-stream 昨晚开播,你…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬