拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

SeaTunnel Paimon Source Connector 使用指南:从 Apache Paimon 批量读取数据的配置、查询过滤与源码实现

SeaTunnel Paimon Source Connector 使用指南:从 Apache Paimon 批量读取数据的配置、查询过滤与源码实现 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文基于 Paimon Source Connector 官方文档 展开并结合当前仓库中connector-paimon模块的源码与测试用例进行纵深解读。读完本文你将掌握在 SeaTunnel 中配置 Paimon 数据源、通过query参数完成下推式过滤与列裁剪、以及理解其底层 Catalog 加载、Split 枚举与谓词转换的完整实现链路。概述Paimon Source 能做什么SeaTunnel 的 Paimon Source Connector 用于从 Apache Paimon 表中批量读取数据Bounded 读取是数据集成任务中将 Paimon 作为读侧存储的关键组件。它在 SeaTunnel 作业中通过 HOCON 格式的source配置块声明支持基于 filesystem 或 hive 两种 Catalog 读取 Paimon 表通过query参数下发WHERE过滤条件与SELECT列投影谓词下推与列裁剪读取本地文件系统或 HDFS 上的 Paimon 仓库支持配置 Hadoop HA 相关参数。连接器的插件标识为Paimon在 plugin-mapping.properties 中映射为connector-paimon模块seatunnel.source.Paimon connector-paimon。特性支持矩阵特性支持情况batch批量读取✅ 支持stream流式读取❌ 不支持exactly-once精确一次❌ 不支持column projection列投影✅ 已实现通过query的 SELECT 列清单文档特性矩阵尚未同步更新parallelism并行度文档矩阵未勾选但源码实现了基于 Split 的枚举与多 Reader 分发机制support user-defined split❌ 不支持关于特性矩阵的说明连接器特性定义可参见 connector-v2-features.md。官方文档的特性勾选列表与 ChangelogSupport projection for Paimon Source存在时间差从 PaimonSource.java 的实现和端到端测试用例来看投影能力已经落地。工作原理从配置到数据读出的完整链路Paimon Source 的执行链路可以拆分为四个环节均能在 connector-paimon 模块中找到对应实现Catalog 加载PaimonSourceFactory根据配置创建并打开PaimonCatalog通过 PaimonCatalogLoader 调用 Apache Paimon 的CatalogFactory.createCatalog()生成底层 Catalog 实例最终定位到目标表。SQL 解析与谓词转换PaimonSource 在构造阶段读取query参数用 JSqlParser 解析为PlainSelect再由SqlToPaimonPredicateConverter将 WHERE 子句转换为 Paimon 原生Predicate、将 SELECT 列清单转换为投影下标数组。Split 枚举PaimonSourceSplitEnumerator 基于 Paimon 表的TableScan能力把待读数据划分为多个 Split 并分配给各个 Reader同时实现了基于PaimonSourceState的 Checkpoint 恢复逻辑。数据读取PaimonSourceReader 通过table.newReadBuilder().withProjection(projection).withFilter(predicate)构建读取器将 Paimon 的InternalRow逐行转换为 SeaTunnel 的SeaTunnelRow并向下游输出。由于PaimonSource.getBoundedness()返回BOUNDED连接器天然是批处理模式读完全部 Split 后Reader 会调用context.signalNoMoreElement()结束任务见 PaimonSourceReader.java。参数详解以下是官方文档给出的完整参数表结合 PaimonConfig.java 与 PaimonSourceConfig.java 中的源码定义逐项说明参数名类型必填默认值说明warehouseString是-Paimon 仓库路径如/tmp/paimon或hdfs:///tmp/paimoncatalog_typeString否filesystemPaimon Catalog 类型支持filesystem和hivecatalog_uriString否-Paimon Catalog 的 URI仅catalog_typehive时需要如thrift://hadoop04:9083databaseString是-要访问的 Paimon 数据库名tableString是-要访问的 Paimon 表名hdfs_site_pathString否-hdfs-site.xml文件路径源码中已标记Deprecated推荐改用paimon.hadoop.conf-pathqueryString否-表读取的过滤条件例如select * from st_test where id 100不指定则读取全表paimon.hadoop.confMap否-Hadoop 配置属性以 Key-Value 形式内联声明paimon.hadoop.conf-pathString否-指定加载core-site.xml、hdfs-site.xml、hive-site.xml的路径必填参数的校验逻辑在 PaimonConfig.java 中warehouse、database、table三项在构造时经过checkArgumentNotBlank强校验任一为空都会抛出SeaTunnelException而当catalog_type为hive时catalog_uri同样被强校验为必填。这一点与 PaimonSourceFactory 中OptionRule的conditional(CATALOG_TYPE, HIVE, CATALOG_URI)声明一致只有选择了 hive Catalog 才要求提供 URI。catalog_typefilesystem 与 hivecatalog_type的合法取值定义在 PaimonCatalogEnum.java 中即filesystem与hive。二者在 PaimonCatalogLoader.java 中的加载差异如下filesystem仅把warehouse与metastorefilesystem写入 PaimonCatalogOptions直接从文件系统解析表元数据hive额外写入CatalogOptions.URI即catalog_uri并且会把paimon.hadoop.conf中配置的所有属性合并进 Options最终通过 Paimon 的CatalogContext.create(options, hadoopConfiguration)完成 Catalog 构建。另外当warehouse以hdfs://开头时加载器会强校验 Hadoop 配置中必须存在fs.defaultFS并自动补充fs.hdfs.impl org.apache.hadoop.hdfs.DistributedFileSystem这是连接 HDFS 仓库的关键前置条件。query 参数SQL 过滤与列投影的实战规则query是 Paimon Source 最核心的进阶参数语法与限制如下支持的语法select * from st_test where id 100仅支持简单 SELECT 语句必须包含 WHERE 子句或省略 WHERE省略则读取全表过滤条件支持、、、、、!、or、and、is null、is not null暂不支持Having、Group By、Order By这些子句 Paimon 本身不支持因此无法转换投影与 limit投影已支持limit 计划在未来支持。字符串与布尔值的引号规则当 WHERE 条件中的字段是字符串或布尔类型时其值必须用单引号包裹否则会报错。例如nameabc or tagtrue支持的数据类型WHERE 条件目前支持以下字段数据类型stringbooleantinyintsmallintintbigintfloatdoubledatetimestamp源码层面的实现细节这些规则并非空谈而是由 SqlToPaimonPredicateConverter.java 逐条保证语句校验convertToPlainSelect()只接受Select类型的语句且SelectBody必须是PlainSelect一旦检测到 Having / Group By / Order By / Limit 任一子句立即抛出IllegalArgumentException(Only SELECT statements with WHERE clause are supported...)运算符映射parseExpressionToPredicate()把 JSqlParser 的EqualsTo、GreaterThan、GreaterThanEquals、MinorThan、MinorThanEquals、NotEqualsTo、IsNullExpression、AndExpression、OrExpression、Parenthesis分别映射为 PaimonPredicateBuilder的equal、greaterThan、greaterOrEqual、lessThan、lessOrEqual、notEqual、isNull、isNotNull、and、or等调用遇到不支持的类型直接抛异常值类型转换convertValueByPaimonDataType()根据 Paimon 字段的DataType.getTypeRoot()将 SQL 字面量转换为 Paimon 内部类型——字符串/布尔/Decimal含精度与小数位/tinyint/smallint/int/bigint/float/double/date/timestamp 均有对应分支其中 date 通过DateTimeUtils.toInternal转为内部天数、timestamp 通过Timestamp.fromLocalDateTime转换列名校验WHERE 或 SELECT 中引用了表中不存在的列时会抛出 The column named [xxx] is not exists 异常避免静默错误投影下推convertSqlSelectToPaimonProjectionIndex()把 SELECT 的列清单映射为表字段下标数组遇到SELECT *返回null表示全列随后在PaimonSource中通过getTableWithProjection构建投影后的CatalogTable真正实现列裁剪。单元测试 SqlToPaimonConverterTest.java 对上述 13 种字段类型char、varchar、boolean、binary、decimal、tinyint、smallint、int、bigint、float、double、date、timestamp的等值谓词转换、IS NULL、IS NOT NULL、AND/OR组合以及投影下标计算SELECT decimal_col, int_col, char_col, timestamp_col, boolean_col→{4, 7, 0, 12, 2}做了完整断言验证可作为编写 query 时对照的权威参考。配置示例1. 简单示例读取本地文件系统仓库source { Paimon { warehouse /tmp/paimon database default table st_test } }2. 过滤示例多种类型混合的下推条件source { Paimon { warehouse /tmp/paimon database full_type table st_test query select c_boolean, c_tinyint from st_test where c_boolean true and c_tinyint 116 and c_smallint 15987 or c_decimal2924137191386439303744.39292213 } }注意其中布尔值true和 decimal 值均使用了单引号该语句同时演示了投影只读c_boolean、c_tinyint两列与谓词下推and/or 混合。3. Hadoop 配置示例HDFS HA 集群当仓库位于 HDFS 且集群启用了 NameNode HA 时必须通过paimon.hadoop.conf内联提供 HA 相关属性source { Paimon { catalog_nameseatunnel_test warehousehdfs:///tmp/paimon databaseseatunnel_namespace1 tablest_test query select * from st_test where pk_id is not null and pk_id 3 paimon.hadoop.conf { fs.defaultFS hdfs://nameservice1 dfs.nameservices nameservice1 dfs.ha.namenodes.nameservice1 nn1,nn2 dfs.namenode.rpc-address.nameservice1.nn1 hadoop03:8020 dfs.namenode.rpc-address.nameservice1.nn2 hadoop04:8020 dfs.client.failover.proxy.provider.nameservice1 org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider dfs.client.use.datanode.hostname true } } }这里catalog_name是可选项默认值为paimon见 PaimonConfig.java示例同时演示了is not null与的组合过滤。4. Hive Catalog 示例当 Paimon 表元数据托管在 Hive Metastore 时使用catalog_typehive并配套catalog_urisource { Paimon { catalog_nameseatunnel_test catalog_typehive catalog_urithrift://hadoop04:9083 warehousehdfs:///tmp/seatunnel databaseseatunnel_test tablest_test3 paimon.hadoop.conf { fs.defaultFS hdfs://nameservice1 dfs.nameservices nameservice1 dfs.ha.namenodes.nameservice1 nn1,nn2 dfs.namenode.rpc-address.nameservice1.nn1 hadoop03:8020 dfs.namenode.rpc-address.nameservice1.nn2 hadoop04:8020 dfs.client.failover.proxy.provider.nameservice1 org.apache.hadoop.hdfs.server.namenode.ha.ConfiguredFailoverProxyProvider dfs.client.use.datanode.hostname true } } }端到端验证仓库内的真实测试资产如果你希望在实际环境验证连接器行为仓库提供了完整的集成测试资产connector-paimon-e2ePaimonIT.java 的流程是先执行fake_to_paimon.conf用 Fake 源向 Paimon 写入 100000 行数据再执行paimon_to_assert.conf全量读取、paimon_projection_to_assert.conf投影读取最后用 Assert Sink 校验行数与字段非空约束paimon_to_assert.conf 展示了最朴素的完整作业骨架env中job.mode BATCHsource 块只声明warehouse/database/table三个必填参数sink 块对接 Assertpaimon_projection_to_assert.conf 演示了query select c_string, c_boolean from st_test where c_string is not null这种投影 过滤组合的实测配置另有paimon_to_assert_with_filter1.confselect * from st_test where c_string is not null、Hive Catalog 版本paimon_to_assert_with_hivecatalog.conf以及 HDFS HA 版本read_from_paimon_with_hdfs_ha_to_assert.conf可供参考。这些配置与官方文档中的示例相互印证是排查为什么读不到数据 / 过滤没生效 / 列对不上时最直接的对照样本。依赖与版本说明从 connector-paimon/pom.xml 可以看到连接器的关键依赖paimon-bundle版本为0.7.0-incubating这是底层 Paimon 能力的来源Split 扫描、谓词执行、行转换均由该库提供jsqlparser负责 SQL 文本解析支撑query参数的实现hive-exec2.3.9provided 作用域支撑 hive Catalog 模式下的元数据交互seatunnel-hadoop3-3.1.4-uberprovided提供 HDFS 客户端能力。因此使用 Hive Catalog 或 HDFS 仓库时请确认运行环境SeaTunnel 发行包的lib/或 connectors 目录已包含对应依赖使用本地 filesystem 仓库则无需额外部署。总结Paimon Source Connector 是 SeaTunnel 打通Paimon 存储 ↔ 下游数据系统批式读取的标准化入口。它把定位表Catalog 加载→ 拆数据Split 枚举→ 下推条件SQL 谓词转换→ 行转换InternalRow → SeaTunnelRow整条链路封装在connector-paimon模块中用户只需在source块中声明warehouse/database/table三个必填参数即可按需通过query实现谓词下推与列投影。理解 PaimonConfig 中参数语义与 SqlToPaimonPredicateConverter 的运算符/类型支持边界是避免踩坑如引号规则、HA 配置缺失、不支持的子句的关键。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel Cassandra Source Connector 使用指南从 Apache Cassandra 批量读取数据SeaTunnel Cassandra Source Connector 使用指南从 Apache Cassandra 批量读取数据 本指南面向 SeaTun数据工程大数据批处理流处理Checkmate 监控部署指南从克隆到告警跑通的完整流程Checkmate 监控部署指南从克隆到告警跑通的完整流程 服务器挂了第一个知道的却是报障的客户——没有专人盯基础设施时你需要一个自己就能跑起来的哨兵。C数据工程大数据批处理流处理SeaTunnel Notion Source Connector 实战指南从 Notion API 批量读取数据SeaTunnel Notion Source Connector 实战指南从 Notion API 批量读取数据 Notion Source 是 SeaTu数据工程大数据批处理流处理上一篇探索SerialPlot高效串口数据可视化的实战指南下一篇Wedecode深度技术解析揭秘微信小程序逆向分析核心架构与实战应用创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表