- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
Apache Avro 是一种面向行存储与数据交换的序列化数据格式,具有紧凑的二进制编码、内嵌 schema 与跨语言互操作等特性。Apache Beam 通过AvroIO模块中的ReadFromAvro/WriteToAvro(Python)、AvroIO.read/AvroIO.write(Java)、avroio.Read/avroio.Write(Go)等转换,为批处理与流处理管道提供对 Avro 文件的读写能力。读完本文,你将掌握 AvroIO 连接器在 Python、Java、Go、TypeScript 各 SDK 中的使用方式、核心参数含义与底层实现原理,并能够直接在 Beam 管道中落地 Avro 数据的读取、写入与 schema 对齐。
一、Apache Beam 对 Avro 格式的支持概览
Apache Beam 将文件 I/O 统一抽象为「读转换(Read transform)+ 写转换(Write transform)」模型,AvroIO 是这一模型在 Avro 格式上的落地实现。Beam 各语言 SDK 均提供了对 Avro 的原生或跨语言支持:
| SDK | 读转换 | 写转换 | 实现位置 |
|---|---|---|---|
| Python | ReadFromAvro/ReadAllFromAvro | WriteToAvro | sdks/python/apache_beam/io/avroio.py |
| Java | AvroIO.read/readAll/readFiles | AvroIO.write/writeGenericRecords | sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java |
| Go | avroio.Read | avroio.Write | sdks/go/pkg/beam/io/avroio/avroio.go |
| TypeScript | readFromAvro | writeToAvro | sdks/typescript/src/apache_beam/io/avroio.ts |
从源码结构看,各 SDK 的 Avro 支持在实现策略上有所分工:Python SDK 默认使用 fastavro 库实现编解码(use_fastavro=True,旧参数仅保留兼容),Java SDK 将 AvroIO 放在独立的extensions/avro模块中,Go SDK 基于fileio框架实现,而 TypeScript SDK 则通过跨语言(cross-language)扩展服务调用schemaio_avro_read:v1/schemaio_avro_write:v1转换实现(见 avroio.ts)。
二、Python SDK:ReadFromAvro 读取 Avro 文件
2.1 最小可运行示例
下面的 Python 管道代码展示了如何从 GCS 通配符路径读取 Avro 文件,并将每条记录打印到日志(示例取自 26_io_avro.md,可复制直接运行):
class ReadAvroOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument( "--path", default="gs://cloud-samples-data/bigquery/us-states/*.avro", help="GCS path to read from", ) options = ReadAvroOptions() with beam.Pipeline(options=options) as p: (p | "Read from Avro" >> ReadFromAvro(options.path) | Map(logging.info))要点说明:
--path使用 Glob 通配符(如*.avro)匹配多个文件,也可指向单个文件;ReadFromAvro的file_pattern参数支持任意 Beam 文件系统(本地、GCS、S3 等)路径;- 默认
validate=True,即管道创建阶段就会校验文件是否存在,若通配符匹配不到任何文件会直接抛错。
2.2 构造参数与行为细节
ReadFromAvro的完整签名(见 avroio.py)为:
ReadFromAvro( file_pattern=None, min_bundle_size=0, validate=True, use_fastavro=True, as_rows=False)各参数含义如下:
| 参数 | 默认值 | 说明 |
|---|---|---|
file_pattern | None | 要读取的文件 Glob 路径,必填 |
min_bundle_size | 0 | 将输入切分为 bundle 时考虑的最小字节数,用于控制并行度与动态分片 |
validate | True | 是否在管道创建阶段校验文件存在 |
use_fastavro | True | 仅为 API 向后兼容保留,已无实际作用,请勿使用 |
as_rows | False | 是否返回带 schema 的 Beam Row PCollection |
记录类型映射规则(源码 docstring 明示):
- 简单类型的记录被映射为对应的 Python 类型;
- Avro
RECORD类型记录被映射为符合 Avro 文件内嵌 schema 的 Python 字典,键为字段名(str),值为 schema 中定义的对应类型; as_rows=True时,读取管道会解析文件中的 writer schema,将其转换为 Beam schema(通过avro_schema_to_beam_schema),并将每条 Avro 记录转换为 Beam Row,从而支持后续使用 SQL、Schema 相关转换等行式 API。
例如,对如下 User schema:
{ "namespace": "example.avro", "type": "record", "name": "User", "fields": [ {"name": "name", "type": "string"}, {"name": "favorite_number", "type": ["int", "null"]}, {"name": "favorite_color", "type": ["string", "null"]} ] }读取得到的每条记录为如下形式的字典:{'name': 'Alyssa', 'favorite_number': 256, 'favorite_color': None}。
2.3 ReadAllFromAvro:动态文件列表读取
当需要先对文件路径本身进行计算(如来自另一个 PCollection)时,使用ReadAllFromAvro:
ReadAllFromAvro(min_bundle_size=0, desired_bundle_size=64 * 1024 * 1024, ...)它接收一个PCollection的 Avro 文件路径或文件模式,输出记录的PCollection。其默认期望 bundle 大小为 64MB(DEFAULT_DESIRED_BUNDLE_SIZE = 64 * 1024 * 1024)。源码注释指出该实现目前仅针对批处理管道做过完整测试,流式场景下因涉及 ReShuffle 可能存在读取延迟(见 avroio.py)。
三、Python SDK:WriteToAvro 写入 Avro 文件
3.1 构造参数
WriteToAvro的签名(见 avroio.py)为:
WriteToAvro( file_path_prefix, schema=None, codec='deflate', file_name_suffix='', num_shards=0, shard_name_template=None, mime_type='application/x-avro', use_fastavro=True)| 参数 | 默认值 | 说明 |
|---|---|---|
file_path_prefix | — | 输出文件路径前缀,生成的文件名为「前缀 + 分片标识 + 后缀」 |
schema | None | 要使用的 Avro schema(dict 形式)。不指定时,若输入 PCollection 带 Beam schema,会自动转换生成 Avro schema |
codec | 'deflate' | 块级(block-level)压缩编码,Avro 规范支持的任何字符串均可,例如'null' |
file_name_suffix | '' | 输出文件的后缀,如'.avro' |
num_shards | 0 | 输出分片数,0 表示由执行服务自动决定最优分片数。手动约束分片数通常会降低管道性能,除非必须指定输出文件数量,否则不建议设置 |
shard_name_template | None | 分片文件名模板,大写S/N分别被替换为 0 填充的分片号与分片总数;默认模板为-SSSSS-of-NNNNN;传''等价于num_shards=1,只生成一个文件 |
mime_type | 'application/x-avro' | 生成文件的 MIME 类型(在文件系统支持时生效) |
use_fastavro | True | 仅向后兼容保留,无实际作用 |
3.2 schema 自动生成规则
WriteToAvro.expand中的实现逻辑(见 avroio.py)值得注意:
- 若显式传入
schema,则直接使用; - 否则尝试通过
schemas.schema_from_element_type从输入 PCollection 的元素类型推断 Beam schema,再经beam_schema_to_avro_schema转换为 Avro schema; - 若输入既无 schema 也无显式 Avro schema,会抛出
ValueError,提示「写入非 schema 化 PCollection 必须显式指定 schema」。
另外,写入时压缩发生在块级(由codec控制),因此底层 sink 将compression_type设为CompressionTypes.UNCOMPRESSED(见 avroio.py)。
3.3 实用写示例
# 方式一:输入为带 schema 的 PCollection,自动生成 Avro schema (p | beam.Create([{"name": "Alyssa", "favorite_number": 256}]) | WriteToAvro(file_path_prefix="gs://my-bucket/output/users", file_name_suffix=".avro")) # 方式二:显式传入 schema(dict),输出固定 schema 的 Avro 文件 schema = { "type": "record", "name": "User", "fields": [ {"name": "name", "type": "string"}, {"name": "favorite_number", "type": ["int", "null"]}, ], } (p | beam.Create([{"name": "Alyssa", "favorite_number": 256}]) | WriteToAvro(file_path_prefix="users", schema=schema, codec='snappy', file_name_suffix='.avro'))四、Java SDK:AvroIO.read / AvroIO.write
Java 的 AvroIO 位于独立模块sdks/java/extensions/avro,类声明为org.apache.beam.sdk.extensions.avro.io.AvroIO。核心 API 由嵌套的Read、ReadFiles、ReadAll与Write、TypedWrite类组成(见 AvroIO.java)。
4.1 读取
基本读法:
PCollection<AvroAutoGenClass> records = p.apply(AvroIO .read(AvroAutoGenClass.class) .from("gs://my_bucket/path/to/records-*.avro"));支持的读取形态(源码 javadoc 汇总):
AvroIO.read(Class):按 Avro 生成的特定类(specific record)读取;AvroIO.readGenericRecords(Schema):按给定 AvroSchema读取为GenericRecord;readAll(Class)/readAllGenericRecords(schema):从「文件模式 PCollection」读取;FileIO.matchAll() + FileIO.readMatches() + AvroIO.readFiles(Class):先匹配再读取的解耦式用法;AvroIO.parseFilesGenericRecords(...):配合自定义解析函数读取。
流式场景下可使用watchForNewFiles持续监听新文件:
PCollection<AvroAutoGenClass> records = p.apply(AvroIO .read(AvroAutoGenClass.class) .from("gs://my_bucket/path/to/records-*.avro") .watchForNewFiles( Duration.standardMinutes(1), // 每分钟检查一次新文件 afterTimeSinceNewOutput(Duration.standardHours(1)))); // 1 小时无新文件则停止监听文件数量极大(数万级)时,可用withHintMatchesManyFiles()提升性能与可扩展性;但若匹配文件很少,该提示反而可能降低性能。
4.2 Beam schema 推断与 AvroCoder
若需要对 Avro 结果的 PCollection 使用 SQL 或基于 schema 的转换,需开启 schema 推断:
PCollection<AvroAutoGenClass> records = p.apply(AvroIO.read(...).from(...).withBeamSchemas(true));对于经由其他途径(如 Kafka)得到的 Avro 对象 PCollection,可通过向pipeline.getSchemaRegistry()注册 schema provider 使其 schema 化:
pipeline.getSchemaRegistry().registerSchemaProvider( AvroAutoGenClass.class, AvroAutoGenClass.getClassSchema());或手动为 PCollection 设置 Avro 支持的 schema coder:
records.setCoder(AvroUtils.schemaCoder(recordClass, schema)); // GenericRecord 场景: records.setCoder(AvroUtils.schemaCoder(avroSchema));4.3 写入
// 特定类写入(本地执行) PCollection<AvroAutoGenClass> records = ...; records.apply(AvroIO.write(AvroAutoGenClass.class).to("/path/to/file.avro")); // GenericRecord 写入(分片 GCS 文件,可远程执行) Schema schema = new Schema.Parser().parse(new File("schema.avsc")); PCollection<GenericRecord> records = ...; records.apply("WriteToAvro", AvroIO.writeGenericRecords(schema) .to("gs://my_bucket/path/to/numbers") .withSuffix(".avro"));关键行为(源码 javadoc 明示):
- 输出文件名由
DefaultFilenamePolicy依据前缀、ShardNameTemplate(withShardNameTemplate)与后缀(withSuffix)分片生成,也可用Write.to(FilenamePolicy)自定义命名策略; - 默认使用
CodecFactory.snappyCodec()压缩,可通过withCodec更改; - 默认所有输入落入全局窗口后再写入;流式 runner 需要按窗口写时用
withWindowedWrites()保留窗口与触发语义,此时必须显式设置withNumShards(int),且自定义FilenamePolicy须保证不同窗口/触发产生唯一文件名; TypedWrite.withSchema(Schema)用于动态目的地(DynamicDestinations)场景下按文件指定 schema,若使用DynamicDestinations而忽略withSchema()会抛出异常(见 AvroIO.java)。
五、Go SDK:avroio.Read / avroio.Write
Go 的 Avro 支持位于 sdks/go/pkg/beam/io/avroio/avroio.go,基于fileio框架实现。
import "github.com/apache/beam/sdks/v2/go/pkg/beam/io/avroio" // 读取:glob 为文件模式,t 为目标 Go 类型(对应 Avro schema) result := avroio.Read(s, "gs://my-bucket/*.avro", reflect.TypeOf(MyStruct{})) // 写入:filename 为输出前缀,schema 为 Avro schema JSON 字符串 avroio.Write(s, "gs://my-bucket/out.avro", schemaJSON, col)实现细节:Read内部先做fileio.MatchAll(允许空匹配),再fileio.ReadMatches(按未压缩方式读取),最后通过avroReadFnDoFn 逐文件打开并解析记录(见 avroio.go)。Write则先AddFixedKey+GroupByKey归并,再由writeAvroFn打开文件系统写入流写出(见 avroio.go)。
六、TypeScript SDK:基于跨语言转换的 Avro 支持
TypeScript SDK 本身不实现 Avro 编解码,而是通过跨语言(cross-language)调用 Java 扩展服务。从 avroio.ts 源码可见:
readFromAvro(filePattern, options)内部通过schemaio调用beam:transform:org.apache.beam:schemaio_avro_read:v1,并将{ location: filePattern, schema }作为参数下发;writeToAvro(filePath, options)调用beam:transform:org.apache.beam:schemaio_avro_write:v1;- 当前要求显式传入
options.schema(源码标注 TODO:允许自动推断 schema)。
因此,TypeScript 使用者需要先具备对应的 Avro schema,再由远程扩展服务完成实际的 Avro 序列化/反序列化。
七、测试与验证:如何确认 Avro 读写正确性
仓库提供了完备的测试覆盖,可作为自建管道的行为参照:
- Python 侧,avroio_test.py 覆盖了带/不带分片读取(
test_read_without_splitting、test_read_with_splitting)、多 block 读取、deflate/snappy 压缩读取、通配符模式读取以及test_read_from_avro等用例; - Go 侧,avroio_test.go 提供读写的端到端测试。
上述测试表明 AvroIO 对分片(splitting)、块级压缩(deflate、snappy)与 Glob 模式均有稳定支持,这正是 Avro 文件可被 Beam 并行、分布式处理的基础。
八、最佳实践小结
- 读取优先用通配符:
ReadFromAvro(".../*.avro")即可并行读取多个文件;文件列表需要动态计算时改用ReadAllFromAvro(Python)或readAll(Java)。 - 写 schema 优先显式指定:Python 中非 schema 化 PCollection 必须显式传
schema,否则抛出ValueError;传 schema 时注意使用 dict 而非 avro-python3 解析出的 schema 对象(fastavro 模式下会报错)。 - 压缩按需选择:Python 默认
deflate,Java 默认snappy;块级压缩不影响 Beam 文件级分片读取,可依据吞吐与压缩率权衡。 - 分片数交给服务决定:
num_shards=0(Python)/ 不设withNumShards(Java)时由 runner 自适应,手动固定分片会牺牲性能。 - 需要 SQL/行式 API 时开启 schema 化:Python 用
as_rows=True,Java 用withBeamSchemas(true)或注册 schema provider / 设置AvroUtils.schemaCoder。 - 流式写窗口数据时:Java 使用
withWindowedWrites()并显式设置分片数、自定义FilenamePolicy,保证每个窗口的文件名唯一。
参考路径速查
- Python 实现:sdks/python/apache_beam/io/avroio.py(
ReadFromAvro见 L78、ReadAllFromAvro见 L186、WriteToAvro见 L365) - Python 测试:sdks/python/apache_beam/io/avroio_test.py
- Java 实现:sdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
- Go 实现与测试:sdks/go/pkg/beam/io/avroio/avroio.go、sdks/go/pkg/beam/io/avroio/avroio_test.go
- TypeScript 实现:sdks/typescript/src/apache_beam/io/avroio.ts
- 批处理
- 流处理
- 大数据
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Midway Web 路由表(Router Table)完全指南:从路由服务注入到动态注册与优先级排序
Midway Web 路由表(Router Table)完全指南:从路由服务注入到动态注册与优先级排序 从 v2.8.0 开始,Midway 提供了内置的路由表
后端微服务云原生使用 Apache Beam 的 AvroIO 读取 Avro 文件:GenericRecord 实战详解
使用 Apache Beam 的 AvroIO 读取 Avro 文件:GenericRecord 实战详解 本指南围绕 Apache Beam Java SDK
大数据批处理流处理数据工程Apache Beam Python 中使用 AvroIO 读取与写入 Avro 文件:ReadFromAvro 实战指南
Apache Beam Python 中使用 AvroIO 读取与写入 Avro 文件:ReadFromAvro 实战指南 导读 Apache Avro 是一种
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考