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

SeaTunnel MySQL 到 HDFS 批量同步实战:JDBC 分片读取 + SQL 清洗 + Snappy Parquet 分区写入

SeaTunnel MySQL 到 HDFS 批量同步实战JDBC 分片读取 SQL 清洗 Snappy Parquet 分区写入【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇指南基于 SeaTunnel 官方 Recipe 文档完整演示一条MySQL 订单表 → HDFS 数仓明细层DWD的批量离线同步链路使用JDBC Source以固定分片split方式并行读取 MySQL 快照经过SQL Transform完成字段重命名、状态值大写归一化、负金额过滤与分区键生成最终由HdfsFile Sink以 Snappy 压缩的 Parquet 格式、按日期目录pt_dtyyyy-MM-dd落盘到 HDFS。读完本文你将掌握依赖与驱动的安装校验、JDBC 分片读参数的语义、SQL Transform 的内存 SQL 用法、HdfsFile 分区写入与保存模式的行为以及一套可复现的验证与排错方法。说明这是一个**批量快照batch snapshot**示例不是 CDC也不是增量同步作业。JDBC Source 是有界bounded数据源读取完查询可见的行后任务即结束若需要持续捕获后续的插入、更新、删除应改用 CDC 连接器。前置条件Prerequisites在开始前需要依次完成以下准备工作。本示例基于 SeaTunnel Zeta 引擎的本地模式local mode运行平台为 Linux。完成首个作业的部署与运行先阅读 Run your first job 并成功跑通一个本地作业。将SEATUNNEL_HOME设置为解压后的 SeaTunnel 发行版目录源码 checkout 不是必需的。部署细节见 Deployment。安装与发行版同版本的连接器在config/plugin_config中追加connector-jdbc与connector-file-hadoop保留你的环境需要的其他连接器条目然后执行插件安装脚本。plugin_config的内容如下--seatunnel-connectors-- connector-jdbc connector-file-hadoop --end--其中connector-file-hadoop内部实际复用了 connector-file-base 的通用文件读写实现并叠加了 Hadoop 文件系统适配是 HdfsFile 连接器的承载插件。放置 MySQL JDBC 驱动将驱动 JAR如mysql-connector-j-8.x.jar放入${SEATUNNEL_HOME}/lib然后校验连接器与驱动均已就位cd ${SEATUNNEL_HOME} sh bin/install-plugin.sh ls connectors | grep -E connector-(jdbc|file-hadoop) ls lib | grep mysql-connector从源码视角看connector-jdbc的驱动类加载与连接管理位于 seatunnel-connectors-v2/connector-jdbcSeaTunnel 出于 JDBC 驱动许可证与版本兼容性考虑默认不捆绑数据库驱动需要自行下载并放置。Zeta 引擎的驱动目录是${SEATUNNEL_HOME}/libSpark/Flink 引擎则为${SEATUNNEL_HOME}/plugins/Jdbc/lib/放置后需重启 SeaTunnel 进程使其生效详见 JDBC Source 文档。Zeta 发行版自带了 Hadoop 相关 JAR添加依赖前先检查lib目录不要混用任意版本的 Hadoop 客户端。HDFS 环境相关的特殊要求Kerberos、HA 等参见 HdfsFile sink。准备 MySQL 实例与账号需要一个可访问的 MySQL 实例以及一个对源表具备SELECT权限的作业账号。用于执行下方种子 SQL 的设置账号还需具备建库、建表与插入权限作业账号不需要这些设置类权限。准备 HDFS 集群与输出目录需要一个可达的 HDFS 集群和一个未使用过的输出目录。SeaTunnel 进程必须对输出目录以及 sink 的临时目录默认/tmp/seatunnel具备写权限。示例假定为非 Kerberos 的 HDFSKerberos 或 HA 场景请按 HdfsFile 文档补充hdfs_site_path、krb5_path、kerberos_principal、kerberos_keytab_path等选项。用于校验的 Hadoop CLI 也必须能访问该集群。准备源数据Prepare source data:::caution 使用隔离的测试数据库请使用设置账号在 MySQL 客户端中一次性执行下面的 SQL。它刻意使用不带IF NOT EXISTS的CREATE DATABASE并且不 drop、不 truncate 任何表。如果trade_db已经存在请停止并换一个未使用的测试库名并在 SQL、JDBC URL、table_path与查询中保持一致替换。不要让 SQL 客户端在出错后继续执行。:::CREATE DATABASE trade_db; USE trade_db; CREATE TABLE orders ( id BIGINT NOT NULL PRIMARY KEY, order_no VARCHAR(64) NOT NULL, user_id BIGINT NOT NULL, amount DECIMAL(10, 2) NOT NULL, status VARCHAR(32) NOT NULL, create_time DATETIME NOT NULL ); INSERT INTO orders (id, order_no, user_id, amount, status, create_time) VALUES (1, ORD-20260823-001, 10001, 99.50, completed, 2026-08-23 10:15:30), (2, ORD-20260823-002, 10002, 199.00, pending, 2026-08-23 14:20:00), (3, ORD-20260824-001, 10003, 49.90, COMPLETED, 2026-08-24 09:00:15), (4, ORD-20260824-002, 10001, 350.00, paid, 2026-08-24 18:45:10), (5, ORD-20260824-003, 10004, -10.00, cancelled, 2026-08-24 20:00:00);这 5 条订单横跨两个日期并故意包含一条负金额记录id5用于演示过滤效果。运行 SeaTunnel 前请为作业账号配置好对该测试表的访问权限。完整配置Complete configuration将下面的配置保存为发行版下的config/mysql-to-hdfs.conf。其中 JDBC 的 host、库名、username、password需要替换为你测试环境的值fs.defaultFS与path需要替换为你的 HDFS 地址和未使用的测试输出路径。这里的localhost指运行 SeaTunnel 的机器namenode必须能从该机器解析。示例凭据均为占位符本隔离示例之外请按你环境的要求配置 TLS。env { job.name mysql_to_hdfs_batch_dw job.mode BATCH parallelism 4 } source { Jdbc { plugin_output src_mysql_orders url jdbc:mysql://localhost:3306/trade_db?useSSLfalseserverTimezoneUTCrewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username test_user password test_password table_path trade_db.orders query select id, order_no, user_id, amount, status, create_time, date(create_time) as create_date from trade_db.orders partition_column id partition_num 4 partition_lower_bound 1 partition_upper_bound 10000000 fetch_size 2000 } } transform { Sql { plugin_input src_mysql_orders plugin_output dwd_orders query select id as order_id, order_no, user_id, amount, upper(status) as order_status, create_time, FORMATDATETIME(create_date, yyyy-MM-dd) as pt_dt from src_mysql_orders where amount 0 } } sink { HdfsFile { plugin_input dwd_orders fs.defaultFS hdfs://namenode:8020 path /user/hive/warehouse/dwd.db/dwd_orders_df file_format_type parquet partition_by [pt_dt] compress_codec snappy schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }配置语义逐项解读env 块job.mode BATCH声明批处理语义parallelism 4决定任务的整体并行度它将影响源端并发读与 sink 端并发写。source 块JDBC Sourceplugin_output src_mysql_orders定义下游可见的表名是 SQL Transform 中from子句引用的来源。url/driver连接串与驱动类名MySQL 驱动类为com.mysql.cj.jdbc.Driver。URL 中rewriteBatchedStatementstrue有助于提升批量写回性能本示例为纯读保留该参数亦无副作用。table_path trade_db.orders给出表全路径与query可以同时使用——既显式声明表身份与元数据来源又由query控制投影列与库侧表达式。query在 MySQL 侧求值date(create_time)得到create_date供下游 SQL Transform 使用。分片读参数partition_column id指定分片键partition_num 4指定固定分片数partition_lower_bound/partition_upper_bound指定闭区间下界/上界。从源码 JdbcSourceOptions.java 可以看到这四个参数分别对应PARTITION_COLUMN、PARTITION_NUM、PARTITION_LOWER_BOUND、PARTITION_UPPER_BOUND四个 Option 定义。只有顶层querypartition_column的组合才会选择固定分片器fixed splitterpartition_num此时控制分片个数分片键必须出现在查询结果中。上下界是可选的——省略时 SeaTunnel 会额外查询MIN/MAX显式给出可以省去这两次查询但错误的边界会导致源行被遗漏所以只在完全掌握数据范围时使用。fetch_size 2000查询的行抓取大小用于减少大结果集下访问数据库的次数以提升性能0表示使用 JDBC 默认值见 JdbcSourceOptions.java。注意本示例的边界值与partition_num 4是为了演示固定分片读配置不是针对 5 行数据的性能调优。放大到真实数据时再适配这些参数极小的数据集不必用完所有 writer也不应期望产生等大小的文件。transform 块SQL TransformSeaTunnel 的 SQL Transform 使用内存 SQL 引擎默认engine ZETAplugin_input必须与上游plugin_output一致、plugin_output必须与下游plugin_input一致query中的表名必须等于plugin_input。其完整选项表见 SQL transform。这里的 SQL 完成四件事id as order_id字段重命名upper(status) as order_status订单状态归一化为大写示例数据中completed、pending、paid、cancelled与COMPLETED都会被统一为大写where amount 0过滤负金额剔除 id5 的-10.00订单FORMATDATETIME(create_date, yyyy-MM-dd) as pt_dt用 SQL 函数 中的FORMATDATETIME(dateAndTime, formatString) - STRING生成pt_dt分区键格式字符含义遵循java.time.format.DateTimeFormattery年、M月、d日、H时、m分、s秒。sink 块HdfsFile Sink完整的选项表见 HdfsFile sink本示例核心参数的行为如下fs.defaultFSHadoop 集群地址标准 HDFS 形如hdfs://hadoopcluster或hdfs://namenode:8020亦支持 ViewFSviewfs://mycluster。path目标目录必须填写。tmp_path默认/tmp/seatunnel结果先写入临时路径再通过mv提交到目标目录——这就是本示例要求进程对临时目录有写权限的原因。file_format_type parquet支持text、csv、parquet、orc、json、excel、xml、binary等格式。partition_by [pt_dt]按选定字段分区同时要求have_partition生效默认partition_dir_expression为${k0}${v0}/...即生成pt_dt2026-08-23/这样的 Hive 风格目录默认is_partition_field_write_in_file false即分区字段不写入数据文件内部。compress_codec snappyparquet 格式支持lzo、snappy、lz4、gzip、brotli、zstd、none。schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST目录不存在则创建已存在则跳过。data_save_mode APPEND_DATA保留目录、保留已有数据文件——这也是重复运行会追加重复数据的原因。数据流与调用链整条链路可以概括为Jdbc(plugin_outputsrc_mysql_orders) → Sql(plugin_inputsrc_mysql_orders, plugin_outputdwd_orders) → HdfsFile(plugin_inputdwd_orders)。每个插件通过plugin_input/plugin_output以显式命名的方式首尾相接配置中必须保持命名一致。从源码结构看sink 侧的写策略与配置校验集中在 connector-file-base含file_format_type、compress_codec、partition_by、schema_save_mode、data_save_mode等选项的声明Hadoop 文件系统适配则位于 connector-file-hadoopZeta 引擎下 sink 默认通过 2PC 提交保证 exactly-once文件先落在tmp_pathcheckpoint 成功后统一 mv 到目标目录。运行作业Run the job首次运行前先在 MySQL 中记录期望行数SELECT DATE(create_time) AS pt_dt, COUNT(*) AS expected_rows FROM trade_db.orders WHERE amount 0 GROUP BY DATE(create_time) ORDER BY pt_dt;针对给定的种子数据预期每个日期 2 行、共 4 行。然后提交作业cd ${SEATUNNEL_HOME} ./bin/seatunnel.sh --config ./config/mysql-to-hdfs.conf -m local等待作业成功结束后再检查输出。:::caution 重复运行会追加数据本作业使用APPEND_DATA对同一源表和同一输出路径重复运行可能产生重复记录。4 行的预期只适用于对未使用过的输出目录成功运行一次的场景。若需再次验证请使用新的测试输出路径不要删除或覆盖已有的数仓数据。:::验证结果Validation result检查分区目录使用你实际的 HDFS 地址与输出路径执行hdfs dfs -ls -R hdfs://namenode:8020/user/hive/warehouse/dwd.db/dwd_orders_df预期目录结构如下/user/hive/warehouse/dwd.db/dwd_orders_df/ ├── pt_dt2026-08-23/ │ └── generated-file.parquet └── pt_dt2026-08-24/ └── generated-file.parquet上面的树是示意性的实际数据文件的名称与数量取决于 writer 实例数、并行度与数据分布。仅凭文件名无法证明压缩编码器——压缩格式必须通过 Parquet 元数据确认。检查记录与文件格式使用你环境中已有的、支持 Parquet 的读取器读取两个分区目录下的全部已提交数据文件。对于给定的种子数据与单次运行应验证以下值行顺序不保证order_idorder_statusamountpt_dt1COMPLETED99.502026-08-232PENDING199.002026-08-233COMPLETED49.902026-08-244PAID350.002026-08-24同时验证order_no、user_id、create_time均被保留。订单 5必须缺失因为其金额为负被 SQL Transform 过滤。检查 Parquet 元数据以确认 Snappy 压缩不要用cat解读二进制文件。两点补充提醒pt_dt的值对应目录分区视读取器的不同可能需要从目录名推断而非从单个文件内读取该列。Hive 风格的目录不会自动创建或注册 Hive 表。目录存在本身不等于数据验证完成——务必核对行内容与文件格式。常见坑Common pitfallsMySQL JDBC 驱动缺失于 SeaTunnel 进程的lib目录或连接器版本与发行版不匹配驱动放置与版本要求见 JDBC Source 文档。localhost指向了错误的机器或 HDFS NameNode/DataNode 主机名从 SeaTunnel 所在机器不可达。作业账号能连 MySQL 但读不到trade_db.orders或 HDFS 身份无法写输出目录与临时文件默认/tmp/seatunnel。对已存在的数据库或输出目录运行本教程导致种子数据或期望行数失效。修改源端query时删掉了create_date而 SQL Transform 仍引用它——两端字段必须保持一致。把 Parquet 文件当作纯文本处理或假定固定的文件名与文件数量——实际文件命名受事务与并行写影响。相关文档Related docsJDBC sourceSQL transformSQL functionsHdfsFile sink【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

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

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