- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文围绕 Apache Beam 中 Schema(模式)这一核心概念展开:Schema 是一种语言无关的类型定义,用于描述PCollection中元素的结构,让 Beam 能够理解数据的字段布局、类型与嵌套关系。本文将从 Schema 的定义与类型体系出发,深入当前仓库的源码实现(如 Schema.java 中的TypeName枚举与LogicalType接口),讲解 Schema 如何被附加到PCollection(来源自动绑定、SchemaRegistry注册、SchemaCoder/RowCoder编解码),并完整覆盖 Schema Transforms 提供的字段选择、分组聚合、连接、过滤、字段增删、重命名、类型转换与增强版 ParDo 等关键能力,最后通过SqlTransform给出可直接复用的实战示例。读完本文,你将掌握 Beam Schema 从定义、绑定到变换的完整技术链路,能够在 Java、Python 等多语言 SDK 中直接构建并操作带 Schema 的管道。
什么是 Apache Beam Schema
在 Apache Beam 中,Schema 是PCollection的一种语言无关(language-independent)类型定义。一个 Schema 将该PCollection的元素定义为一个有序的命名字段列表(ordered list of named fields)——字段的先后顺序、字段名与字段类型共同构成了元素的完整类型描述。
这意味着 Beam 管道不再只是处理"不透明"的元素对象,而是能够理解每个元素内部的结构:哪个字段叫什么名字、是什么类型、是否可空、是原子类型还是嵌套结构。这种结构化的理解是后续所有 Schema Transforms(字段选择、过滤、连接、聚合等)能够自动化的前提。
为什么需要 Schema:可内省的结构化数据
在大多数实际业务中,PCollection中的元素类型本身具有可以被内省(introspect)的结构。典型的例子包括:
- JSON对象
- Protocol Buffer消息
- Avro记录
- 数据库行对象(database row objects)
所有这些格式都可以被转换为 Beam Schema。以仓库中的实现为例,Java SDK 提供了RowJson(见 RowJson.java)与 JsonUtils.java 等工具,用于在 JSON 与带 Schema 的Row之间互转;RowCoder(见 RowCoder.java)则负责按 Schema 对Row进行高效编码与解码。
Schema 的附加方式
要利用 Schema 的能力,PCollection必须先挂上 Schema。在大多数情况下,数据源(source)本身就会为PCollection附加 Schema——例如读取 Avro 文件、BigQuery 表或数据库结果时,Beam 会依据外部格式的元数据自动推导出 Schema,你无需手工声明。
除此之外,还可以通过Schema Registry显式地为用户自定义类型注册 Schema。仓库中的 SchemaRegistry.java 提供了完整 API:
createDefault():创建默认 Registry,内置对常见 Java 类型(POJO、JavaBean、AutoValue 等)的 Schema 推断支持;registerSchemaForClass(...)/registerSchemaForType(...):为具体类型注册 Schema;registerSchemaProvider(...):注册自定义的 SchemaProvider(提供者的抽象,用于描述"如何从一个类型推导/获取 Schema");getSchema(...):按类型查询已注册的 Schema,找不到时抛出NoSuchSchemaException;getSchemaCoder(...):直接获取与该类型关联的 SchemaCoder。
从 SchemaCoder.java 可以看到,SchemaCoder是CustomCoder的子类,of(Schema)与coderForFieldType(FieldType)等静态方法负责构造"按 Schema 编解码"的 Coder,配合RowCoder完成实际序列化。因此,当元素带 Schema 时,Beam 可以自动为其选择合适的 Coder,无需手工指定。
Schema 的类型体系:从原子类型到嵌套结构
Schema 的核心数据类型定义在 Schema.java 中。其内部TypeName枚举(第 522 行起)列出了全部类型构造器(type constructor):
| 类型类别 | 取值 | 说明(来自源码注释) |
|---|---|---|
| 整数 | BYTE/INT16/INT32/INT64 | 1/2/4/8 字节有符号整数 |
| 高精度数值 | DECIMAL | 任意精度十进制数 |
| 浮点 | FLOAT/DOUBLE | 单精度 / 双精度浮点 |
| 字符串 | STRING | 字符串 |
| 日期时间 | DATETIME | 日期和时间 |
| 布尔 | BOOLEAN | 布尔值 |
| 字节 | BYTES | 字节数组 |
| 集合 | ARRAY/ITERABLE | 数组;Iterable 与 Array 不同,可能无法整体装入内存 |
| 键值对 | MAP | 映射 |
| 嵌套行 | ROW | 字段本身是一个嵌套 Row |
| 自定义 | LOGICAL_TYPE | 用户自定义逻辑类型 |
该枚举还定义了若干类型分组的语义集合:NUMERIC_TYPES(BYTE、INT16、INT32、INT64、DECIMAL、FLOAT、DOUBLE)、STRING_TYPES、DATE_TYPES、COLLECTION_TYPES(ARRAY、ITERABLE)、MAP_TYPES、COMPOSITE_TYPES(ROW),并通过isPrimitiveType()、isNumericType()、isCollectionType()、isCompositeType()等辅助方法进行判断。
数值类型的隐式拓宽
源码中isSupertypeOf(TypeName other)定义了数值类型之间的兼容关系,这也是类型转换(Cast)与 Schema 兼容性判断的基础:
INT16是BYTE的超类型;INT32是BYTE、INT16的超类型;INT64是BYTE、INT16、INT32的超类型;DOUBLE是FLOAT的超类型;DECIMAL是FLOAT、DOUBLE的超类型。
嵌套与容器字段
Schema.Builder提供了丰富的字段添加方法(见 Schema.java),除了addInt32Field、addStringField等原子字段外,还支持:
addNullableField/ 各类addNullableXxxField:声明可空字段;addArrayField/addIterableField:集合字段;addMapField:键值对字段(需同时给出 key 与 value 类型);addRowField:嵌套 Row 字段(需要传入子 Schema);addLogicalTypeField:挂接自定义逻辑类型。
LogicalType:自定义类型机制
Schema.LogicalType<InputT, BaseT>接口(Schema.java)允许用户定义全新的 Schema 类型:
getIdentifier():返回全局唯一的类型标识符;getBaseType():声明底层存储所用的基础FieldType(通常为标准类型,最终必须解析到标准 Schema 类型,且不允许递归引用);toBaseType(InputT)/toInputType(BaseT):在用户 Java 类型与底层存储类型之间双向转换;- 可选
getArgument():为类型提供配置参数(例如定长字节数组的长度)。
仓库的 logicaltypes 目录内置了Date、DateTime、EnumerationType、FixedBytes、MicrosInstant、NanosDuration、OneOfType、SqlTypes、UuidLogicalType等丰富的逻辑类型实现,可直接使用。
Schema 的相等性与兼容性
Schema是不可变对象,并缓存了 hashCode(见 Schema.java)。其相等性与兼容性语义在源码中有明确定义:
equals(...):字段同名、同序、同类型且 Options 一致才算相等;若两个 Schema 均带 UUID 则直接比较 UUID(每个SchemaCoder都有 UUID,同 UUID 的 Schema 必然相等,可短路比较);typesEqual(...):忽略字段名与描述,仅比较类型;equivalent(...):字段可以顺序不同(按字段名排序后比较),可配合EquivalenceNullablePolicy(SAME/WEAKEN/IGNORE)控制是否把可空性纳入等价判断;assignableTo(...):采用WEAKEN策略判断是否可赋值。
此外,Schema内部维护fieldIndices(字段名到索引的双向映射)与encodingPositions(编码位置),保证按字段名访问与按位置编码都能高效完成。
Schema 与 SDK 语言的自然嵌入
Schema 虽然语言无关,但设计上被刻意做成了"自然地嵌入到各 Beam SDK 编程语言中"的形式,让你可以继续使用原生类型,同时享受 Beam 理解元素 Schema 带来的好处:
- Java:通过
SchemaProvider自动从 POJO、JavaBean、AutoValue 等类型推导 Schema。仓库的 schemas 目录中,AutoValueSchema.java、JavaBeanSchema.java、JavaFieldSchema.java 以及GetterBasedSchemaProvider分别实现了对不同 Java 类型风格的 Schema 推导;annotations子目录还提供了@DefaultSchema、@SchemaCreate、@SchemaFieldName、@SchemaIgnore等注解,用于精确控制 Schema 的生成。 - Python:
apache_beam的 schemas.py 实现了类型与 Schema 的双向翻译,例如named_tuple_to_schema将 Python namedtuple 转换为 Schema,typing_to_runner_api/typing_from_runner_api在 Python 类型与 Runner API 表示之间转换,让你在 Python 管道中用原生类型(如typing.NamedTuple、typing.Optional)声明带 Schema 的PCollection。 - Go / TypeScript:同样提供对原生类型的 Schema 映射支持,保证跨语言管道(Cross-language)中 Schema 语义一致。
这种"原生类型 + 自动 Schema"的设计,意味着你可以写出直观的领域模型代码,却让 Beam 底层获得完全结构化的数据视图,从而解锁各类通用 Schema 变换。
Schema Transforms:Schema 驱动的通用变换能力
Beam 提供了一套直接操作 Schema的变换集合(schema transforms)。其 Java 实现位于 transforms 目录,核心能力包括:
| 能力 | 对应 Transform(源码路径) | 用途说明 |
|---|---|---|
| 字段选择(field selection) | Select.java | 从 Row 中选择一个或多个字段,生成新的精简 Row |
| 分组与聚合 | Group.java、SchemaAggregateFn.java | 按字段分组并对组内聚合(计数、求和、求均值等) |
| 连接操作 | Join.java、CoGroup.java | 基于公共字段对多个 PCollection 做内连接、外连接或共组 |
| 过滤数据 | Filter.java | 按字段条件过滤 Row |
| 添加字段 | AddFields.java | 为 Schema 追加字段 |
| 移除字段 | DropFields.java | 丢弃指定字段 |
| 重命名字段 | RenameFields.java | 修改字段名 |
| 类型转换 | Cast.java | 在 Schema 类型之间做类型转换(如数值拓宽) |
| 结构转换 | Convert.java | 在 Row 与自定义 Java 类型之间互转 |
| 附加键 | WithKeys.java | 将字段提升为 Key,便于后续分组/连接 |
其中SchemaTransform(SchemaTransform.java)是这类变换的统一抽象,SchemaTransformProvider则负责将配置转换为具体的变换实例——这也是一系列"声明式"Schema 变换(如通过 YAML 描述管道)得以落地的底层机制,相关 Provider 见 transforms/providers 目录(JavaFilterTransformProvider、JavaMapToFieldsTransformProvider、JavaExplodeTransformProvider等)。
增强版 ParDo 功能
除了上述通用变换,Schema 还显著增强了ParDo的能力:当DoFn的输入输出带有 Schema 时(相关支持见 DoFnSchemaInformation.java 与 ParDo.java),Beam 可以在DoFn内部通过字段名直接访问元素字段、按需声明只消费部分字段(利于投影优化)、甚至通过注解自动将字段映射到DoFn参数上,减少样板代码。
实战示例:使用 SqlTransform 处理带 Schema 的 PCollection
原文档明确以SqlTransform作为 Schema Transforms 的示例。SqlTransform位于 SqlTransform.java,它把 Beam 管道中的PCollection当作"表",直接用 SQL 做声明式数据处理,底层由 Calcite 查询规划器解析并翻译为 Beam 变换——其前提正是每个输入PCollection都有 Schema。
核心 API(见 SqlTransform.java):
SqlTransform.query(String queryString):以 SQL 字符串构建变换;withTableProvider(name, tableProvider)/withDefaultTableProvider(...):注册自定义表提供者;withQueryPlannerClass(Class<? extends QueryPlanner>):覆盖全局查询规划器(全局可通过BeamSqlPipelineOptions的 planner 选项指定);withNamedParameters(Map<String, ?>)/withPositionalParameters(List<?>):为 SQL 绑定命名 / 位置参数。
一个最小可运行的 Java 示例(Java SDK 核心 API,参考 learning/katas/java 中的常见用法):
import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.extensions.sql.SqlTransform; import org.apache.beam.sdk.schemas.Schema; import org.apache.beam.sdk.transforms.Create; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.Row; // 1. 定义 Schema:有序的命名字段列表 Schema userSchema = Schema.builder() .addInt32Field("id") .addStringField("name") .addInt32Field("age") .build(); Row alice = Row.withSchema(userSchema).addValues(1, "Alice", 30).build(); Row bob = Row.withSchema(userSchema).addValues(2, "Bob", 25).build(); Pipeline pipeline = Pipeline.create(); // 2. Create 会根据传入的 Row 自动携带 Schema PCollection<Row> users = pipeline.apply(Create.of(alice, bob)); // 3. 用 SQL 做声明式过滤与字段选择 PCollection<Row> adults = users.apply( SqlTransform.query("SELECT id, name FROM PCOLLECTION WHERE age >= 26")); pipeline.run();要点说明:
Create.of(...)会根据元素自动推断并附加SchemaCoder;- SQL 中
PCOLLECTION是输入PCollection的固定表名; - 输出
PCollection依然带 Schema(字段为id、name),可以继续衔接其他 Schema Transforms 或写出到带 Schema 的目标(如数据库、Avro、BigQuery 等); - 多张输入表时可通过
withTableProvider/withDefaultTableProvider注册命名表,SQL 中按名称引用。
这一模式说明:只要PCollection带 Schema,Beam 就能把 SQL 的过滤、投影、连接、聚合等操作自动翻译为高效的管道执行,而无需手写DoFn。
结构化数据的使用建议
关于 Apache Beam 中处理结构化数据的最佳实践,仓库的 learning/prompts/documentation-lookup/06_basic_schema.md 引导读者进一步参考 Schema Usage Patterns。结合仓库实现可以总结出几条要点:
- 尽早让数据"带 Schema":优先选择自动附加 Schema 的 I/O(Avro、BigQuery、数据库等),或通过
SchemaRegistry/SchemaProvider为自定义类型注册 Schema,避免在管道中途手工组装Row; - 善用声明式变换替代手写 DoFn:字段选择、过滤、类型转换等操作优先使用
Select、Filter、Cast等 Schema Transforms,代码更简洁且利于 Beam 做投影优化; - 类型转换注意数值拓宽规则:
Cast等操作遵循TypeName.isSupertypeOf定义的隐式兼容链(如INT32→INT64、FLOAT→DOUBLE),超出范围的转换需要显式处理; - 嵌套与自定义类型:复杂的领域结构可用
ROW嵌套或LOGICAL_TYPE建模,例如利用logicaltypes中的EnumerationType、FixedBytes、MicrosInstant等开箱即用的类型; - 关注可空性语义:Schema 的可空性会影响相等性(
EquivalenceNullablePolicy)与编码,声明字段时明确使用addNullableXxxField,避免隐式假设。
总结
Apache Beam Schema 是连接"结构化数据格式(JSON、Protobuf、Avro、数据库行)"与"统一批流处理模型"的桥梁:它以语言无关的有序字段列表描述PCollection元素,在 Java、Python、Go、TypeScript 各 SDK 中原生嵌入,并通过SchemaRegistry、SchemaCoder/RowCoder完成从类型到 Schema 再到编码的完整闭环。在此基础上,字段选择、分组聚合、连接、过滤、字段增删、重命名、类型转换与增强版 ParDo 等 Schema Transforms,以及SqlTransform的声明式 SQL 处理,让开发者可以用极少的样板代码完成绝大多数结构化数据操作。理解 Schema.java 中的类型体系与兼容性规则,是深入使用这些能力的关键起点。
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Schema 完全指南:PCollection 的语言无关类型系统与 Schema Transforms 实战
Apache Beam Schema 完全指南:PCollection 的语言无关类型系统与 Schema Transforms 实战 Apache Beam
大数据批处理流处理数据工程Apache Beam Schema 完全指南:理解 PCollection 的结构化类型系统与 Schema Transform 实战
Apache Beam Schema 完全指南:理解 PCollection 的结构化类型系统与 Schema Transform 实战 Apache Beam
批处理流处理大数据Apache Beam Schema 创建指南:通过 Java POJO、JavaBean 与 AutoValue 构建类型化 PCollection
Apache Beam Schema 创建指南:通过 Java POJO、JavaBean 与 AutoValue 构建类型化 PCollection 导读 本
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考