多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Beam SQL CLI 添加新数据源实战:基于 TableProvider 扩展自定义 SQL 表

Apache Beam SQL CLI 添加新数据源实战:基于 TableProvider 扩展自定义 SQL 表 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载Apache Beam 的 SQL 扩展允许直接用 SQL 编写数据处理管道无论是 Batch 还是 Streaming 场景。本文以 Beam SQL CLIBeamSQL交互式 shell为背景围绕「如何给 Beam SQL 添加一个全新的数据源」这一核心主题展开你会完整掌握TableProvider与BeamSqlTable两个关键抽象的职责分工亲手实现一个基于GenerateSequence的「连续整数流」数据表并学会在 SQL CLI 中通过CREATE EXTERNAL TABLE注册它、用SELECT与窗口聚合查询它。文中的实现代码、命令与配置均取自当前仓库可直接对照源码阅读与验证。背景Beam SQL 与CREATE EXTERNAL TABLEBeam 的 SQL 能力通过 Java SDK 中的SqlTransform源码位于 SqlTransform.java接入管道。而 SQL CLI 则是 Beam 提供的交互式 shell它以sqlline为交互前端把用户输入的 SQL 语句翻译成 Beam 管道默认使用DirectRunner在本地执行并返回结果表格。SQL 管道中的数据来源并不是「持久化」的物理表而是通过CREATE EXTERNAL TABLE语句声明的虚拟表。在 Beam SQL 中所有表都是 external 的——该语句只负责把某个外部数据源文件、消息队列、数据库等注册成一张可供查询的表。该语句的完整语法见仓库文档 create-external-table.md核心结构如下CREATE EXTERNAL TABLE [ IF NOT EXISTS ] tableName (tableElement [, tableElement ]*) TYPE type [LOCATION location] [TBLPROPERTIES tblProperties]其中tableElement声明列形如columnName fieldType [ NOT NULL ]fieldType支持TINYINT、SMALLINT、INTEGER、BIGINT、FLOAT、DOUBLE、DECIMAL、BOOLEAN、DATE、TIME、TIMESTAMP、CHAR、VARCHAR等简单类型以及MAP、ARRAY、ROW等复合类型TYPE标识「由哪个表提供者TableProvider来支撑这张表」例如bigquery、bigtable、pubsub、kafka、text、mongodb、parquet、avro、pubsublite等LOCATION是 I/O 特定的地址描述如文件路径、topic 名TBLPROPERTIES是以字符串字面量给出的 JSON 对象用于传递该 I/O 特有的额外配置。新数据源 新的TYPE。内置支持之外的数据源就是通过新增一个TableProvider来接入的——这正是本文要解决的问题。核心抽象TableProvider与BeamSqlTableSqlTransform在遇到CREATE EXTERNAL TABLE语句时会依赖一组TableProvider来完成建表元数据操作与实际数据读写。TableProvider的接口定义位于 TableProvider.java它负责方法职责getTableType()返回该 provider 处理的表类型标识对应 SQL 中TYPE的值createTable(Table)/dropTable(String)表元数据的创建与删除CRUDgetTables()/getTable(name)列出 / 获取该 provider 管理的所有表buildBeamSqlTable(Table)依据元数据构建真正的读写实现返回BeamSqlTable而从接口注释可以确认所有标注了AutoService(TableProvider.class)的实现都会在使用默认连接参数启动JdbcDriverSQL CLI 正是如此时被自动加载。BeamSqlTable接口见 BeamSqlTable.java则定义了表的具体读写行为buildIOReader(PBegin)从数据源构建一个PCollectionRow读路径buildIOWriter(PCollectionRow)把PCollectionRow写入目标写路径isBounded()返回IsBounded.BOUNDED或IsBounded.UNBOUNDED决定 SQL 优化器按批还是按流处理getSchema()返回表的SchemagetTableStatistics(PipelineOptions)提供行数有界表或速率无界表的估算。Table抽象类见 Table.java承载了解析CREATE EXTERNAL TABLE后的全部元数据getType()、getName()、getSchema()、getLocation()以及用于读取TBLPROPERTIES的getProperties()。因此给 Beam SQL CLI 添加一个数据源的完整路径就是实现一个TableProvider 一个BeamSqlTable。仓库中所有内置数据源都遵循这一模式位于 provider 目录下每个子目录对应一个数据源avro、bigquery、bigtable、datastore、kafka、mongodb、parquet、pubsub、pubsublite、seqgen、text等。实战目标一个生成连续整数的 Streaming 数据表接下来实现一个完整的自定义数据源它基于 Beam SDK 的GenerateSequencePTransform源码见 GenerateSequence.java生成一条连续无界的整数流。最终效果是在 SQL CLI 中执行CREATE EXTERNAL TABLE sequenceTable -- 查询中使用的表别名 ( sequence BIGINT, -- 序列号 event_timestamp TIMESTAMP -- 生成事件的时间戳 ) TYPE sequence -- type 标识对应的表提供者 TBLPROPERTIES { elementsPerSecond : 12 } -- 每秒生成元素速率的可选参数然后即可查询SELECT sequence FROM sequenceTable;值得注意的是这个示例并非「纸上谈兵」——当前仓库的 seqgen 包中已经存在完整的、与本文讲解同构的生产实现下文代码即基于该目录下的真实源码。第一步实现GenerateSequenceTableProviderTableProvider只需完成两件事声明getTableType()返回的类型标识以及实现buildBeamSqlTable返回对应的表实现。源码见 GenerateSequenceTableProvider.javaAutoService(TableProvider.class) public class GenerateSequenceTableProvider extends InMemoryMetaTableProvider { Override public String getTableType() { return sequence; } Override public BeamSqlTable buildBeamSqlTable(Table table) { return new GenerateSequenceTable(table); } }这里有两个关键点AutoService(TableProvider.class)借助 Google AutoService 在编译期生成META-INF/services注册文件使该 provider 能被JdbcDriver/ SQL CLI 自动发现并加载无需手动注册。继承InMemoryMetaTableProvider该基类见 InMemoryMetaTableProvider.java把createTable、dropTable实现为 no-op因为 external 表不做持久化并把getTables()返回空集合让实现者只需关心类型标识与表构建。仓库中BigQueryTableProvider、TextTableProvider、KafkaTableProvider等均采用同样的继承方式可以直接对照阅读。第二步实现GenerateSequenceTableGenerateSequenceTable是真正干活的部分。源码见 GenerateSequenceTable.java完整实现如下class GenerateSequenceTable extends SchemaBaseBeamTable implements Serializable { public static final Schema TABLE_SCHEMA Schema.of(Field.of(sequence, FieldType.INT64), Field.of(event_time, FieldType.DATETIME)); Integer elementsPerSecond 5; GenerateSequenceTable(Table table) { super(TABLE_SCHEMA); if (table.getProperties().has(elementsPerSecond)) { elementsPerSecond table.getProperties().get(elementsPerSecond).asInt(); } } Override public PCollection.IsBounded isBounded() { return IsBounded.UNBOUNDED; } Override public PCollectionRow buildIOReader(PBegin begin) { return begin .apply(GenerateSequence.from(0).withRate(elementsPerSecond, Duration.standardSeconds(1))) .apply( MapElements.into(TypeDescriptor.of(Row.class)) .via(elm - Row.withSchema(TABLE_SCHEMA).addValues(elm, Instant.now()).build())) .setRowSchema(getSchema()); } Override public BeamTableStatistics getTableStatistics(PipelineOptions options) { return BeamTableStatistics.createUnboundedTableStatistics((double) elementsPerSecond); } Override public POutput buildIOWriter(PCollectionRow input) { throw new UnsupportedOperationException(buildIOWriter unsupported!); } }逐段拆解表结构TABLE_SCHEMA定义两列——sequenceFieldType.INT64即 BIGINT与event_timeFieldType.DATETIME即 TIMESTAMP。它通过SchemaBaseBeamTable见 SchemaBaseBeamTable.java提供getSchema()的默认实现。类的 Javadoc 表明Beam 中的每个 I/O 都有对应的表 Schema扩展SchemaBaseBeamTable即为此而生。TBLPROPERTIES 解析构造器从table.getProperties()中读取elementsPerSecond键默认值为每秒 5 个元素用户可以在TBLPROPERTIES里用 JSON 覆盖它如{ elementsPerSecond : 12 }。getProperties()返回的是 Jackson 的ObjectNode因此可以用has()判断存在性、asInt()取整数值。有界性isBounded()返回IsBounded.UNBOUNDED明确告诉 Beam 这是一条无限流必须以 Streaming 语义处理对应文档中「所有表都是 external」、读取无界源不会自然结束的行为。读路径buildIOReader这是整张表的核心。它依次执行两个 transformGenerateSequence.from(0).withRate(elementsPerSecond, Duration.standardSeconds(1))——从 0 开始、以每秒elementsPerSecond个的速度生成递增整数。从 GenerateSequence.java 源码可知from(long)要求起始值非负withRate(numElements, periodLength)要求元素数为正且周期非空两者会组合成底层CountingSource的限速配置。该类还支持to(long)限定最大值默认-1表示无界、withTimestampFn自定义元素时间戳、withMaxReadTime限制读取时长等均可按需组合。MapElements把每个Long映射为RowRow.withSchema(TABLE_SCHEMA).addValues(elm, Instant.now()).build()即序列号本身作为sequence列Instant.now()作为事件时间戳。最后调用setRowSchema(getSchema())让PCollectionRow携带 Schema供 SQL 引擎按列访问。统计信息getTableStatistics返回BeamTableStatistics.createUnboundedTableStatistics(elementsPerSecond)向优化器提供该无界表的元素速率估算。写路径buildIOWriter直接抛出UnsupportedOperationException——这是一个只读数据源明确表达「不支持写入」。从源码结构看原始博客中的示例直接继承BaseBeamTable与仓库当前实现继承SchemaBaseBeamTable并新增统计信息略有演进但核心思路完全一致实现一个有界性声明 读路径必要时写路径 Schema 的表类。在 SQL CLI 中体验新数据源构建并启动 SQL ShellBeam 官方提供了独立的 SQL shell 模块位于 shell 目录。根据仓库文档 shell.md从仓库根目录执行以下命令即可构建并启动./gradlew -p sdks/java/extensions/sql/shell -Pbeam.sql.shell.bundled:runners:flink:1.13,:sdks:java:io:kafka installDist ./sdks/java/extensions/sql/shell/build/install/shell/bin/shell要点说明-Pbeam.sql.shell.bundled参数用逗号分隔的 Gradle 项目 ID 列表决定把哪些 runner 与 I/O 组件打包进 shell。默认不带该参数只含DirectRunner上面示例额外引入了 Flink 1.13 runner 与 KafkaIO。首次构建需要几分钟因为 Gradle 要先编译所有依赖。启动后进入交互式环境提示符为0: BeamSQL查询默认由DirectRunner本地执行结果以表格形式返回。若希望以长跑任务方式提交到远程 runner可在 shell 内用SET runnerFlinkRunner;切换 runner并用SET projectId...;、SET tempLocation...;配置PipelineOptions提交后结果不再回显需到对应 runner 的 UI 查看。也可以使用distZip/distTar任务打包成独立发行物。注册并查询序列表在 shell 中执行CREATE EXTERNAL TABLE注册新数据源注意博客中演示的类型标识带引号TYPE sequence而仓库当前GenerateSequenceTableProvider.getTableType()返回sequence两种写法在语法上都可接受0: BeamSQL CREATE EXTERNAL TABLE input_seq ( . . . . . sequence BIGINT COMMENT this is the primary key, . . . . . event_time TIMESTAMP COMMENT this is the element timestamp . . . . . ) . . . . . TYPE sequence; No rows affected (0.005 seconds)随后执行SELECT * FROM input_seq LIMIT 5;0: BeamSQL SELECT * FROM input_seq LIMIT 5; --------------------------------- | sequence | event_time | --------------------------------- | 0 | 2019-05-21 00:36:33 | | 1 | 2019-05-21 00:36:33 | | 2 | 2019-05-21 00:36:33 | | 3 | 2019-05-21 00:36:33 | | 4 | 2019-05-21 00:36:33 | --------------------------------- 5 rows selected (1.138 seconds)这里LIMIT 5至关重要对于无界源如本文的sequence、以及内置的pubsub不加LIMIT的查询永远不会结束。这正是 shell.md 中「使用无界源开发」一节强调的开发要点——先用LIMIT x快速迭代验证逻辑确认无误后再去掉LIMIT作为长跑作业提交。用滚动窗口验证流式时间戳数据源是否正确提供了事件时间戳可以通过窗口聚合来验证——这正是本文实现中event_time Instant.now()的意义所在0: BeamSQL SELECT . . . . . COUNT(sequence) as elements, . . . . . TUMBLE_START(event_time, INTERVAL 2 SECOND) as window_start . . . . . FROM input_seq . . . . . GROUP BY TUMBLE(event_time, INTERVAL 2 SECOND) LIMIT 5; ----------------------------------- | elements | window_start | ----------------------------------- | 6 | 2019-06-05 00:39:24 | | 10 | 2019-06-05 00:39:26 | | 10 | 2019-06-05 00:39:28 | | 10 | 2019-06-05 00:39:30 | | 10 | 2019-06-05 00:39:32 | ----------------------------------- 5 rows selected (10.142 seconds)TUMBLE(event_time, INTERVAL 2 SECOND)按 2 秒的固定窗口对事件时间做滚动切分TUMBLE_START返回每个窗口的起始时间。可以看到除首个不完整窗口外每个 2 秒窗口恰好包含约 10 行对应默认速率每秒 5 个元素证明事件时间戳被正确写入、窗口聚合按预期工作。配合TBLPROPERTIES { elementsPerSecond : 12 }将速率提高到每秒 12 个元素后每窗口元素数会相应变为约 24。更多可借鉴的内置 Provider如果想把上述思路推广到真实数据源文件、数据库、消息队列仓库的 provider 目录提供了最直接的参考实现数据源Provider 实现读/写典型LOCATION/TBLPROPERTIESBigQueryBigQueryTableProvider.java读 写LOCATION [PROJECT_ID]:[DATASET].[TABLE]TBLPROPERTIES {method: DIRECT_READ}BigtableBigtableTableProvider.java读 写LOCATION googleapis.com/bigtable/projects/...KafkaKafkaTableProvider.java读 写LOCATION broker:port/topicTBLPROPERTIES中可配bootstrap_servers、topics、formatcsv/avro/json/proto/thriftPub/SubPubsubTableProvider.java读 写LOCATION projects/[PROJECT]/topics/[TOPIC]MongoDBMongoDbTableProvider.java读 写LOCATION mongodb://[HOST]:[PORT]/[DATABASE]/[COLLECTION]文本文件TextTableProvider.java读 写LOCATION /path/to/fileTBLPROPERTIES中format可选default/rfc4180/excel/tdf/mysqlAvro / ParquetAvroTableProvider.java、ParquetTableProvider.java读文件路径型LOCATIONPub/Sub LitePubsubLiteTableProvider.java读 写topic / subscription 型LOCATION这些实现展示了不同复杂度的模式有的直接继承InMemoryMetaTableProvider并手写BeamSqlTable如TextTableProvider、KafkaTableProvider有的通过SchemaIOTableProviderWrapper包装通用的SchemaIOProvider如PubsubTableProvider。无论哪种getTableType()返回的字符串就是用户 SQL 中TYPE的值且都依赖AutoService(TableProvider.class)实现自动发现。更完整的各数据源语法与参数说明可直接查阅 create-external-table.md。总结接入一个新数据源的完整清单综合前文向 Beam SQL / Beam SQL CLI 添加一个新的数据源只需完成三步实现BeamSqlTable继承SchemaBaseBeamTable或直接实现BeamSqlTable接口定义 Schema实现isBounded()、buildIOReader(PBegin)只读源在buildIOWriter中抛出UnsupportedOperationException并尽量提供getTableStatistics帮助优化器做速率/行数估算。实现TableProvider继承InMemoryMetaTableProvider用getTableType()声明 SQL 中的TYPE标识用buildBeamSqlTable(Table)把解析好的元数据列、LOCATION、TBLPROPERTIES翻译成上一步的表实例。标注AutoService(TableProvider.class)并随 shell 一起构建这样 SQL CLI 启动时会自动加载该 provider用户即可执行CREATE EXTERNAL TABLE ... TYPE 你的类型完成注册并像使用内置数据源一样用SELECT、JOIN、INSERT INTO操作它。本质上Beam SQL 的数据源扩展就是一个「声明式元数据Table 命令式读写BeamSqlTable 自动发现AutoService」的组合。掌握了seqgen这个最小可运行的样例再对照 provider 目录中各类内置实现无论是接入一个自研存储还是包装一个第三方 I/O都能快速落地。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam SQL CLI 数据源扩展实战从零实现自定义 TableProvider 接入新数据源Apache Beam SQL CLI 数据源扩展实战从零实现自定义 TableProvider 接入新数据源 Apache Beam 为批处理与流处理提供了大数据批处理流处理数据工程Apache Beam 深度解析为 Beam SQL CLI 添加新数据源的 TableProvider 扩展机制Apache Beam 深度解析为 Beam SQL CLI 添加新数据源的 TableProvider 扩展机制 本文以 Apache Beam 官方博客《批处理流处理大数据自定义SQL函数基于MyBatis-Plus扩展数据库特定函数支持自定义SQL函数基于MyBatis Plus扩展数据库特定函数支持 引言为什么需要自定义SQL函数 在日常开发中我们经常会遇到这样的场景项目需要支持特后端ORM上一篇嵌入式系统测试终极指南5大策略提升单元测试与集成测试效率下一篇终极指南如何使用ipsw AI反编译器集成Claude到OpenAI的5种AI模型创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表