- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本文以 Apache SeaTunnel 的 JDBC Source 连接器(插件标识Jdbc)为核心,系统讲解如何通过 JDBC 从 MySQL、PostgreSQL、Oracle、SQL Server 等关系型数据库批量读取数据:从驱动部署、完整参数表,到partition_column与split.*两种并行分片机制,再到单表/多表/带边界条件的真实配置示例。读完本文,你将能独立编写一套可运行的 JDBC Source 作业,并理解其"自动分片、并行拉取"的底层实现逻辑。
JDBC Source 是什么
JDBC Source 是 SeaTunnel 中用于读取外部数据源的批式(Batch)输入插件。它不关心具体是哪种数据库,只要该数据库提供 JDBC 驱动,就可以通过统一的url + driver + query/table_path三要素接入,读取整表或执行任意查询语句得到的结果集。
插件的能力矩阵如下(能力定义可参考 connector-v2-features):
| 能力 | 支持情况 |
|---|---|
| batch(批式) | ✅ 支持 |
| stream(流式) | ❌ 不支持 |
| exactly-once(精确一次) | ✅ 支持 |
| column projection(列投影) | ✅ 支持,通过查询 SQL 即可实现投影效果 |
| parallelism(并行度) | ✅ 支持 |
| support user-defined split(用户自定义分片) | ✅ 支持 |
| support multiple table read(多表读取) | ✅ 支持 |
在源码层面,该插件的入口为 JdbcSourceFactory,其factoryIdentifier()返回"Jdbc",与配置文件中source { Jdbc { ... } }的插件名一一对应;实际读取逻辑封装在 JdbcSource 中,其getBoundedness()返回BOUNDED,印证了这是一个纯批式(有界)数据源。
环境准备:数据库驱动必须自行提供
SeaTunnel 出于 License 合规考虑,不内置任何数据库驱动,需要使用者自行准备并拷贝到指定目录:
- 使用 SeaTunnel Zeta 引擎时,将驱动 jar 拷贝到
${SEATUNNEL_HOME}/lib/目录; - 使用 Spark/Flink 引擎时,将驱动 jar 拷贝到
${SEATUNNEL_HOME}/plugins/目录(Spark 还需放入$SPARK_HOME/jars/,Flink 需放入$FLINK_HOME/lib/)。
以 MySQL 为例,需要下载并拷贝mysql-connector-java-xxx.jar。其他数据库对应的驱动类名与连接 URL 可参考下文附录表,各驱动 jar 均可从其对应厂商的 Maven 中央仓库坐标(例如mysql:mysql-connector-java、org.postgresql:postgresql等)获取,注意版本需与目标数据库服务端兼容。
选项(Options)全解
以下为 JDBC Source 的完整参数表,均已在 JdbcSourceFactory.optionRule() 中注册,其中url与driver为必填,其余均为可选:
| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接地址,例如jdbc:postgresql://localhost/test |
| driver | String | 是 | - | 连接远程数据源的 JDBC 驱动类名,如 MySQL 为com.mysql.cj.jdbc.Driver |
| user | String | 否 | - | 用户名 |
| password | String | 否 | - | 密码 |
| query | String | 否 | - | 查询语句 |
| compatible_mode | String | 否 | - | 数据库兼容模式,当数据库支持多种兼容模式时必须设置,如 OceanBase 需设置为mysql或oracle |
| connection_check_timeout_sec | Int | 否 | 30 | 用于校验连接的数据操作最长等待秒数 |
| partition_column | String | 否 | - | 用于数据分片的列名 |
| partition_upper_bound | Long | 否 | - | 扫描的partition_column最大值,不设置时 SeaTunnel 会查询数据库获取最大值 |
| partition_lower_bound | Long | 否 | - | 扫描的partition_column最小值,不设置时 SeaTunnel 会查询数据库获取最小值 |
| partition_num | Int | 否 | 作业并行度 | 分片数量(不推荐使用),正确做法是通过split.size控制分片数,仅支持正整数 |
| use_select_count | Boolean | 否 | false | 动态分片阶段是否用select count统计行数(当前仅 Oracle 可用),当用 analyze 语句更新统计信息较慢时可直接使用 select count |
| skip_analyze | Boolean | 否 | false | 动态分片阶段跳过表行数分析(当前仅 Oracle 可用),适用于已定时执行 analyze 或表数据变化不频繁的场景 |
| fetch_size | Int | 否 | 0 | 查询返回大量对象时配置每次抓取的行数以提升性能,0 表示使用 JDBC 默认值 |
| properties | Map | 否 | - | 附加连接参数;当 properties 与 url 中参数同名时,优先级由驱动实现决定,如 MySQL 中 properties 优先于 URL |
| table_path | String | 否 | - | 表的完整路径,可替代query。示例:MySQL"testdb.table1"、Oracle"test_schema.table1"、SQL Server"testdb.test_schema.table1"、PostgreSQL"testdb.test_schema.table1"、InterSystems IRIS"test_schema.table1" |
| table_list | Array | 否 | - | 要读取的表列表,可替代table_path |
| where_condition | String | 否 | - | 作用于所有表/查询的公共行过滤条件,必须以where开头,例如where id > 100 |
| split.size | Int | 否 | 8096 | 每个分片包含的行数,读取表时按该值切分 |
| split.even-distribution.factor.lower-bound | Double | 否 | 0.05 | (不推荐使用)分片键分布因子下界,用于判断表数据是否均匀分布。分布因子 = (MAX(id) − MIN(id) + 1) / 行数,若大于等于该下界则按均匀分布优化分块,否则视为非均匀分布,在预估分片数超过sample-sharding.threshold时采用采样分片策略 |
| split.even-distribution.factor.upper-bound | Double | 否 | 100 | (不推荐使用)分片键分布因子上界,分布因子小于等于该上界时按均匀分布优化分块,否则视为非均匀分布并可能触发采样分片 |
| split.sample-sharding.threshold | Int | 否 | 1000 | 触发采样分片策略的预估分片数阈值。当分布因子落在上下界之外,且预估分片数(近似行数 / chunk 大小)超过该阈值时,启用采样分片策略,以更高效地处理大数据集 |
| split.inverse-sampling.rate | Int | 否 | 1000 | 采样分片策略中采样率的倒数。例如设为 1000 表示按 1/1000 采样,用于控制采样粒度、影响最终分片数,非常适合超大数据集 |
| common-options | - | 否 | - | Source 插件通用参数,详见 Source Common Options |
参数背后的源码实现
- 默认值一致性:上述
connection_check_timeout_sec = 30、fetch_size = 0定义于 JdbcOptions;split.size = 8096、split.even-distribution.factor.lower-bound = 0.05、split.even-distribution.factor.upper-bound = 100、split.sample-sharding.threshold = 1000、split.inverse-sampling.rate = 1000、use_select_count = false、skip_analyze = false定义于 JdbcSourceOptions,与文档表格完全对应,可放心按默认值使用。 - where_condition 校验:在 JdbcSourceConfig.of() 中,
where_condition被强制要求以小写where开头,否则抛出IllegalArgumentException;实际生成 SQL 时它会以SELECT * FROM (子查询) tmp where ...的形式包裹在查询外层。 - table_list 与 query/table_path 互斥:在 JdbcSourceTableConfig.of() 中,若同时配置了
table_list和query/table_path,会直接报错;且当表数量大于 1 时,要求table_path非空且不重复。 - properties 处理:
properties会被解析为Map<String, String>(见 JdbcConnectionConfig),在连接建立时合并进 JDBC URL;compatible_mode用于 OceanBase 这类多兼容模式数据库,帮助JdbcDialectLoader选择正确的方言实现。
并行读取与分片(Split)机制
JDBC Source 支持从表中并行读取数据。SeaTunnel 会依据一定规则将表数据切分成多个 Split,交给多个 Reader 并行消费;Reader 的数量由作业parallelism决定。分片键(Split Key)的选择规则如下:
- 若配置了
partition_column,则直接用该列参与分片计算,该列必须属于支持的分片数据类型; - 若未配置
partition_column,SeaTunnel 会读取表结构(Schema),依次查找主键(Primary Key)和唯一索引(Unique Index),取其中第一个属于支持的分片数据类型的列作为分片键。例如某表主键为(nn guid, name varchar),因guid不在支持类型内,会退而选择name列参与分片。
支持的分片数据类型:
- String(字符串)
- Number(int、bigint、decimal 等数值类型)
- Date(日期)
这一规则在 ChunkSplitter.findSplitKey() 中完整实现:先校验partition_column是否存在且属于支持类型,再遍历主键列、唯一索引列;支持类型(TINYINT/SMALLINT/INT/BIGINT/DOUBLE/FLOAT/DECIMAL/STRING/DATE)由isSupportSplitColumn()判定。
固定分片(Fixed)与动态分片(Dynamic)
源码中分片器由 ChunkSplitter.create() 统一创建,选择逻辑非常关键(见 JdbcSourceConfig.of()):
- 旧版固定分片(FixedChunkSplitter):当同时配置了
query和partition_column时,走固定分片路径,按partition_lower_bound、partition_upper_bound、partition_num等差数列切分; - 动态分片(DynamicChunkSplitter):其余情况(即未同时提供
query+partition_column,例如使用table_path/table_list)默认启用动态分片,依据split.*系列参数自适应切分。
动态分片的核心流程(见 DynamicChunkSplitter.splitTableIntoChunks()):
- 查询分片键列的
MIN/MAX值;若表为空或只有一行,则整表作为一个 Chunk; - 查询表的近似行数(
queryApproximateRowCnt); - 计算分布因子
(MAX − MIN + 1) / 行数(见calculateDistributionFactor()),并与上下界比较判断数据是否均匀分布:- 均匀分布:走均匀分块优化(
splitEvenlySizedChunks),按distributionFactor * split.size动态步长切块; - 非均匀分布:当预估分片数
近似行数 / split.size超过split.sample-sharding.threshold时,按split.inverse-sampling.rate的倒数(如 1/1000)对分片键列采样,再基于采样点切分(efficientShardingThroughSampling),避免逐块查询压垮数据库; - 否则退化为逐块查询式的不均匀分块(
splitUnevenlySizedChunks),每次通过queryNextChunkMax取下一块边界,并且每 10 次查询会短暂 sleep 100ms 以保护源数据库(maySleep)。
- 均匀分布:走均匀分块优化(
无法分片时:单并发兜底
如果表既无主键/唯一索引、也未配置partition_column,则无法分片,该表会以单并发方式整体读取(对应ChunkSplitter.generateSplits()中findSplitKey返回空、创建单 Split 的分支)。
附录:常见数据源参考配置
下表汇总了 JDBC Source 常见数据源的驱动类名与连接 URL 参考值(驱动 jar 请从对应厂商的 Maven 仓库坐标获取,例如 MySQL 对应mysql:mysql-connector-java、PostgreSQL 对应org.postgresql:postgresql、Oracle 对应com.oracle.database.jdbc:ojdbc8):
| 数据源 | driver | url |
|---|---|---|
| mysql | com.mysql.cj.jdbc.Driver | jdbc:mysql://localhost:3306/test |
| postgresql | org.postgresql.Driver | jdbc:postgresql://localhost:5432/postgres |
| dm(达梦) | dm.jdbc.driver.DmDriver | jdbc:dm://localhost:5236 |
| phoenix | org.apache.phoenix.queryserver.client.Driver | jdbc:phoenix:thin:url=http://localhost:8765;serialization=PROTOBUF |
| sqlserver | com.microsoft.sqlserver.jdbc.SQLServerDriver | jdbc:sqlserver://localhost:1433 |
| oracle | oracle.jdbc.OracleDriver | jdbc:oracle:thin:@localhost:1521/xepdb1 |
| sqlite | org.sqlite.JDBC | jdbc:sqlite:test.db |
| gbase8a | com.gbase.jdbc.Driver | jdbc:gbase://e2e_gbase8aDb:5258/test |
| starrocks | com.mysql.cj.jdbc.Driver | jdbc:mysql://localhost:3306/test |
| db2 | com.ibm.db2.jcc.DB2Driver | jdbc:db2://localhost:50000/testdb |
| tablestore | com.alicloud.openservices.tablestore.jdbc.OTSDriver | jdbc:ots:http://myinstance.cn-hangzhou.ots.aliyuncs.com/myinstance |
| saphana | com.sap.db.jdbc.Driver | jdbc:sap://localhost:39015 |
| doris | com.mysql.cj.jdbc.Driver | jdbc:mysql://localhost:3306/test |
| teradata | com.teradata.jdbc.TeraDriver | jdbc:teradata://localhost/DBS_PORT=1025,DATABASE=test |
| Snowflake | net.snowflake.client.jdbc.SnowflakeDriver | jdbc |
| Redshift | com.amazon.redshift.jdbc42.Driver | jdbc:redshift://localhost:5439/testdb?defaultRowFetchSize=1000 |
| Vertica | com.vertica.jdbc.Driver | jdbc:vertica://localhost:5433 |
| Kingbase | com.kingbase8.Driver | jdbc:kingbase8://localhost:54321/db_test |
| OceanBase | com.oceanbase.jdbc.Driver | jdbc:oceanbase://localhost:2881 |
| Hive | org.apache.hive.jdbc.HiveDriver | jdbc:hive2://localhost:10000 |
| xugu(虚谷) | com.xugu.cloudjdbc.Driver | jdbc:xugu://localhost:5138 |
| InterSystems IRIS | com.intersystems.jdbc.IRISDriver | jdbc:IRIS://localhost:1972/%SYS |
注意:
compatible_mode参数仅在数据库本身支持多兼容模式时需要设置(如 OceanBase 的mysql/oracle模式)。
实战配置示例
以下示例均为完整可用的 HOCON 配置片段,可直接套用于 SeaTunnel 作业文件(完整作业还需env块与sink块,可参考 config 模板)。
示例 1:最简单的 JDBC 读取
Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" query = "select * from type_bin" }示例 2:动态分片阶段使用 select count 统计行数
适用于 Oracle 场景下 analyze 统计信息更新较慢、直接select count更快的情形:
Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" use_select_count = true query = "select * from type_bin" }示例 3:跳过动态分片阶段的表行数分析
适用于已定时执行 analyze 更新统计信息、或表数据变化不频繁的场景(当前仅 Oracle 可用):
Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" skip_analyze = true query = "select * from type_bin" }示例 4:通过 partition_column 并行读取
env { parallelism = 10 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" query = "select * from type_bin" partition_column = "id" split.size = 10000 # 读取起始边界(可选) #partition_lower_bound = ... # 读取结束边界(可选) #partition_upper_bound = ... } } sink { Console {} }示例 5:显式指定并行边界(推荐更高效)
显式给定上下界后,SeaTunnel 可以按你配置的边界直接切分,省去查询 MIN/MAX 的开销,读取更高效:
source { Jdbc { url = "jdbc:mysql://localhost:3306/test?serverTimezone=GMT%2b8&useUnicode=true&characterEncoding=UTF-8&rewriteBatchedStatements=true" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" # 按需定义查询逻辑 query = "select * from type_bin" partition_column = "id" # 读取起始边界 partition_lower_bound = 1 # 读取结束边界 partition_upper_bound = 500 partition_num = 10 properties { useSSL=false } } }示例 6:通过主键 / 唯一索引自动并行
配置table_path会自动开启动态分片(auto split),并通过split.*调整分片策略:
env { parallelism = 10 job.mode = "BATCH" } source { Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" table_path = "testdb.table1" query = "select * from testdb.table1" split.size = 10000 } } sink { Console {} }示例 7:多表读取(table_list)
配置table_list同样会自动开启动态分片,可按表粒度独立设置query(实现行/列过滤),并支持公共的where_condition:
Jdbc { url = "jdbc:mysql://localhost/test?serverTimezone=GMT%2b8" driver = "com.mysql.cj.jdbc.Driver" connection_check_timeout_sec = 100 user = "root" password = "123456" table_list = [ { # 例如 table_path = "testdb.table1"、table_path = "test_schema.table1"、table_path = "testdb.test_schema.table1" table_path = "testdb.table1" }, { table_path = "testdb.table2" # 使用 query 过滤行与列 query = "select id, name from testdb.table2 where id > 100" } ] #where_condition = "where id > 100" #split.size = 10000 #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 }配置要点速记
- 单表读取:优先用
table_path替代query;需要读取多张表时使用table_list(两种方式二选一,不可混用,源码层面已做互斥校验); partition_num不推荐使用:控制分片粒度建议直接配置split.size(默认 8096 行/分片);where_condition必须以where开头,否则作业启动即报错;- 无法分片的表(无主键/唯一索引且未设置
partition_column)会退化为单并发读取。
版本演进(Changelog)
该连接器能力持续演进,各版本变更如下:
- 2.2.0-beta(2022-09-26):新增 ClickHouse Source Connector;
- 2.3.0-beta(2022-10-20):新增 Phoenix、SQL Server、Oracle、StarRocks、GBase8a、DB2 等 JDBC Source 支持;
- 后续版本:修复 JDBC 分片 Bug;新增 Sqlite、Tablestore、Teradata、Doris、Redshift Source 支持;新增
fetch_size配置;修复连接重置问题;新增 Vertica 连接器。
小结
JDBC Source 是 SeaTunnel 接入关系型数据库最通用的入口:只需一个驱动、一段url/driver配置,即可完成从简单全表查询到多表并行读取的全部诉求。理解partition_column(固定分片)与table_path/table_list(动态分片)两条分片路径的区别,是合理配置并行度、稳定高效拉取大表数据的关键;若再配合split.*系列参数,即可应对主键分布不均、数据量极大等复杂场景。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel JDBC Source 连接器完全指南:批量读取、并行分片与多表配置实战
SeaTunnel JDBC Source 连接器完全指南:批量读取、并行分片与多表配置实战 本文是 Apache SeaTunnel 官方文档 docs/en
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Snowflake JDBC Source 连接器:从配置到并行分片读取的完整实战指南
SeaTunnel Snowflake JDBC Source 连接器:从配置到并行分片读取的完整实战指南 本文聚焦 SeaTunnel 生态中通过 JDBC
数据工程大数据批处理流处理SeaTunnel JDBC Source 完全指南:多表读取、并行快照与动态分片实战
SeaTunnel JDBC Source 完全指南:多表读取、并行快照与动态分片实战 本文是 SeaTunnel 项目中 JDBC Source 连接器的权威
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考