尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel Doris 源连接器完全指南:配置参数、多表读取与源码级原理剖析

SeaTunnel Doris 源连接器完全指南:配置参数、多表读取与源码级原理剖析 数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载本指南以 Apache SeaTunnel 官方文档 docs/zh/connectors/source/Doris.md 为骨架结合 connector-doris 模块 的源码实现与 connector-doris-e2e 集成测试用例进行深度扩充讲解如何在 SeaTunnel 中配置 Doris 源连接器完成批式数据读取并解释其底层分片、并行与类型映射机制。概述为什么需要 Doris 源连接器Apache Doris 是一个基于 MPP 架构的高性能实时分析型数据库常用于数仓与报表场景。当需要把 Doris 中的数据同步到其他存储如 Kafka、Hive、其他数据库或参与离线 ETL 时就可以使用 SeaTunnel 的Doris 源连接器插件名Doris见 DorisSource.java。该连接器通过 Doris FE 的_query_planHTTP 接口获取查询执行计划将表数据按 tablet 划分为多个分片split再并发地从 BE 节点拉取数据Arrow 列式格式从而天然支持高并行度的批式读取并支持单表与多表两种使用形态。阅读完本文你将掌握Doris 源连接器的全部配置项、默认值与适用场景单表读取、列投影、下推过滤、多表读取的完整配置写法连接器底层如何发现分片、如何按 tablet 切分、如何并发拉取 Arrow 数据Doris 与 SeaTunnel 之间的数据类型映射规则。支持的引擎与功能矩阵Doris 源连接器支持以下执行引擎与 Doris.md 保持一致SparkFlinkSeaTunnel Zeta其功能支持情况如下各特性说明可参考 connector-v2-features.md特性支持情况批处理Batch✅ 支持流处理Streaming❌ 不支持精确一次Exactly Once❌ 不支持列投影Column Projection✅ 支持并行度Parallelism✅ 支持用户自定义分片✅ 支持多表读Multi-Table Read✅ 支持从源码可以印证批处理这一约束DorisSource.java 中getBoundedness()固定返回Boundedness.BOUNDED即 Doris 源连接器是一个有界一次性读完的数据源任务应配置为job.mode BATCH。环境准备与依赖安装在编写配置文件之前需要根据运行引擎安装 JDBC 驱动。Doris 兼容 MySQL 协议因此使用的是 MySQL 驱动包Spark / Flink 引擎下载mysql-connector-javaMySQL Connector/J驱动 jar放入${SEATUNNEL_HOME}/plugins/目录SeaTunnel Zeta 引擎下载同样的驱动 jar放入${SEATUNNEL_HOME}/lib/目录。驱动可从 Maven 中央仓库mysql组下的mysql-connector-java构件获取注意选择与运行环境JDK 版本匹配的版本。支持的数据源版本数据源支持版本驱动UrlMavenDoris仅支持 Doris 2.0 及以上版本---版本约束来自 Doris.md 的官方声明。此外从 DorisTypeConverterV2.java 的类注释可以确认连接器针对 Doris 2.x 提供了专门的类型转换器负责解析 2.x 引入的复杂类型ARRAY、MAP、DECIMALV3 等。数据类型映射Doris 源连接器内置类型转换器将 Doris 数据类型映射为 SeaTunnel 数据类型完整映射关系如下Doris 数据类型SeaTunnel 数据类型INTINTTINYINTTINYINTSMALLINTSMALLINTBIGINTBIGINTLARGEINTSTRINGBOOLEANBOOLEANDECIMALDECIMAL精度 列的指定长度 1小数位 小数点右侧位数FLOATFLOATDOUBLEDOUBLECHAR / VARCHAR / STRING / TEXTSTRINGJSONSTRINGVARIANTSTRINGDATEDATEDATETIME / DATETIME(p)TIMESTAMPARRAYARRAY映射规则的源码印证DATETIME → TIMESTAMP在 DorisTypeConverterV2.java 中DATETIME被转换为LocalTimeType.LOCAL_DATE_TIME_TYPE并继承列的 scale小数秒精度缺省为 0DECIMAL → DECIMALDorisTypeConverterV2.java 读取 Doris 列的 precision 与 scale 构造DecimalTypeARRAY 复杂类型转换器会递归解析ARRAY...内部的元素类型如tinyint(1)→ BOOLEAN 数组、DECIMAL(p,s)→ DECIMAL 数组、date→ DATE 数组等见 DorisTypeConverterV2.javaMAP 复杂类型Doris 2.x 的MAPK,V会被解析为 SeaTunnel 的MapType解析逻辑见 DorisTypeConverterV2.java。反向转换SeaTunnel → Doris用于写入或建表同样由DorisTypeConverterV2承担例如TIMESTAMP会按 scale 还原为DATETIME(p)JSON/VARIANT源类型会被保留。源选项详解基础配置配置项类型是否必须默认值说明fenodesstring是-Doris FE 地址格式fe_host:fe_http_port多个 FE 用英文逗号分隔usernamestring是-访问 Doris 的用户名passwordstring是-访问 Doris 的密码doris.request.retriesint否3请求 Doris FE 的重试次数doris.request.read.timeout.msint否30000请求 Doris BE 的 socket 读取超时时间毫秒doris.request.connect.timeout.msint否30000请求 Doris FE 或 BE 的连接超时时间毫秒query-portint否9030Doris 查询端口FE 的 MySQL 协议端口doris.request.query.timeout.sint否3600Doris 扫描数据的超时时间单位秒doris.request.tablet.sizeint否Integer.MAX_VALUE每个 SeaTunnel split 包含的 Doris tablet 数量最小值为1doris.deserialize.arrow.asyncboolean否false是否异步反序列化 Arrow 数据doris.request.retriesdoris.deserialize.queue.sizeint否64异步反序列化 Arrow 数据时使用的队列大小table_listArray否-要读取的 Doris 表清单多表读取时使用重要提示doris.request.retriesdoris.deserialize.queue.size这个配置名看起来像笔误但它正是当前运行时实际使用的配置名。调整异步 Arrow 反序列化队列大小时必须按这个完整名称配置。这一事实可以从 DorisSourceOptions.java 中的Options.key(doris.request.retriesdoris.deserialize.queue.size)得到源码级确认。表清单配置table_list 内部字段配置项类型是否必须默认值说明databasestring是-数据库名tablestring是-表名doris.read.fieldstring否-选择要读取的 Doris 表字段逗号分隔的列名列表doris.filter.querystring否-数据过滤条件格式字段 值例如doris.filter.query F_ID 2会原样下推给 Doris 执行doris.request.tablet.sizeint否Integer.MAX_VALUE当前表每个 SeaTunnel split 包含的 Doris tablet 数量最小值为1表级覆盖doris.batch.sizeint否1024每次从 BE 读取的最大行数doris.exec.mem.limitlong否2147483648单个 BE 扫描请求可使用的最大内存默认 2G2147483648 字节展平规则当只读单张表时table_list中的配置项如doris.read.field、doris.filter.query、doris.request.tablet.size可以展平到 source 外层直接配置此时必须在外层配置database和table。若使用table_list则每张表在内部配置database、table及其表级参数。配置解析的源码行为在 DorisTableConfig.java 中可以看到解析期的实际校验逻辑未配置table_list时会用外层database/table构造单表配置并读取外层doris.read.field、doris.filter.query、doris.batch.size、doris.request.tablet.size、doris.exec.mem.limit作为该表的配置多表配置时database与table均不允许为空缺失会抛出IllegalArgumentException且不允许出现重复或缺失的database.table组合doris.batch.size、doris.exec.mem.limit、doris.request.tablet.size若被配置为 ≤ 0会自动回退为默认值1024 / 2147483648 / Integer.MAX_VALUE无需担心非法输入导致任务失败。高级参数提示不建议随意修改doris.request.retries、doris.request.read.timeout.ms、doris.request.connect.timeout.ms、doris.request.query.timeout.s、doris.deserialize.arrow.async、doris.request.retriesdoris.deserialize.queue.size等高级参数仅在出现明显超时、重试或内存问题时结合集群实际情况做针对性调整。源码级原理连接器如何读取 Doris 数据理解底层执行路径有助于合理设置并行度与分片参数。Doris 源连接器的读取链路可以概括为FE 取计划 → 按 tablet 分片 → BE 并行拉取 Arrow 数据1. 分片发现Split Discovery在 DorisSourceSplitEnumerator.java 中枚举器启动后会对每张表调用RestService.findPartitions(...)发现分片。核心逻辑位于 RestService.java根据读取字段构造 SQL未配置列投影时默认select *配置了doris.read.field则拼接对应列名doris.filter.query会以where ...形式拼入 SQL通过 HTTP POST 调用 FE 的http://fe_host:fe_http_port/api/{database}/{table}/_query_plan接口携带{sql: ...}JSON 体使用 Basic Auth 认证解析返回的查询计划将 tablet 按所在 BE 节点分组最终把(be地址, tabletId 列表, 查询计划, 数据库, 表名)封装为PartitionDefinition即一个读取分片。值得注意的是多个 FE 地址会先被Collections.shuffle打乱顺序后逐个尝试避免请求始终压在某一个 FE 上RestService.java。2. Split 分配与并行度发现的所有 split 会先按 splitId 排序再以轮询round-robin方式平均分配给各个 reader即并行子任务实现负载均衡见 DorisSourceSplitEnumerator.java。同时枚举器支持snapshotState可将待分配分片写入状态供失败恢复时重新调度同文件 DorisSourceSplitEnumerator.java。实践含义doris.request.tablet.size控制每个 split 打包多少个 tablet。表的分片总数约等于表内 tablet 数 ÷ 每 split tablet 数分片越多可被并行度消化的程度越高env.parallelism决定同时运行多少个 reader。两者配合即可放大读取吞吐。3. 数据读取与 Arrow 反序列化每个 reader 针对分配到的分片创建 DorisValueReader通过 Thrift 协议在 BE 上打开 scanneropenScanner传入 database、table、tabletId 列表、batchSize、查询超时、内存上限等参数见 DorisValueReader.java之后循环调用getNext拉取TScanBatchResult将 Arrow 列式数据反序列化为SeaTunnelRow当doris.deserialize.arrow.async true时会启动一个后台线程提前拉取并反序列化数据放入有界阻塞队列容量由doris.request.retriesdoris.deserialize.queue.size决定实现拉取与消费的流水线并行见 DorisValueReader.java。默认 false 表示同步反序列化。实战示例示例一单表读取从 Doris 读取单表数据并输出到控制台env { parallelism 2 job.mode BATCH } source { Doris { fenodes doris_e2e:8030 username root password database e2e_source table doris_e2e_table } } transform { # 如需了解更多 transform 插件配置请参考 SQL 转换插件文档 } sink { Console {} }要点fenodes中的端口是 FE 的HTTP 端口示例为 8030按实际部署调整query-port默认 9030是 FE 的查询端口两者不同。示例二列投影doris.read.field使用doris.read.field只读取需要的字段减少网络与内存开销env { parallelism 2 job.mode BATCH } source { Doris { fenodes doris_e2e:8030 username root password database e2e_source table doris_e2e_table doris.read.field F_ID,F_INT,F_BIGINT,F_TINYINT,F_SMALLINT } } transform { } sink { Console {} }示例三下推过滤doris.filter.querydoris.filter.query的值会作为过滤条件原样传递给 Doris由 Doris 在源端完成数据裁剪而不是把全表拉到 SeaTunnel 再过滤env { parallelism 2 job.mode BATCH } source { Doris { fenodes doris_e2e:8030 username root password database e2e_source table doris_e2e_table doris.filter.query F_ID 2 } } transform { } sink { Console {} }示例四多表读取table_listtable_list支持一次读取多张表每张表可独立配置字段、过滤条件、分片与内存参数非常适合需要一次任务同步多张表的场景。下游 sink 可使用${table_name}动态获取当前表名env { parallelism 1 job.mode BATCH } source { Doris { fenodes xxxx:8030 username root password table_list [ { database st_source_0 table doris_table_0 doris.read.field F_ID,F_INT,F_BIGINT,F_TINYINT doris.filter.query F_ID 50 doris.request.tablet.size 1 doris.exec.mem.limit 2147483648 }, { database st_source_1 table doris_table_1 } ] } } transform {} sink { Doris { fenodes xxxx:8030 schema_save_mode RECREATE_SCHEMA username root password database st_sink table ${table_name} sink.enable-2pc true sink.label-prefix test_json doris.config { format json read_json_by_line true } } }多表读取的集成测试参考仓库的 e2e 测试真实演示了多表读取与断言校验的组合用法可作为配置范本。见 doris_multi_source_to_assert.confsource { Doris { fenodes doris_e2e:8030 username root password table_list [ { database e2e_source_0 table doris_e2e_unique_table_0 doris.read.field F_ID,F_INT,F_BIGINT,F_TINYINT,F_SMALLINT, ... doris.filter.query F_ID 50 }, { database e2e_source_1 table doris_e2e_unique_table_1 doris.read.field F_ID,F_INT,F_BIGINT,F_TINYINT,F_SMALLINT, ... doris.filter.query F_ID 40 } ] } }该测试同时验证了复杂字段类型含 MAP 系列字段在多表读取场景下的正确性对应测试类为 DorisMultiReadIT.java。更多读写组合配置可查阅 connector-doris-e2e 测试资源目录 下的doris_source_*.conf系列文件。常见问题与调优建议任务报连接/超时错误优先检查fenodes的 FE HTTP 端口是否正确、query-port是否为 FE 查询端口大表读取可将doris.request.read.timeout.ms与doris.request.query.timeout.s适当调大并行度上不去确认doris.request.tablet.size未设置过大保证分片数量不小于env.parallelism否则部分并行子任务会无分片可读源码中 split 按轮询分配分片数少于 reader 数时部分 reader 拿不到数据吞吐优化在内存充足时可开启doris.deserialize.arrow.async true并用doris.request.retriesdoris.deserialize.queue.size调节异步队列深度同时可适当调大doris.batch.size内存控制单 BE 扫描内存由doris.exec.mem.limit默认 2G控制多表并行读取时要为每个任务预留足够的内存预算。总结SeaTunnel 的 Doris 源连接器通过FE 查询计划 tablet 分片 BE Arrow 并行拉取的架构为 Doris 2.0 提供了高并行、支持列投影与过滤下推、支持多表读取的批式数据接入能力。配置层面只需掌握fenodes/username/password三个必填项与table_list或databasetable两种表指定方式即可快速完成单表与多表同步如需深入调优可结合本文梳理的源码路径DorisSourceOptions.java、DorisSourceSplitEnumerator.java、RestService.java、DorisValueReader.java进一步定位参数生效位置与影响范围。变更日志连接器的版本更新记录见仓库文档 connector-doris 变更日志其中记录了每次发版对 Doris 源/宿连接器的功能新增与问题修复。赞分享数据集成ETL大数据批处理流处理变更数据捕获【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/GitHub_Trending/se/seatunnel点击查看免费下载相关推荐SeaTunnel RocketMQ 源连接器完全指南参数配置、多表读取与精确一次消费原理SeaTunnel RocketMQ 源连接器完全指南参数配置、多表读取与精确一次消费原理 本文以 docs/zh/connectors/source/Roc数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel JDBC MySQL 源连接器完全指南参数配置、并行读取拆分与多表同步实战SeaTunnel JDBC MySQL 源连接器完全指南参数配置、并行读取拆分与多表同步实战 导读 本文以 SeaTunnel 官方文档中 MySQL 源连数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel MaxCompute 源连接器实战配置详解、多表读取与并行读取实现原理SeaTunnel MaxCompute 源连接器实战配置详解、多表读取与并行读取实现原理 本文基于 SeaTunnel 仓库中的 MaxCompute 源连数据集成ETL大数据批处理流处理变更数据捕获创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表