ARTICLE DETAIL

资讯详情

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

SeaTunnel JDBC Source 连接器完全指南:配置详解、并行分片原理与多表读取实战

SeaTunnel JDBC Source 连接器完全指南:配置详解、并行分片原理与多表读取实战 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文以 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为必填其余均为可选参数类型必填默认值说明urlString是-JDBC 连接地址例如jdbc:postgresql://localhost/testdriverString是-连接远程数据源的 JDBC 驱动类名如 MySQL 为com.mysql.cj.jdbc.DriveruserString否-用户名passwordString否-密码queryString否-查询语句compatible_modeString否-数据库兼容模式当数据库支持多种兼容模式时必须设置如 OceanBase 需设置为mysql或oracleconnection_check_timeout_secInt否30用于校验连接的数据操作最长等待秒数partition_columnString否-用于数据分片的列名partition_upper_boundLong否-扫描的partition_column最大值不设置时 SeaTunnel 会查询数据库获取最大值partition_lower_boundLong否-扫描的partition_column最小值不设置时 SeaTunnel 会查询数据库获取最小值partition_numInt否作业并行度分片数量不推荐使用正确做法是通过split.size控制分片数仅支持正整数use_select_countBoolean否false动态分片阶段是否用select count统计行数当前仅 Oracle 可用当用 analyze 语句更新统计信息较慢时可直接使用 select countskip_analyzeBoolean否false动态分片阶段跳过表行数分析当前仅 Oracle 可用适用于已定时执行 analyze 或表数据变化不频繁的场景fetch_sizeInt否0查询返回大量对象时配置每次抓取的行数以提升性能0 表示使用 JDBC 默认值propertiesMap否-附加连接参数当 properties 与 url 中参数同名时优先级由驱动实现决定如 MySQL 中 properties 优先于 URLtable_pathString否-表的完整路径可替代query。示例MySQLtestdb.table1、Oracletest_schema.table1、SQL Servertestdb.test_schema.table1、PostgreSQLtestdb.test_schema.table1、InterSystems IRIStest_schema.table1table_listArray否-要读取的表列表可替代table_pathwhere_conditionString否-作用于所有表/查询的公共行过滤条件必须以where开头例如where id 100split.sizeInt否8096每个分片包含的行数读取表时按该值切分split.even-distribution.factor.lower-boundDouble否0.05不推荐使用分片键分布因子下界用于判断表数据是否均匀分布。分布因子 (MAX(id) − MIN(id) 1) / 行数若大于等于该下界则按均匀分布优化分块否则视为非均匀分布在预估分片数超过sample-sharding.threshold时采用采样分片策略split.even-distribution.factor.upper-boundDouble否100不推荐使用分片键分布因子上界分布因子小于等于该上界时按均匀分布优化分块否则视为非均匀分布并可能触发采样分片split.sample-sharding.thresholdInt否1000触发采样分片策略的预估分片数阈值。当分布因子落在上下界之外且预估分片数近似行数 / chunk 大小超过该阈值时启用采样分片策略以更高效地处理大数据集split.inverse-sampling.rateInt否1000采样分片策略中采样率的倒数。例如设为 1000 表示按 1/1000 采样用于控制采样粒度、影响最终分片数非常适合超大数据集common-options-否-Source 插件通用参数详见 Source Common Options参数背后的源码实现默认值一致性上述connection_check_timeout_sec 30、fetch_size 0定义于 JdbcOptionssplit.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会被解析为MapString, String见 JdbcConnectionConfig在连接建立时合并进 JDBC URLcompatible_mode用于 OceanBase 这类多兼容模式数据库帮助JdbcDialectLoader选择正确的方言实现。并行读取与分片Split机制JDBC Source 支持从表中并行读取数据。SeaTunnel 会依据一定规则将表数据切分成多个 Split交给多个 Reader 并行消费Reader 的数量由作业parallelism决定。分片键Split Key的选择规则如下若配置了partition_column则直接用该列参与分片计算该列必须属于支持的分片数据类型若未配置partition_columnSeaTunnel 会读取表结构Schema依次查找主键Primary Key和唯一索引Unique Index取其中第一个属于支持的分片数据类型的列作为分片键。例如某表主键为(nn guid, name varchar)因guid不在支持类型内会退而选择name列参与分片。支持的分片数据类型String字符串Numberint、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其余情况即未同时提供querypartition_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数据源driverurlmysqlcom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testpostgresqlorg.postgresql.Driverjdbc:postgresql://localhost:5432/postgresdm达梦dm.jdbc.driver.DmDriverjdbc:dm://localhost:5236phoenixorg.apache.phoenix.queryserver.client.Driverjdbc:phoenix:thin:urlhttp://localhost:8765;serializationPROTOBUFsqlservercom.microsoft.sqlserver.jdbc.SQLServerDriverjdbc:sqlserver://localhost:1433oracleoracle.jdbc.OracleDriverjdbc:oracle:thin:localhost:1521/xepdb1sqliteorg.sqlite.JDBCjdbc:sqlite:test.dbgbase8acom.gbase.jdbc.Driverjdbc:gbase://e2e_gbase8aDb:5258/teststarrockscom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testdb2com.ibm.db2.jcc.DB2Driverjdbc:db2://localhost:50000/testdbtablestorecom.alicloud.openservices.tablestore.jdbc.OTSDriverjdbc:ots:http://myinstance.cn-hangzhou.ots.aliyuncs.com/myinstancesaphanacom.sap.db.jdbc.Driverjdbc:sap://localhost:39015doriscom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testteradatacom.teradata.jdbc.TeraDriverjdbc:teradata://localhost/DBS_PORT1025,DATABASEtestSnowflakenet.snowflake.client.jdbc.SnowflakeDriverjdbc//account_name.snowflakecomputing.comRedshiftcom.amazon.redshift.jdbc42.Driverjdbc:redshift://localhost:5439/testdb?defaultRowFetchSize1000Verticacom.vertica.jdbc.Driverjdbc:vertica://localhost:5433Kingbasecom.kingbase8.Driverjdbc:kingbase8://localhost:54321/db_testOceanBasecom.oceanbase.jdbc.Driverjdbc:oceanbase://localhost:2881Hiveorg.apache.hive.jdbc.HiveDriverjdbc:hive2://localhost:10000xugu虚谷com.xugu.cloudjdbc.Driverjdbc:xugu://localhost:5138InterSystems IRIScom.intersystems.jdbc.IRISDriverjdbc:IRIS://localhost:1972/%SYS注意compatible_mode参数仅在数据库本身支持多兼容模式时需要设置如 OceanBase 的mysql/oracle模式。实战配置示例以下示例均为完整可用的 HOCON 配置片段可直接套用于 SeaTunnel 作业文件完整作业还需env块与sink块可参考 config 模板。示例 1最简单的 JDBC 读取Jdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%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?serverTimezoneGMT%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?serverTimezoneGMT%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?serverTimezoneGMT%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?serverTimezoneGMT%2b8useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue 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 { useSSLfalse } } }示例 6通过主键 / 唯一索引自动并行配置table_path会自动开启动态分片auto split并通过split.*调整分片策略env { parallelism 10 job.mode BATCH } source { Jdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%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_conditionJdbc { url jdbc:mysql://localhost/test?serverTimezoneGMT%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-beta2022-09-26新增 ClickHouse Source Connector2.3.0-beta2022-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.*系列参数即可应对主键分布不均、数据量极大等复杂场景。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐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大数据批处理流处理变更数据捕获上一篇Godot Pck Tool5个核心技巧掌握Godot游戏资源管理终极方案下一篇复古游戏机风格重现Pyxelate内置调色板PICO-8/Apple II使用指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表