news 2026/10/10 1:45:07

Apache Beam 中的 Avro 文件读写:AvroIO 连接器全解析

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam 中的 Avro 文件读写:AvroIO 连接器全解析
  • 批处理
  • 流处理
  • 大数据

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam15/beam
点击查看免费下载

导读

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读转换写转换实现位置
PythonReadFromAvro/ReadAllFromAvroWriteToAvrosdks/python/apache_beam/io/avroio.py
JavaAvroIO.read/readAll/readFilesAvroIO.write/writeGenericRecordssdks/java/extensions/avro/src/main/java/org/apache/beam/sdk/extensions/avro/io/AvroIO.java
Goavroio.Readavroio.Writesdks/go/pkg/beam/io/avroio/avroio.go
TypeScriptreadFromAvrowriteToAvrosdks/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_patternNone要读取的文件 Glob 路径,必填
min_bundle_size0将输入切分为 bundle 时考虑的最小字节数,用于控制并行度与动态分片
validateTrue是否在管道创建阶段校验文件存在
use_fastavroTrue仅为 API 向后兼容保留,已无实际作用,请勿使用
as_rowsFalse是否返回带 schema 的 Beam Row PCollection

记录类型映射规则(源码 docstring 明示):

  • 简单类型的记录被映射为对应的 Python 类型;
  • AvroRECORD类型记录被映射为符合 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—输出文件路径前缀,生成的文件名为「前缀 + 分片标识 + 后缀」
schemaNone要使用的 Avro schema(dict 形式)。不指定时,若输入 PCollection 带 Beam schema,会自动转换生成 Avro schema
codec'deflate'块级(block-level)压缩编码,Avro 规范支持的任何字符串均可,例如'null'
file_name_suffix''输出文件的后缀,如'.avro'
num_shards0输出分片数,0 表示由执行服务自动决定最优分片数。手动约束分片数通常会降低管道性能,除非必须指定输出文件数量,否则不建议设置
shard_name_templateNone分片文件名模板,大写S/N分别被替换为 0 填充的分片号与分片总数;默认模板为-SSSSS-of-NNNNN;传''等价于num_shards=1,只生成一个文件
mime_type'application/x-avro'生成文件的 MIME 类型(在文件系统支持时生效)
use_fastavroTrue仅向后兼容保留,无实际作用

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 并行、分布式处理的基础。

八、最佳实践小结

  1. 读取优先用通配符:ReadFromAvro(".../*.avro")即可并行读取多个文件;文件列表需要动态计算时改用ReadAllFromAvro(Python)或readAll(Java)。
  2. 写 schema 优先显式指定:Python 中非 schema 化 PCollection 必须显式传schema,否则抛出ValueError;传 schema 时注意使用 dict 而非 avro-python3 解析出的 schema 对象(fastavro 模式下会报错)。
  3. 压缩按需选择:Python 默认deflate,Java 默认snappy;块级压缩不影响 Beam 文件级分片读取,可依据吞吐与压缩率权衡。
  4. 分片数交给服务决定:num_shards=0(Python)/ 不设withNumShards(Java)时由 runner 自适应,手动固定分片会牺牲性能。
  5. 需要 SQL/行式 API 时开启 schema 化:Python 用as_rows=True,Java 用withBeamSchemas(true)或注册 schema provider / 设置AvroUtils.schemaCoder。
  6. 流式写窗口数据时: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.

项目地址:https://gitcode.com/gh_mirrors/beam15/beam
点击查看免费下载
上一篇:Qwen3-VL-30B-A3B-Thinking-FP8模型发布:FP8量化技术赋能多模态大模型高效部署
下一篇:react-native-image-picker无障碍测试:VoiceOver与TalkBack兼容性

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/10 1:44:51

作物害虫识别实战:从数据集预处理到迁移学习模型训练全流程

简介&#xff1a;面向计算机、人工智能、数据科学及相关专业的同学和从业者&#xff0c;这套基于机器学习的作物害虫识别与分类项目包&#xff0c;覆盖从数据加载、模型训练到分类结果输出的完整流程&#xff0c;既可用来练手入门&#xff0c;也可作为大作业、课程设计或毕业设…

作者头像 李华
网站建设 2026/10/10 1:44:40

波士顿房价预测实战:线性回归从数据预处理到模型评估全流程

简介&#xff1a;一份以波士顿房价预测为主线、线性回归从原理到实战的代码合集&#xff0c;适合机器学习初学者与需要快速搭建回归预测流程的开发者。资源共20个文件&#xff0c;含11个Python脚本与9个CSV数据文件&#xff1a;脚本承担数据加载、特征工程、模型训练、评估与可…

作者头像 李华