使用追加模式写入Avro文件时如何修改其Schema?
如何修改Avro文件Schema并追加记录(基于fastavro)
问题场景
你尝试用fastavro写入初始Avro文件后,想要修改Schema(让val字段支持null)并追加包含None值的记录,但直接修改Schema后追加时报错:
ValueError: Provided schema {'type': 'record', 'name': 'test', 'fields': [{'name': 'id', 'type': 'int'}, {'name': 'val', 'type': ['long', 'null']}], '__fastavro_parsed': True, '__named_schemas': {'test': {'type': 'record', 'name': 'test', 'fields': [{'name': 'id', 'type': 'int'}, {'name': 'val', 'type': ['long', 'null']}]}}} does not match file writer_schema {'type': 'record', 'name': 'test', 'fields': [{'name': 'id', 'type': 'int'}, {'name': 'val', 'type': 'long'}], '__fastavro_parsed': True, '__named_schemas': {'test': {'type': 'record', 'name': 'test', 'fields': [{'name': 'id', 'type': 'int'}, {'name': 'val', 'type': 'long'}]}}}
初始写入代码:
from fastavro import writer, parse_schema schema = { 'name': 'test', 'type': 'record', 'fields': [ {'name': 'id', 'type': 'int'}, {'name': 'val', 'type': 'long'}, ], } records = [ {'id': 1, 'val': 0.2}, {'id': 2, 'val': 3.1}, ] with open('test.avro', 'wb') as f: writer(f, parse_schema(schema), records)
尝试修改Schema并追加的代码:
more_records = [ {'id': 3, 'val': 1.5}, {'id': 2, 'val': None}, ] schema['fields'][1]['type'] = ['long', 'null'] with open('test.avro', 'a+b') as f: writer(f, parse_schema(schema), more_records)
原因
Avro文件的Schema是不可变的,写入时会被嵌入文件头部。fastavro的writer在追加模式下会校验传入的Schema与文件原始Schema是否完全一致,不允许直接修改原始Schema后追加。
可行解决方案
由于无法直接修改已存在的Avro文件Schema,只能通过重写整个文件的方式实现需求,步骤如下:
- 读取原始文件的所有记录和原始Schema
- 定义兼容的新Schema(将
val字段类型改为可空的联合类型) - 合并原始记录与新记录
- 使用新Schema重新写入文件
具体代码实现
from fastavro import writer, reader, parse_schema # 1. 读取原始Avro文件的记录和Schema original_records = [] with open('test.avro', 'rb') as f: avro_reader = reader(f) original_schema = avro_reader.writer_schema for record in avro_reader: original_records.append(record) # 2. 定义新的兼容Schema new_schema = { 'name': 'test', 'type': 'record', 'fields': [ {'name': 'id', 'type': 'int'}, {'name': 'val', 'type': ['long', 'null']}, # 修改为支持null的联合类型 ], } parsed_new_schema = parse_schema(new_schema) # 3. 合并原始记录和新记录 more_records = [ {'id': 3, 'val': 1.5}, {'id': 2, 'val': None}, ] all_records = original_records + more_records # 4. 用新Schema重写整个文件 with open('test_updated.avro', 'wb') as f: writer(f, parsed_new_schema, all_records)
注意事项
- 新Schema必须与原始Schema向前兼容,否则读取原始记录时会出错。这里将字段类型从
long改为['long', 'null']是符合Avro兼容规则的(允许字段新增可空选项)。 - 如果原始文件很大,一次性读取所有记录可能会占用较多内存,此时可以考虑分批次处理,但最终仍需重写整个文件。
内容的提问来源于stack exchange,提问作者egriffiths
相关产品推荐
相关产品推荐

