Apache Beam SQL 完全指南:用 SQL 查询有界与无界 PCollection

发布时间:2026/10/12 2:10:11
Apache Beam SQL 完全指南:用 SQL 查询有界与无界 PCollection
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载Apache Beam SQL 是 Apache Beam 提供的声明式查询层它允许开发者目前支持 Beam Java 与 Beam Python用标准 SQL 语句直接查询有界bounded与无界unbounded的PCollection并将查询翻译为一个PTransform融入既有 Beam 管道。本文基于当前仓库的官方文档与源码系统讲解 Beam SQL 的两种方言、SqlTransform与Row两大核心概念、方言切换方式、交互式 SQL Shell 以及窗口化、UDF/UDAF、外部表等扩展能力帮助读者将 SQL 无缝嵌入 Beam 批流一体的数据管道。Beam SQL 是什么SQL 即 PTransformBeam SQL 的核心设计思想是把一段 SQL 查询字符串当作一个可复用的PTransform。在管道中任何需要写ParDo、GroupByKey等变换的地方都可以换成SqlTransform.query(...)并且可以自由地与普通PTransform混合编排——SQL 变换的输入输出都是标准的PCollection不存在“SQL 专用通道”。这一设计在源码中有直接体现SqlTransform.java 的类注释将其定义为 the DSL interface of Beam SQL. It translates a SQL query as a PTransformSqlTransform本身继承自PTransformPInput, PCollectionRow见 SqlTransform.java因此它天然可以apply到PCollection或PCollectionTuple上输出也必然是PCollectionRow。Beam SQL 整体上包含两大方言dialect均由当前仓库实现Beam Calcite SQLApache Calcite 的一个变体。Calcite 在大数据处理领域被广泛采用因此其 SQL 语法、函数与操作符生态成熟Beam Calcite SQL 是 Beam SQL 的默认方言。对应的查询规划器实现为CalciteQueryPlanner从 BeamSqlPipelineOptions.java 可以看出plannerName的默认值正是org.apache.beam.sdk.extensions.sql.impl.CalciteQueryPlanner。Beam ZetaSQLGoogle ZetaSQL 语言的一个变体与 BigQuery 的 SQL 框架高度兼容因此尤其适合读写 BigQuery 表的管道底层对接BigQueryIO。其查询规划器实现位于 ZetaSQLQueryPlanner.java。核心概念SqlTransform 与 Row要使用 Beam SQL必须理解两个概念1. SqlTransform从 SQL 字符串创建 PTransformSqlTransform.query(queryString)是创建PTransform的唯一入口定义于 SqlTransform.java。它可应用于单个PCollection查询中可用表名PCOLLECTION引用该集合常量定义见 SqlTransform.javaPCollectionTuple元组中每个PCollection的TupleTag的 ID 即为其在查询中的表名。注意这些表名绑定于特定的PCollectionTuple只在作用于它的查询上下文中有效。从实现上看expand() 方法 会把单个PCollection包装为名为PCOLLECTION的表ReadOnlyTableProvider把PCollectionTuple展开后以 tag 名注册为多张表toTableMap()随后交由BeamSqlEnv完成解析与执行计划生成。2. RowBeam SQL 操作的元素类型Row是 Beam SQL 操作的元素类型一个PCollectionRow在 SQL 语境中扮演“表”的角色。SQL 查询只能应用于已注册 schema 的PCollectionT或PCollectionRow关于如何为类型T注册 schema可参考 Beam Programming Guide 中的 schema 章节。切换方言通过 PipelineOptions 设置 Planner默认使用 Calcite 方言若需要 ZetaSQL通过PipelineOptions设置plannerName即可官方示例代码如下PipelineOptions options ...; options .as(BeamSqlPipelineOptions.class) .setPlannerName(org.apache.beam.sdk.extensions.sql.zetasql.ZetaSQLQueryPlanner);其背后机制是SqlTransform.expand()在构建BeamSqlEnv时会读取管道级BeamSqlPipelineOptions中的getPlannerName()作为查询规划器见 SqlTransform.java。也就是说方言选择是全局管道级的配置且可以在单个变换上用withQueryPlannerClass(Class)覆盖见 SqlTransform.java。此外BeamSqlPipelineOptions.java 还提供了zetasqlDefaultTimezoneZetaSQL 分析器默认时区默认UTC对 Calcite 无影响与verifyRowValues调试用行值校验默认关闭两个辅助选项。依赖注意使用ZetaSQLQueryPlanner需要在beam-sdks-java-extensions-sqlCalciteQueryPlanner所在包之外额外引入beam-sdks-java-extensions-sql-zetasql依赖。两种方言的内置函数清单各有独立来源Calcite 的支持情况见 calcite/overview.mdZetaSQL 的内置函数由 SupportedZetaSqlBuiltinFunctions.java 列出注释掉的部分表示尚未支持。实操入门从 Row 构造到 SQL 查询第一步获取 PCollectionRowSQL 查询只能作用在带 schema 的集合上。若没有现成的类型T获取PCollectionRow有三种典型途径途径一从内存数据创建常用于单元测试。必须显式指定Row的 coder示例中通过Create.of(..)完成// 为记录定义 schema Schema appSchema Schema .builder() .addInt32Field(appId) .addStringField(description) .addDateTimeField(rowtime) .build(); // 构造一条符合该 schema 的具体 Row Row row Row .withSchema(appSchema) .addValues(1, Some cool app, new Date()) .build(); // 创建仅包含该行的源 PCollection PCollectionRow testApps PBegin .in(p) .apply(Create .of(row) .withCoder(RowCoder.of(appSchema)));途径二从其他类型的PCollectionT转换即T不是Row通过ParDo将输入记录映射为Row// 一个示例 POJO 类 class AppPojo { Integer appId; String description; Date timestamp; } // 通过某种方式获取 POJO 集合 PCollectionAppPojo pojos ... // 通过 DoFn 转换为与上述 schema 一致的 Rows PCollectionRow apps pojos .apply( ParDo.of(new DoFnAppPojo, Row() { ProcessElement public void processElement(ProcessContext c) { AppPojo pojo c.element(); Row appRow Row .withSchema(appSchema) .addValues( pojo.appId, pojo.description, pojo.timestamp) .build(); c.output(appRow); } })).setRowSchema(appSchema);途径三作为另一个SqlTransform的输出从而形成 SQL 链式管道详见下文。第二步用 SqlTransform 执行查询场景 A作用于单个 PCollection在查询中通过表名PCOLLECTION引用PCollectionRow filteredNames testApps.apply( SqlTransform.query( SELECT appId, description, rowtime FROM PCOLLECTION WHERE appId1));场景 B作用于 PCollectionTuple实现多表 JOIN。以“按 app 统计评论总数与平均评分”为例// 创建 reviews 的 schema Schema reviewSchema Schema .builder() .addInt32Field(appId) .addInt32Field(reviewerId) .addFloatField(rating) .addDateTimeField(rowtime) .build(); // 获取符合该 schema 的评论记录 PCollectionRow reviewsRows ... // 将两个 PCollection 放入元组TupleTag ID 作为 SQL 中的表名 PCollectionTuple namesAndFoods PCollectionTuple .of(new TupleTag(Apps), appsRows) // appsRows 来自上一示例 .and(new TupleTag(Reviews), reviewsRows); // 通过 JOIN 两个 PCollection 计算每个 app 的评论总数与平均评分 PCollectionRow output namesAndFoods.apply( SqlTransform.query( SELECT Apps.appId, COUNT(Reviews.rating), AVG(Reviews.rating) FROM Apps INNER JOIN Reviews ON Apps.appId Reviews.appId GROUP BY Apps.appId));仓库中的 BeamSqlExample.java 完整演示了两种 API 的用法先用Create.of(...).withRowSchema(type)构造内存输入再对PCOLLECTION执行SELECT c1, c2, c3 FROM PCOLLECTION WHERE c1 1最后把结果放入PCollectionTupletag 为CASE1_RESULT执行GROUP BY聚合。可直接用以下命令在本地 DirectRunner 运行./gradlew :sdks:java:extensions:sql:runBasicExamplePython 端的 Beam SQLBeam Python 通过跨语言xlang方式复用 Java 的SqlTransform。在 sql.py 中SqlTransform继承自ExternalTransformURN 为beam:external:java:sql:v1通过默认的 expansion service:sdks:java:extensions:sql:expansion-service:shadowJar展开为 Java 变换。输入PCollection必须带 schema可通过注册typing.NamedTuple类型的RowCoder或直接产出beam.Row实现purchases | SqlTransform( SELECT item_name, COUNT(*) AS count FROM PCOLLECTION GROUP BY item_name) | beam.Map(lambda row: Weve sold %d %ss! % (row.count, row.item_name))如需 ZetaSQL 方言构造时传入dialectzetasql。源码级原理一条 SQL 查询如何变成 Beam 管道结合 SqlTransform.expand() 的实现一条 SQL 查询的展开链路可以概括为输入注册将PCollection以PCOLLECTION为名或PCollectionTuple以 tag 名为名包装为BeamPCollectionTable挂载到InMemoryMetaStore函数注册通过registerUdf/registerUdaf注册的用户函数UDF/UDAF被写入BeamSqlEnv若开启autoLoading默认开启还会通过ServiceLoader自动加载TableProvider实现规划器选择优先使用withQueryPlannerClass指定的类否则读取管道选项中的getPlannerName()DDL 执行withDdlString(...)传入的CREATE EXTERNAL TABLE等 DDL 语句在BeamSqlEnv中逐一执行解析与转换sqlEnv.parseQuery(...)解析 SQL最终由BeamSqlRelUtils.toPCollection(...)生成执行计划并输出PCollectionRow。这种“SQL 字符串 → RelNode → Beam PCollection”的链路正是 Beam SQL 能把声明式查询与命令式变换统一在同一个管道模型中的根本原因。Beam SQL Shell无需 Java SDK 的交互式查询自 Beam 2.6.0 起Beam SQL 附带交互式 shell即 Beam SQL shell允许不写 Java SDK 代码、直接以 SQL 查询构建管道。默认使用DirectRunner在本地执行。详细说明见 shell.md。快速启动克隆 Beam 仓库后在仓库根目录依次执行./gradlew -p sdks/java/extensions/sql/shell -Pbeam.sql.shell.bundled:runners:flink:1.17,:sdks:java:io:kafka installDist ./sdks/java/extensions/sql/shell/build/install/shell/bin/shell注若此前未构建过项目首次执行 Gradle 命令会因需要构建全部依赖而耗时数分钟。启动后即可输入查询shell 会把查询编译为 Beam 管道、用DirectRunner执行并在管道结束时以表格形式返回结果0: BeamSQL SELECT foo AS NAME, bar AS TYPE, num AS NUMBER; -------------------- | NAME | TYPE | NUMBER | -------------------- | foo | bar | num | -------------------- 1 row selected (0.826 seconds)声明表CREATE EXTERNAL TABLE读写数据前必须先用CREATE EXTERNAL TABLE声明一张虚拟表。例如为当前目录下的本地 CSV 文件test-file.csv建表0: BeamSQL CREATE EXTERNAL TABLE csv_file (field1 VARCHAR, field2 INTEGER) TYPE text LOCATION test-file.csv; No rows affected (0.042 seconds)该语句只向 Beam SQL 描述数据源/汇schema 与位置不会直接创建持久化的物理表表会在写入发生时按需创建。完整的语法、TYPE可选值bigquery、bigtable、pubsub、pubsublite、kafka、mongodb、text、LOCATION与TBLPROPERTIES说明见 create-external-table.md。读写数据读取上一步声明的 CSV 表0: BeamSQL SELECT field1 AS field FROM csv_file;写入则使用INSERT INTO … SELECT ...0: BeamSQL INSERT INTO csv_file SELECT foo, bar;读写行为取决于表类型例如text表基于TextIO实现写入可能产生多个编号文件pubsub表是无界源读操作永远不会自然结束。开发无界源务必加 LIMIT检查无界源数据时必须在SELECT末尾加LIMIT x限制输出条数否则管道永不结束0: BeamSQL SELECT field1 FROM unbounded_source LIMIT 10 ;本地快速查询适合探查数据、迭代设计确认 SQL 逻辑无误后可去掉LIMIT将语句作为长跑任务提交此时若含无界源管道可能无限运行。指定 Runner 与 PipelineOptions默认DirectRunner之外的 runner 需要两步构建 shell 时用-Pbeam.sql.shell.bundled打入目标 runner 与额外组件逗号分隔可打包多个也可追加更多 I/O例如./gradlew -p sdks/java/extensions/sql/shell -Pbeam.sql.shell.bundled:runners:flink:1.17,:sdks:java:io:kafka installDist在 shell 内用SET命令指定 runner0: BeamSQL SET runnerFlinkRunner;此后所有INSERT语句都会作为管道提交到该 runnershell 不再直接回显查询结果需要通过对应 runner 的 UI 或命令行管理作业。配置 runner 所需的PipelineOptions同样用SET设置例如0: BeamSQL SET projectIdgcpProjectId; 0: BeamSQL SET tempLocation/tmp/tempDir;SET/RESET语句的语法SET option valueflag 类选项值用trueRESET option恢复默认值见 set.md。打包独立分发可用distZip或distTar任务构建 shell 的独立发布包./gradlew -p sdks/java/extensions/sql/shell -Pbeam.sql.shell.bundled:runners:flink:1.17,:sdks:java:io:kafka distZip ls ./sdks/java/extensions/sql/shell/build/distributions/ beam-sdks-java-extensions-sql-shell-2.6.0-SNAPSHOT.tar beam-sdks-java-extensions-sql-shell-2.6.0-SNAPSHOT.zipBeam SQL 扩展批流模型与复杂类型的 SQL 化Beam SQL 基于 Beam 统一的批/流模型与复杂数据类型处理能力提供了一系列对所有方言通用的扩展完整索引见 overview.md 及 extensions 目录下的各参考页。窗口化与触发Windowing and triggering可以在两个层面使用窗口语义一是在输入PCollection上预先配置 windowing二是在 SQL 的GROUP BY中直接使用窗口函数此时会覆盖输入集合的窗口设置。触发trigger只能通过输入PCollection设置SQL 侧没有对应扩展。使用窗口函数时要求有TIMESTAMP字段作为行的事件时间。支持的窗口函数见 windowing-and-triggering.md-- TUMBLE固定窗口时长 1 小时 SELECT f_int, COUNT(*) FROM PCOLLECTION GROUP BY f_int, TUMBLE(f_timestamp, INTERVAL 1 HOUR) -- HOP滑动窗口每 30 分钟滑动一次、时长 1 小时 SELECT f_int, COUNT(*) FROM PCOLLECTION GROUP BY f_int, HOP(f_timestamp, INTERVAL 30 MINUTE, INTERVAL 1 HOUR) -- SESSION会话窗口间隔 5 分钟 SELECT f_int, COUNT(*) FROM PCOLLECTION GROUP BY f_int, SESSION(f_timestamp, INTERVAL 5 MINUTE)注意查询中未指定窗口函数时输入PCollection的窗口策略保持不变指定后窗口会被更新但 trigger 不变。用户自定义函数UDF 与 UDAF当内置函数不满足需求时可以用 Java 编写标量函数UDF与聚合函数UDAF并在 SQL 中调用详见 user-defined-functions.md。UDF 可以是任意“输入零或多个标量字段、返回一个标量值”的 Java 方法或一个SerializableFunctionpublic static class CubicInteger implements BeamSqlUdf { public static Integer eval(Integer input){ return input * input * input; } } public static class CubicIntegerFn implements SerializableFunctionInteger, Integer { Override public Integer apply(Integer input) { return input * input * input; } } String sql SELECT f_int, cubic1(f_int), cubic2(f_int) FROM PCOLLECTION WHERE f_int 2; PCollectionRow result input.apply( udfExample, SqlTransform .query(sql) .registerUdf(cubic1, CubicInteger.class) .registerUdf(cubic2, new CubicIntegerFn()));UDAF 则接受CombineFn通过registerUdaf注册例如实现“平方和”聚合public static class SquareSum extends CombineFnInteger, Integer, Integer { Override public Integer createAccumulator() { return 0; } Override public Integer addInput(Integer accumulator, Integer input) { return accumulator input * input; } Override public Integer mergeAccumulators(IterableInteger accumulators) { ... } Override public Integer extractOutput(Integer accumulator) { return accumulator; } } PCollectionRow result input.apply( udafExample, SqlTransform .query(SELECT f_int1, squaresum(f_int2) FROM PCOLLECTION GROUP BY f_int1) .registerUdaf(squaresum, new SquareSum()));JOIN 支持范围Beam SQL 支持INNER、LEFT OUTER、RIGHT OUTER三种 JOIN且仅支持等值连接equijoinCROSS JOIN无ON的全笛卡尔积与FULL OUTER JOIN不受支持。依据输入有界性分为三种场景详见 joins.md有界 JOIN 有界使用标准 JOIN 实现不做窗口与触发无界 JOIN 无界输入窗口必须兼容否则抛IllegalArgumentException每个输入的触发只能每窗口触发一次目前仅支持DefaultTrigger且零允许迟到否则抛UnsupportedOperationException即按窗口做 JOIN无界 JOIN 有界有界输入被当作 side-input 处理窗口/触发继承自上游。Calcite 方言函数支持概览Beam Calcite SQL 对 Apache Calcite 运算符/函数的支持状态见 calcite/overview.md比较、逻辑、算术、字符串、日期时间、条件、类型转换等类别均受支持具体函数清单分别对应 scalar-functions、aggregate-functions 等参考页二进制字符串、系统函数、窗口函数、分组函数、空间/几何、JSON 函数、MATCH_RECOGNIZE等暂不支持。ZetaSQL 方言的语法、词法、数据类型、运算符及数学/字符串/聚合函数等参考页位于 zetasql 目录。小结Beam SQL 把 SQL 的声明式表达力与 Beam 的统一批流执行模型结合起来SqlTransform将查询字符串封装为标准PTransformPCollectionRow即 SQL 意义上的“表”PlannerName选项在 Calcite默认与 ZetaSQLBigQuery 友好之间自由切换交互式 SQL Shell 让开发者免去 Java 编码即可快速迭代管道而窗口化、UDF/UDAF、外部表与 JOIN 等扩展则覆盖了从本地 CSV 探查到无界流实时计算的典型场景。继续深入可阅读仓库中的 walkthrough.md完整示例、shell.mdshell 使用以及 extensions 目录各类扩展参考并对照 SqlTransform.java 与 BeamSqlExample.java 动手验证。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam SQL 完全指南用标准 SQL 查询有界与无界 PCollectionApache Beam SQL 完全指南用标准 SQL 查询有界与无界 PCollection Apache Beam SQL 是 Beam 内置的 SQLApache Beam SQL 指南用标准 SQL 查询有界与无界 PCollectionApache Beam SQL 指南用标准 SQL 查询有界与无界 PCollection Beam SQL 是 Apache Beam 内置的 SQL 方言大数据批处理流处理数据工程Apache Beam SQL Walkthrough用 SqlTransform 对 PCollectionRow 执行 SQL 查询的完整实战指南Apache Beam SQL Walkthrough用 SqlTransform 对 PCollectionRow 执行 SQL 查询的完整实战指南 本篇批处理流处理大数据上一篇《鸣潮》自动化工具让重复性游戏操作从负担变为乐趣下一篇ok-ww鸣潮自动化工具从零到精通的深度实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考