SeaTunnel JDBC MySQL 源连接器完全指南:参数配置、并行读取拆分与多表同步实战

📅 发布时间:2026/9/19 1:26:26
SeaTunnel JDBC MySQL 源连接器完全指南:参数配置、并行读取拆分与多表同步实战
SeaTunnel JDBC MySQL 源连接器完全指南参数配置、并行读取拆分与多表同步实战【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel导读本文以 SeaTunnel 官方文档中 MySQL 源连接器JdbcSource基于 JDBC 驱动为骨架系统讲解如何通过 SeaTunnel 从 MySQL 5.58.4 读取数据。你将掌握驱动依赖的安装位置Spark/Flink 与 Zeta 引擎的差异、完整的连接参数与拆分策略参数含义、MySQL 到 SeaTunnel 的类型映射规则、并行读取的拆分键选择与分片算法以及单表查询、分区并行、按主键/唯一索引并行、指定边界、多表读取五类可直接复用的作业配置。文中同时结合仓库内 connector-jdbc 的源码实现说明参数校验规则与底层分片执行流程帮助你在生产环境中做出正确的配置决策。连接器概述MySQL 源连接器通过 JDBC 读取外部数据源数据在 SeaTunnel 配置中以Jdbc作为插件名对应源码中 JdbcSourceFactory 的factoryIdentifier()返回值。支持的 MySQL 版本5.5 / 5.6 / 5.7 / 8.0 / 8.1 / 8.2 / 8.3 / 8.4支持的引擎Spark、Flink、SeaTunnel Zeta即该连接器可运行在三种引擎之上文档示例默认使用 Zeta 引擎的job.mode与parallelism配置语法。依赖安装JDBC 驱动 jar 包的放置位置使用 MySQL 源连接器前必须保证 MySQL JDBC 驱动mysql-connector-java可被加载。官方文档按引擎给出两种放置路径Spark / Flink 引擎将驱动 jar 包放置在${SEATUNNEL_HOME}/plugins/目录中SeaTunnel Zeta 引擎将驱动 jar 包放置在${SEATUNNEL_HOME}/lib/目录中。从源码看驱动加载发生在 JdbcSource 构造阶段通过Class.forName(driverName)将驱动注册到DriverManager若加载失败会记录 WARN 日志。这也是文档强调确保 jar 包已放置在正确目录的原因——驱动缺失会在作业启动阶段直接暴露。功能特性总览特性支持情况批处理✅ 支持流处理❌ 不支持精确一次✅ 支持列投影✅ 支持并行度✅ 支持支持用户定义的拆分✅ 支持支持多表读取✅ 支持支持 SQL 查询并能实现列投影效果。源码层面JdbcSource 实现了SupportParallelism与SupportColumnProjection两个接口分别对应并行度与列投影能力SupportColumnProjection允许引擎在 SQL 查询或table_path/table_list场景下只投影所需列减少网络与内存开销。数据源信息数据源支持的版本驱动器网址Maven 下载链接Mysql不同的依赖版本具有不同的驱动程序类com.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/test从 Maven 仓库下载 mysql-connector-java连接 URL 的标准格式为jdbc:mysql://localhost:3306/test其中3306为 MySQL 默认端口test为数据库名。实际使用时通常还会附加serverTimezone、useUnicode、characterEncoding、rewriteBatchedStatements等参数详见后文示例。MySQL 数据类型映射读取 MySQL 表时SeaTunnel 会依据下表将 MySQL 列类型转换为 SeaTunnel 数据类型。该映射直接影响下游 Transform 与 Sink 所接收到的数据类型是排查类型不符类问题的第一现场。MySQL 数据类型SeaTunnel 数据类型BIT(1)、TINYINT(1)BOOLEANTINYINTBYTETINYINT UNSIGNED、SMALLINTSMALLINTSMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINTINT UNSIGNED、INTEGER UNSIGNED、BIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20,0)DECIMAL(x,y)获取指定列的列大小 38DECIMAL(x,y)DECIMAL(x,y)获取指定列的列大小 38DECIMAL(38,18)DECIMAL UNSIGNEDDECIMAL(列大小1, 小数点右侧位数)FLOAT、FLOAT UNSIGNEDFLOATDOUBLE、DOUBLE UNSIGNEDDOUBLECHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、JSON、ENUMSTRINGDATEDATETIME(s)TIME(s)DATETIME、TIMESTAMP(s)TIMESTAMP(s)TINYBLOB、MEDIUMBLOB、BLOB、LONGBLOB、BINARY、VARBINAR、BIT(n)、GEOMETRYBYTES需要注意TINYINT(1)映射为BOOLEAN的行为受int_type_narrowing参数控制详见下文专题小节BIT(1)则始终映射为BOOLEAN。数据源参数详解以下为官方文档给出的全部源参数结合 JdbcSourceOptions 与 JdbcCommonOptions 源码中的定义与默认值整理名称类型是否必填默认值描述urlString是-JDBC 连接的 URL。示例jdbc:mysql://localhost:3306/testdriverString是-用于连接远程数据源的 JDBC 类名MySQL 下为com.mysql.cj.jdbc.DriverusernameString否-连接实例用户名passwordString否-连接实例密码queryString否-查询语句。当未配置table_path和table_list时必填connection_check_timeout_secInt否30验证数据库连接所使用的操作完成的等待时间秒partition_columnString否-用于并行度分区的列名仅支持数字类型主键并且只能配置一列partition_lower_boundBigDecimal否-扫描时partition_column的最小值未设置时 SeaTunnel 查询数据库获取最小值partition_upper_boundBigDecimal否-扫描时partition_column的最大值未设置时 SeaTunnel 查询数据库获取最大值partition_numInt否作业并行度分区数量仅支持正整数默认值为作业并行度fetch_sizeInt否0查询的行提取大小通过减少满足选择条件所需的数据库访问次数来提高性能0 表示使用 JDBC 默认值propertiesMap否-额外的连接配置参数与 URL 中参数冲突时优先级由驱动实现决定MySQL 中属性优先于 URLuse_regexBoolean否false控制表路径的正则表达式匹配true 时table_path被视为正则表达式模式table_pathString否-表的完整路径可用此配置代替query。示例testdb.table1table_listArray否-要读取的表的列表可用此配置代替table_pathwhere_conditionString否-所有表/查询的通用行过滤条件必须以where开头例如where id 100split.sizeInt否8096表的分割大小行数读取时捕获的表会被分割成多个分片split.even-distribution.factor.lower-boundDouble否0.05分片键分布因子的下限用于判断表数据分布是否均匀split.even-distribution.factor.upper-boundDouble否100分片键分布因子的上限split.sample-sharding.thresholdInt否1000触发样本分片策略的估算分片数阈值split.allow-samplingBoolean否true是否允许对分布不均匀的分片键使用采样分片策略false 时退回迭代式不均匀分片use_select_countBoolean否false是否使用select count(*)在分片前估算表行数skip_analyzeBoolean否false是否跳过分片前的表行数分析split.inverse-sampling.rateInt否1000样本分片策略中采样率的倒数1000 表示应用 1/1000 的采样率int_type_narrowingBoolean否trueInt 类型收窄true 时tinyint(1)在无精度损失情况下收窄为 boolean目前仅支持 MySQLcommon-options-否-源插件的常见参数详见源常见参数参数校验规则源码级除参数表本身外连接器还在提交阶段对参数组合做静态校验相关逻辑位于 JdbcSourceFactory必填项url与driver为必填optionRule()中通过required(URL, DRIVER)声明互斥校验TableListExclusiveValidator保证table_list与table_path/query互斥即一张表选择模式只能二选一不能同时配置JdbcSourceFactory.java#L151-L163前缀校验WhereConditionPrefixValidator要求where_condition必须小写不敏感地以where开头避免拼接出畸形 SQLJdbcSourceFactory.java#L169-L179表路径唯一性多表读取时JdbcSourceTableConfig.of() 会校验table_path不允许重复或为空否则抛出IllegalArgumentException单表默认分区数JdbcSourceTableConfig中内置DEFAULT_PARTITION_NUMBER 10当单表未显式指定partition_num时采用该值。int_type_narrowing 专题Int 类型收窄当值为true时tinyint(1)类型在无精度损失的情况下被收窄为boolean类型目前仅支持 MySQL。源码定义见 JdbcCommonOptions.java#L107-L112int_type_narrowing trueMySQLSeaTunnelTINYINT(1)Booleanint_type_narrowing falseMySQLSeaTunnelTINYINT(1)TINYINT该开关的意义在于MySQL 中TINYINT(1)常被用作布尔标志位开启收窄后下游可直接按BOOLEAN处理若你的业务把TINYINT(1)当作数值使用如 0/1 之外的取值应关闭此开关以避免精度损失。其他源码中可用的拆分/读取选项从 JdbcSourceOptions 可以看到除文档参数表外源码还提供以下进阶选项以当前仓库实现为准split.string_split_mode默认sample字符串拆分算法可切换为charset_based以启用基于字符集的字符串拆分split.string-strategyString 分区列的拆分策略range范围拆分、hash哈希拆分、none禁用、auto优先 range必要时回退 hashrange/auto目前要求 MySQL 二进制排序规则与定长可打印 ASCII 键值split.string_split_mode_collatecharset_based模式下可指定排序规则缺省使用数据库默认排序规则enable_concurrent_read默认true快照阶段是否启用基于拆分的高并发读取设为false时源以单个分片读取整张表不做任何分片分析适合无索引的表。并行读取与拆分机制JDBC 源连接器支持从表中并行读取数据。SeaTunnel 使用特定规则将表中的数据分割再交给多个读取器并行消费读取器的数量由parallelism选项决定。拆分键规则如果partition_column不为空则使用它计算数据分片该列必须属于支持的分片数据类型如果partition_column为空SeaTunnel 从表中读取模式并获取主键和唯一索引。若主键或唯一索引包含多列则选择第一个属于支持的分片数据类型的列来分片。例如主键为(guid, name varchar)时因为guid不属于支持的分片数据类型会退而选择name列进行分片。支持的拆分数据类型StringNumberint、bigint、decimal 等Date与拆分相关的参数split.size每个分片中的行数捕获的表在读取时会被分成多个分片默认 8096。split.even-distribution.factor.lower-bound不推荐使用分片键分布因子的下限。该因子用于判断表数据是否均匀分布若计算出的分布因子大于等于此下限即(MAX(id) - MIN(id) 1) / 行数表的分片将被优化为均匀分布否则视为不均匀分布。若估算的分片数量超过sample-sharding.threshold将使用基于采样的分片策略。默认值 0.05。split.even-distribution.factor.upper-bound不推荐使用分片键分布因子的上限判断规则同上分布因子小于等于此上限(MAX(id) - MIN(id) 1) / 行数时采用均匀分布优化否则视为不均匀。默认值 100.0。split.sample-sharding.threshold触发采样分片策略的估算分片数量阈值。当分布因子超出上下限范围且估算分片数大致行数 / 分片大小超过该阈值时使用采样分片策略以更高效地处理大数据集。默认值 1000 个分片。split.inverse-sampling.rate采样分片策略中采样率的倒数。值 1000 表示采样过程中应用 1/1000 的采样率处理非常大的数据集时通常首选较低采样率。默认值 1000。partition_column [string]拆分数据的列名称。partition_upper_bound [BigDecimal]扫描时partition_column的最大值。未设置时 SeaTunnel 查询数据库获取最大值。partition_lower_bound [BigDecimal]扫描时partition_column的最小值。未设置时 SeaTunnel 查询数据库获取最小值。partition_num [int]不推荐使用正确做法是通过split.size控制分片数量。需要拆分成多少个分片仅支持正整数默认值为作业并行度。分片执行流程源码级从源码可以完整还原并行读取的实现链路ChunkSplitter.create() 根据配置选择DynamicChunkSplitter动态拆分或FixedChunkSplitter固定拆分若enable_concurrent_read为falsegenerateSplits() 直接返回单个全表分片跳过昂贵的 MIN/MAX 扫描否则通过findSplitKey查找分片键找不到时同样退化为单分片拆分出的分片进入 JdbcSourceSplitEnumerator 的pendingTables/pendingSplits队列run()循环逐表调用generateSplits并按HashUtils.bucketIndex(splitId.hashCode(), readerCount)将分片哈希分配到各读取器JdbcSourceSplitEnumerator.java#L176-L178实现负载均衡分片快照状态snapshotState可支持精确一次语义下的故障恢复。任务示例以下示例均为可直接套用的 HOCON 配置运行时环境统一由env块声明。简单的例子以单线程并行方式查询测试数据库中type_bin为 table 的 16 条数据并查询所有字段也可以指定查询字段最终结果输出到控制台# 定义运行时环境 env { parallelism 4 job.mode BATCH } source{ Jdbc { url jdbc:mysql://localhost:3306/test?serverTimezoneGMT%2b8useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 query select * from type_bin limit 16 } } transform { # 如需了解 Transform 插件的完整配置可参考仓库中的 SQL Transform 文档 } sink { Console {} }说明URL 中serverTimezoneGMT%2b8指定时区为东八区useUnicode/characterEncodingUTF-8保证中文不乱码rewriteBatchedStatementstrue可优化批量写入对读取场景影响有限按需保留。按partition_column并行env { parallelism 4 job.mode BATCH } source { Jdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%2b8 driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 query select * from type_bin partition_column id split.size 10000 # Read start boundary #partition_lower_bound ... # Read end boundary #partition_upper_bound ... } } sink { Console {} }partition_column必须为数字类型且只能配置一列配合split.size 10000控制每个分片的行数从而间接决定分片数量。按主键或唯一索引并行配置table_path将启用自动拆分可配置split.*调整拆分策略。env { parallelism 4 job.mode BATCH } source { Jdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%2b8 driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 table_path testdb.table1 query select * from testdb.table1 split.size 10000 } } sink { Console {} }此场景下无需显式指定partition_columnSeaTunnel 会自动从主键或唯一索引中选取支持的分片数据类型列进行拆分。并行的同时指定边界指定数据的上下边界查询会更加高效根据配置的上下边界读取数据源。source { Jdbc { url jdbc:mysql://localhost:3306/test?serverTimezoneGMT%2b8useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 # Define query logic as required query select * from type_bin partition_column id # Read start boundary partition_lower_bound 1 # Read end boundary partition_upper_bound 500 partition_num 10 properties { useSSLfalse } } }partition_lower_bound/partition_upper_bound显式给出扫描边界后SeaTunnel 不再需要查询MIN(id)/MAX(id)可减少一次数据库往返properties中的useSSLfalse用于关闭 SSL 连接按安全要求谨慎使用。多表读取配置table_list将启用自动拆分可配置split.*调整拆分策略。env { job.mode BATCH parallelism 4 } source { Jdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%2b8 driver com.mysql.cj.jdbc.Driver connection_check_timeout_sec 100 username root password 123456 table_list [ { table_path testdb.table1 }, { table_path testdb.table2 # Use query filter rows columns query select id, name from testdb.table2 where id 100 } ] #where_condition where id 100 #split.size 8096 #split.even-distribution.factor.upper-bound 100 #split.even-distribution.factor.lower-bound 0.05 #split.sample-sharding.threshold 1000 #split.inverse-sampling.rate 1000 } } sink { Console {} }要点table_list中每项由table_path标识可选的query只作用于对应单表如示例中对table2做行/列过滤全局的where_condition以where开头会作为所有表/查询的通用过滤条件拆分相关参数split.size、split.even-distribution.factor.*、split.sample-sharding.threshold、split.inverse-sampling.rate可按需放开并作用于所有表多张表的table_path必须唯一否则作业提交失败见上文JdbcSourceTableConfig校验逻辑从变更记录看2.3.12 版本起还支持通过正则表达式读取多张表配合use_regex使用。使用提示与注意事项如果表无法拆分例如表没有主键或唯一索引且未设置partition_column则将以单线程并发方式运行。使用table_path替代query进行单表读取如需读取多个表请使用table_list。当基于query推断主键时主键继承自结果集中第一列所在的底层表如果query包含多表 JOIN 或同时从多张表读取该主键对整个 JOIN 结果集的唯一性不作严格保证。这三点提示对应源码中的两条退化路径enable_concurrent_readfalse或找不到可用分片键时都会退化为单个全表分片ChunkSplitter.java#L108-L120此时并行度无法生效性能受限于单连接读取务必提前为表设计主键/唯一索引或配置partition_column。变更日志与进一步阅读连接器完整变更记录见 connector-jdbc 变更日志其中与 MySQL 源读取相关的代表性改进包括tinyint(1)到 byte/tinyint 的读取支持、MySQLtinyint(1)类型映射修复、基于字符集的 String 列拆分支持、通过正则表达式读取多张表2.3.12等源插件通用参数如result_table_name、parallelism等参见源常见参数连接器实现源码位于 connector-jdbc核心类包括参数定义 JdbcSourceOptions、工厂与校验 JdbcSourceFactory、分片器 ChunkSplitter含DynamicChunkSplitter/FixedChunkSplitter两种实现以及分片调度器 JdbcSourceSplitEnumerator。在实际使用中建议遵循单表优先table_path、多表用table_list、复杂查询用query、大表务必保证可分片的选型原则并结合 MySQL 数据分布情况调整split.size与采样相关参数即可在批处理场景下获得稳定、可预期的并行读取效果。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考