十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel JDBC to JDBC 实战:MySQL 到 PostgreSQL 批式迁移与 SQL Transform 数据重塑

SeaTunnel JDBC to JDBC 实战:MySQL 到 PostgreSQL 批式迁移与 SQL Transform 数据重塑 SeaTunnel JDBC to JDBC 实战MySQL 到 PostgreSQL 批式迁移与 SQL Transform 数据重塑【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文是一份完整的 SeaTunnel 实战配方Recipe指南面向需要在两个关系型数据库之间做批量迁移、并在写入前对数据进行过滤与重塑的场景。文章以仓库文档 docs/en/getting-started/recipes/jdbc-to-jdbc.md 为主线完整复刻其从环境准备、建表、配置、运行到结果验证的全过程并结合 connector-jdbc 与 SQLTransform.java 源码说明底层原理。读完本文你将掌握如何用一条 SeaTunnel 批处理作业完成 MySQL → PostgreSQL 的数据抽取、Sql变换与落地写入并理解plugin_input/plugin_output三阶段串联机制与参数绑定顺序。何时使用这条 Recipe在 Recipe 总览 的规划中JDBC to JDBC 被定位为“关系型数据库之间的批量迁移且带行级转换”的管道形态。它的典型适用点源库与目标库都是通过 JDBC 驱动访问的关系型数据库本文以 MySQL → PostgreSQL 为例需要一个有界bounded批处理作业读完当前可见数据后作业即结束而非持续捕获增量写入前需要对行做过滤如只保留status PAID的订单或重塑改名、改类型、补充常量字段。本文示例的完整数据流如下MySQL source_db.orders │ Jdbc sourceplugin_output mysql_orders ▼ Sql transformplugin_input mysql_orders → plugin_output paid_orders │ 过滤 statusPAID、UPPER 客户名、金额改精度、补 source_system ▼ PostgreSQL target_db.public.paid_orders ▲ Jdbc sinkplugin_input paid_orders前置准备1. 完成首个本地任务本 Recipe 假设你已经能跑通一个最简作业。请先按 Run your first job 完成本地部署并确认bin/seatunnel.sh可用这能提前排除安装、配置解析和执行引擎方面的问题。2. 安装 JDBC 连接器插件SeaTunnel 从 2.2.0-beta 起二进制发行包默认不再捆绑连接器依赖需要按 Deployment Download The Connector Plugins 的方式安装。先在 config/plugin_config 中只保留本作业需要的插件--seatunnel-connectors-- connector-jdbc --end--然后执行安装并确认插件落盘cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | rg connector-jdbc需要特别说明的是Sql变换transform随 SeaTunnel 发行版内置不需要在plugin_config中单独添加 connector 条目。这一点可以从源码得到印证——SQLTransform.java 中PLUGIN_NAME Sql它属于seatunnel-transforms-v2模块而非独立的连接器插件。3. 放置数据库驱动JDBC 驱动的许可证与再分发条款因数据库厂商而异驱动版本还必须同时兼容数据库与 Java 运行时因此 SeaTunnel不捆绑任何 JDBC 驱动需要自行下载后放到引擎对应的目录Zeta 引擎本 Recipe 使用的默认引擎放在${SEATUNNEL_HOME}/lib/并重启相关 SeaTunnel 进程使驱动被加载Spark / Flink 引擎放在每个执行节点${SEATUNNEL_HOME}/plugins/Jdbc/lib/。放置后确认 SeaTunnel 能看到两个驱动ls ${SEATUNNEL_HOME}/lib | rg mysql-connector|postgresql常见驱动文件名为mysql-connector-j-8.x.x.jarMySQL与postgresql-42.x.x.jarPostgreSQL。详细的驱动类名、URL 模板与各厂商对照可参考 JDBC Source 的 Driver reference 和 JDBC Sink 的 Driver reference。4. 准备 MySQL 源表创建源库与源表并插入三条确定性的数据行便于后续精确验证CREATE DATABASE IF NOT EXISTS source_db; CREATE TABLE IF NOT EXISTS source_db.orders ( id BIGINT PRIMARY KEY, customer_name VARCHAR(100) NOT NULL, amount DECIMAL(16, 4) NOT NULL, status VARCHAR(20) NOT NULL ); TRUNCATE TABLE source_db.orders; INSERT INTO source_db.orders (id, customer_name, amount, status) VALUES (1001, alice chen, 120.5000, PAID), (1002, bob li, 80.0000, CREATED), (1003, carol wu, 42.0000, PAID);注意amount在源表是DECIMAL(16, 4)4 位小数目标表是DECIMAL(12, 2)2 位小数这一步刻意制造了“精度重塑”的场景。5. 准备 PostgreSQL 目标库与目标表创建目标用户与数据库。sink 用户需要对目标表拥有INSERT权限CREATE USER test WITH PASSWORD test; CREATE DATABASE target_db OWNER test;以test用户重新连接到target_db然后建表CREATE TABLE IF NOT EXISTS public.paid_orders ( id BIGINT PRIMARY KEY, customer_name VARCHAR(100) NOT NULL, amount DECIMAL(12, 2) NOT NULL, source_system VARCHAR(20) NOT NULL ); TRUNCATE TABLE public.paid_orders;目标表比源表多出一个source_system字段用于承接变换阶段注入的常量值。如果用户或数据库已存在直接复用即可无需重复执行对应的CREATE语句。完整配置逐段解析将以下配置保存为config/jdbc-to-jdbc.conf。整个配置由env、source、transform、sink四段组成三处插件通过plugin_output与plugin_input首尾相连构成完整的数据链。env { parallelism 1 job.mode BATCH } source { Jdbc { plugin_output mysql_orders driver com.mysql.cj.jdbc.Driver url jdbc:mysql://mysql.example.com:3306/source_db?useSSLfalseallowPublicKeyRetrievaltrue username root password password query SELECT id, customer_name, amount, status FROM orders } } transform { Sql { plugin_input mysql_orders plugin_output paid_orders query SELECT id, UPPER(customer_name) AS customer_name, CAST(amount AS DECIMAL(12, 2)) AS amount, MYSQL AS source_system FROM dual WHERE status PAID } } sink { Jdbc { plugin_input paid_orders driver org.postgresql.Driver url jdbc:postgresql://postgresql.example.com:5432/target_db username test password test query INSERT INTO public.paid_orders (id, customer_name, amount, source_system) VALUES (?, ?, ?, ?) } }env批处理模式env { parallelism 1 job.mode BATCH }job.mode BATCH声明这是一个有界批处理作业源端读完后作业即正常结束。parallelism 1表示单并发读取与写入。若数据量较大可提升并行度并使用 JDBC Source 的分区能力详见下文“进阶”小节。sourceJdbc 源JDBC Source 文档 说明它是一个有界源读取查询可见的所有行后结束如果作业需要持续捕获后续的增删改应改用 CDC 连接器。本段用query控制只读源表的部分列driverJDBC 驱动类名MySQL 为com.mysql.cj.jdbc.Driverurl连接串。示例中useSSLfalseallowPublicKeyRetrievaltrue仅用于简单开发环境生产环境需按组织安全要求配置 TLSusername/password账号信息。username是首选键旧配置user仍作为兜底兼容query自定义查询 SQL可按需做列裁剪column projection与库端过滤plugin_output mysql_orders把读出的行以表名mysql_orders暴露给下游变换。从源码看该参数由 JdbcSourceOptions.java 定义query与table_path可以单独或组合使用使用query时 SeaTunnel 会从首列所在底层表继承元数据。对多表 JOIN 或复杂查询推断出的分区键不一定在整个结果集上唯一本 Recipe 的单表查询不存在该问题。transformSql 变换Sql变换使用内存 SQL 引擎对输入行做变换。其三个必选参数为参数类型必填默认值说明plugin_inputstring是-输入表名query中的表名必须与之匹配plugin_outputstring是-输出表名querystring是-变换 SQL支持基础函数与条件过滤enginestring否ZETASQL 引擎可选ZETA与INTERNAL对应源码见 SQLTransform.javaquery无默认值必须配置engine默认取ZETA。变换在open()阶段通过 SQLEngineFactory 初始化 SQL 引擎并把plugin_input指定的表名与上游目录表绑定后执行查询。本示例的query依次完成四件事WHERE status PAID——过滤剔除订单1002CREATED状态UPPER(customer_name) AS customer_name——名字规范化CAST(amount AS DECIMAL(12, 2))——金额从 4 位小数改到 2 位小数MYSQL AS source_system——注入常量字段标识数据来源系统。其中FROM dual是默认 SeaTunnel SQL 变换引擎的虚拟输入表它并不指向 MySQL 或 PostgreSQL 中的任何真实表只是让 SELECT 语法成立请勿将其替换成真实表名。需要留意的是该变换引擎当前支持基础函数与条件过滤但尚不支持多源表 JOIN 与 AGGREGATE 聚合等复杂 SQLSQL 变换文档 有明确说明。sinkJdbc 目标JDBC Sink 文档 说明它支持批处理与流式作业、并行写入、生成式或自定义 SQL、多表写入与 CDC 事件。本段使用自定义 SQL 模式query INSERT INTO public.paid_orders (id, customer_name, amount, source_system) VALUES (?, ?, ?, ?)要点如下当未配置generate_sink_sql true时默认false必须提供queryquery中的?占位符按上游字段顺序绑定。本示例变换输出的列顺序是id, customer_name, amount, source_system与 INSERT 的四个占位符一一对应当前限制自定义query模式下不会执行 save mode 处理schema_save_mode、data_save_mode与custom_sql均不生效。如需自动建表与自动生成 SQL应改用generate_sink_sql truedatabasetable下文“进阶”小节会给出对照写法。字段顺序对齐最容易踩的坑原文特别强调变换输出的列顺序必须与 sinkquery的占位符顺序保持一致。例如把 SELECT 写成amount, id, ...而 INSERT 仍按(id, customer_name, amount, source_system)数据就会错位写入。这一点在自定义 SQL 模式下完全由用户负责SeaTunnel 不会做重排。运行任务替换配置中的主机名与凭据为你自己的环境值然后在本地模式运行cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/jdbc-to-jdbc.conf -m local-m local表示使用本地模式Zeta 引擎单进程运行适合开发与验证。验证结果作业结束后查询 PostgreSQL 目标表SELECT id, customer_name, amount, source_system FROM public.paid_orders ORDER BY id;预期结果idcustomer_nameamountsource_system1001ALICE CHEN120.50MYSQL1003CAROL WU42.00MYSQL这张结果表可以独立验证每一个变换动作订单1002因status CREATED被过滤未出现在结果中客户名被UPPER转为大写alice chen→ALICE CHEN金额精度从 4 位小数收敛为 2 位120.5000→120.50每行都带有常量字段source_system MYSQL。常见陷阱与规避驱动缺失或版本不兼容MySQL / PostgreSQL 驱动不在${SEATUNNEL_HOME}/lib或驱动版本与数据库、Java 运行时不兼容。按上文第 3 步放置并重启进程后再运行。示例主机名未替换示例中的mysql.example.com、postgresql.example.com只是占位符SeaTunnel 进程必须能解析并访问到你实际的主机名网络路由、防火墙、TLS、凭据都要核对。开发环境 URL 直接上生产示例 MySQL URL 关闭了 TLS 并允许公钥检索仅适用于简单开发环境生产环境务必按组织安全要求配置 TLS。字段顺序变更未同步到 INSERT变换输出的列顺序变化时sink 的INSERT占位符顺序必须同步调整。目标表非空导致主键冲突教程中已TRUNCATE目标表保证可重复运行生产环境应选择 upsert 策略如primary_keys 数据库原生 UPSERT。DECIMAL 精度溢出若源DECIMAL(16, 4)的真实数据超过目标DECIMAL(12, 2)可表达的精度/范围写入会失败。迁移前先评估真实数据范围选择合适的目标类型。批处理期间源数据变化批量作业只读取启动时刻可见的数据。若需要持续捕获源端变化应改用 CDC 源例如 Multi-Table CDC 配方。进阶同一管道的其他写法用generate_sink_sql替代自定义 query如果你希望 SeaTunnel 根据上游 schema 自动生成 INSERT / UPSERT 语句甚至自动建表可以去掉自定义query改为sink { Jdbc { plugin_input paid_orders driver org.postgresql.Driver url jdbc:postgresql://postgresql.example.com:5432/target_db username test password test generate_sink_sql true database target_db table public.paid_orders primary_keys [id] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }两种写入模式互斥要么generate_sink_sql truedatabase 通常的table要么generate_sink_sql falsequery。生成式模式配合primary_keys可以在 PostgreSQL 上生成INSERT ... ON CONFLICT (...) DO UPDATE的原生 upsert 语义。源端并行读取如果orders表数据量大可以为 JDBC Source 配置分区读取使用querypartition_columnpartition_num固定分区器或改用table_pathsplit.size走动态拆分依赖主键/唯一索引做分片键。启用并行时记得同步提升env.parallelism。分区的行为细节与边界语义参见 JDBC Source 的 Parallel Reader 章节。相关文档JDBC source 连接器文档SQL transform 文档JDBC sink 连接器文档Multi-Table CDC 配方需要持续捕获源端变化时的替代方案Recipe 总览按管道形态选择其他配方如 JDBC to S3、HTTP to JDBC【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表