使用Apache Beam写入嵌套JSON到BigQuery报错,求解决方法
问题分析与解决:Apache Beam加载嵌套JSON到BigQuery报错
问题背景
有一个包含多层嵌套JSON数据的文件,计划使用Apache Beam将其加载至BigQuery,但运行代码时出现字段不匹配的报错。
1. JSON数据文件内容
{"id":1,"name":"A","status":"ACTIVE","dataProvider":{"name":"Exelate 3PD"},"endDateTime":{"date":{"year":2038,"month":1,"day":19},"hour":14,"minute":14,"second":7,"timeZoneId":"Australia/Sydney"}} {"id":2,"name":"B","status":"ACTIVE","dataProvider":{"name":"Exelate 3PD"},"endDateTime":{"date":{"year":2038,"month":1,"day":19},"hour":14,"minute":14,"second":7,"timeZoneId":"Australia/Sydney"}} {"id":3,"name":"C","status":"ACTIVE","dataProvider":{"name":"Exelate 3PD"},"endDateTime":{"date":{"year":2038,"day":19},"hour":14,"minute":14,"second":7}}
2. BigQuery表结构
{ "fields": [ { "mode": "NULLABLE", "name": "id", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "status", "type": "STRING" }, { "fields": [ { "mode": "NULLABLE", "name": "name", "type": "STRING" } ], "mode": "NULLABLE", "name": "dataProvider", "type": "RECORD" }, { "fields": [ { "fields": [ { "mode": "NULLABLE", "name": "year", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "month", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "day", "type": "INTEGER" } ], "mode": "NULLABLE", "name": "date", "type": "RECORD" }, { "mode": "NULLABLE", "name": "hour", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "minute", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "second", "type": "INTEGER" }, { "mode": "NULLABLE", "name": "timeZoneId", "type": "STRING" } ], "mode": "NULLABLE", "name": "endDateTime", "type": "RECORD" } ] }
3. 加载代码
from apache_beam.io.gcp.bigquery_tools import parse_table_schema_from_json import json import apache_beam as beam import re schema_data = json.dumps(json.load(open("schema.json"))) table_schema = parse_table_schema_from_json(schema_data) def parse_json(element): row = json.loads(element) return row inputs_pattern = 'data/orderrecords.txt' with beam.Pipeline() as pipeline: out= ( pipeline | 'Take in Dataset' >> beam.io.ReadFromText(inputs_pattern) | beam.Map(parse_json) | beam.io.WriteToBigQuery( 'apt-ent-45:test.order' , schema=table_schema, # write_disposition=beam.io.BigQueryDisposition.WRITE_TRUNCATE, method="STREAMING_INSERTS") )
4. 报错信息
WARNING:apache_beam.io.gcp.bigquery:There were errors inserting to BigQuery. Will retry. Errors were [{'index': 0, 'errors': [{'reason': 'invalid', 'location': 'name', 'debugInfo': '', 'message': 'no such field: name.'}]}, {'index': 1, 'errors': [{'reason': 'invalid', 'location': 'name', 'debugInfo': '', 'message': 'no such field: name.'}]}, {'index': 2, 'errors': [{'reason': 'invalid', 'location': 'name', 'debugInfo': '', 'message': 'no such field: name.'}]}]
问题原因
你的BigQuery表结构里没有定义顶层的name字段,但JSON数据里每条记录都包含name字段(比如"name":"A")。当Beam将解析后的JSON数据写入BigQuery时,BigQuery发现传入的name字段不在表的schema定义中,因此抛出该错误。
另外,第三条JSON记录的endDateTime.date缺少month字段,后续写入时也可能引发字段缺失的报错,需要一并处理。
解决方法
方法一:修改BigQuery表结构,添加name字段
在表schema的fields数组中添加name字段的定义,确保数据和表结构匹配:
{ "mode": "NULLABLE", "name": "name", "type": "STRING" }
可以将其放在id字段之后、status字段之前的位置。
方法二:在Beam解析阶段过滤掉name字段
如果不需要将name字段写入BigQuery,可以修改parse_json函数,移除JSON中的name字段,同时处理month字段缺失的问题:
def parse_json(element): row = json.loads(element) # 移除顶层的name字段 row.pop('name', None) # 处理endDateTime.date中缺失的month字段,设置为NULL date_obj = row.get('endDateTime', {}).get('date', {}) if 'month' not in date_obj: date_obj['month'] = None return row
其他加载嵌套JSON到BigQuery的方法
1. 自动推断表结构
如果是测试环境,可以让Beam自动推断表结构,无需手动维护schema。修改WriteToBigQuery的参数,移除schema=table_schema,添加创建表的配置:
| beam.io.WriteToBigQuery( 'apt-ent-45:test.order', create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, method="STREAMING_INSERTS" )
注意:生产环境不建议使用该方法,自动推断的字段类型可能不符合预期。
2. 显式构建Beam Row对象
对于复杂的嵌套结构,可以显式构建Row对象,精准匹配表结构,同时更灵活处理缺失字段:
from apache_beam import Row def parse_json_to_row(element): data = json.loads(element) # 构建dataProvider子对象 data_provider = Row(name=data['dataProvider']['name']) # 处理date字段,补全缺失的month date_data = data['endDateTime']['date'] date_row = Row( year=date_data.get('year'), month=date_data.get('month'), day=date_data.get('day') ) # 构建endDateTime子对象 end_date_time = Row( date=date_row, hour=data['endDateTime']['hour'], minute=data['endDateTime']['minute'], second=data['endDateTime']['second'], timeZoneId=data['endDateTime'].get('timeZoneId') ) # 构建主Row,过滤不需要的name字段 return Row( id=data['id'], status=data['status'], dataProvider=data_provider, endDateTime=end_date_time )
在Pipeline中用该函数替代原有的parse_json即可。
内容的提问来源于stack exchange,提问作者Amarjeet Kushwaha
相关产品推荐
相关产品推荐

