【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
Apache Beam 的 SQL 扩展允许直接用 SQL 编写数据处理管道,无论是 Batch 还是 Streaming 场景。本文以 Beam SQL CLI(BeamSQL交互式 shell)为背景,围绕「如何给 Beam SQL 添加一个全新的数据源」这一核心主题展开:你会完整掌握TableProvider与BeamSqlTable两个关键抽象的职责分工,亲手实现一个基于GenerateSequence的「连续整数流」数据表,并学会在 SQL CLI 中通过CREATE EXTERNAL TABLE注册它、用SELECT与窗口聚合查询它。文中的实现代码、命令与配置均取自当前仓库,可直接对照源码阅读与验证。
背景:Beam SQL 与CREATE EXTERNAL TABLE
Beam 的 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与BeamSqlTable
SqlTransform在遇到CREATE EXTERNAL TABLE语句时,会依赖一组TableProvider来完成建表元数据操作与实际数据读写。TableProvider的接口定义位于 TableProvider.java,它负责:
| 方法 | 职责 |
|---|---|
getTableType() | 返回该 provider 处理的表类型标识(对应 SQL 中TYPE的值) |
createTable(Table)/dropTable(String) | 表元数据的创建与删除(CRUD) |
getTables()/getTable(name) | 列出 / 获取该 provider 管理的所有表 |
buildBeamSqlTable(Table) | 依据元数据构建真正的读写实现,返回BeamSqlTable |
而从接口注释可以确认:所有标注了@AutoService(TableProvider.class)的实现,都会在使用默认连接参数启动JdbcDriver(SQL CLI 正是如此)时被自动加载。
BeamSqlTable接口(见 BeamSqlTable.java)则定义了表的具体读写行为:
buildIOReader(PBegin):从数据源构建一个PCollection<Row>(读路径);buildIOWriter(PCollection<Row>):把PCollection<Row>写入目标(写路径);isBounded():返回IsBounded.BOUNDED或IsBounded.UNBOUNDED,决定 SQL 优化器按批还是按流处理;getSchema():返回表的Schema;getTableStatistics(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 包中已经存在完整的、与本文讲解同构的生产实现,下文代码即基于该目录下的真实源码。
第一步:实现GenerateSequenceTableProvider
TableProvider只需完成两件事:声明getTableType()返回的类型标识,以及实现buildBeamSqlTable返回对应的表实现。源码见 GenerateSequenceTableProvider.java:
@AutoService(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等均采用同样的继承方式,可以直接对照阅读。
第二步:实现GenerateSequenceTable
GenerateSequenceTable是真正干活的部分。源码见 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 PCollection<Row> 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(PCollection<Row> input) { throw new UnsupportedOperationException("buildIOWriter unsupported!"); } }逐段拆解:
- 表结构:
TABLE_SCHEMA定义两列——sequence(FieldType.INT64,即 BIGINT)与event_time(FieldType.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:这是整张表的核心。它依次执行两个 transform:GenerateSequence.from(0).withRate(elementsPerSecond, Duration.standardSeconds(1))——从 0 开始、以每秒elementsPerSecond个的速度生成递增整数。从 GenerateSequence.java 源码可知:from(long)要求起始值非负;withRate(numElements, periodLength)要求元素数为正且周期非空,两者会组合成底层CountingSource的限速配置。该类还支持to(long)限定最大值(默认-1表示无界)、withTimestampFn自定义元素时间戳、withMaxReadTime限制读取时长等,均可按需组合。MapElements把每个Long映射为Row:Row.withSchema(TABLE_SCHEMA).addValues(elm, Instant.now()).build(),即序列号本身作为sequence列,Instant.now()作为事件时间戳。最后调用setRowSchema(getSchema())让PCollection<Row>携带 Schema,供 SQL 引擎按列访问。
- 统计信息:
getTableStatistics返回BeamTableStatistics.createUnboundedTableStatistics(elementsPerSecond),向优化器提供该无界表的元素速率估算。 - 写路径:
buildIOWriter直接抛出UnsupportedOperationException——这是一个只读数据源,明确表达「不支持写入」。
从源码结构看,原始博客中的示例(直接继承BaseBeamTable)与仓库当前实现(继承SchemaBaseBeamTable并新增统计信息)略有演进,但核心思路完全一致:实现一个有界性声明 + 读路径(必要时写路径)+ Schema 的表类。
在 SQL CLI 中体验新数据源
构建并启动 SQL Shell
Beam 官方提供了独立的 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 runner='FlinkRunner';切换 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/TBLPROPERTIES |
|---|---|---|---|
| BigQuery | BigQueryTableProvider.java | 读 + 写 | LOCATION '[PROJECT_ID]:[DATASET].[TABLE]';TBLPROPERTIES '{"method": "DIRECT_READ"}' |
| Bigtable | BigtableTableProvider.java | 读 + 写 | LOCATION 'googleapis.com/bigtable/projects/...' |
| Kafka | KafkaTableProvider.java | 读 + 写 | LOCATION 'broker:port/topic';TBLPROPERTIES中可配bootstrap_servers、topics、format(csv/avro/json/proto/thrift) |
| Pub/Sub | PubsubTableProvider.java | 读 + 写 | LOCATION 'projects/[PROJECT]/topics/[TOPIC]' |
| MongoDB | MongoDbTableProvider.java | 读 + 写 | LOCATION 'mongodb://[HOST]:[PORT]/[DATABASE]/[COLLECTION]' |
| 文本文件 | TextTableProvider.java | 读 + 写 | LOCATION '/path/to/file';TBLPROPERTIES中format可选default/rfc4180/excel/tdf/mysql |
| Avro / Parquet | AvroTableProvider.java、ParquetTableProvider.java | 读 | 文件路径型LOCATION |
| Pub/Sub Lite | PubsubLiteTableProvider.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,都能快速落地。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
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
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考