Apache Beam写入BigQuery遇格式错误:Dataflow动态分表写入失败
问题描述
在Dataflow中构建了一个数据流管道,目标是根据数据中的event_name和event_date将流数据拆分到动态命名的BigQuery表中。目前表已按正确名称创建,但数据写入BigQuery时失败,报错如下:
"Unknown name "json" at 'rows[0]': Proto field is not repeating, cannot start list"
调用WriteToBigQuery前的打印日志显示记录格式看似正常:
About to write to BigQuery - Table: PROJECT_ID:DATASET_NAME.TABLE_NAME, Record: [{'event_name': 'scroll', 'event_date': '20241118', 'user_id': '', 'platform': 'WEB'}]
(已尝试移除方括号,但结果相同)
管道代码如下:
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions from apache_beam.io.gcp.bigquery import WriteToBigQuery from apache_beam.transforms.window import FixedWindows import logging def log_before_write(element): table_name, record = element logging.info(f"About to write to BigQuery - Table: {table_name}, Record: {record}") return element class SplitByParameter(beam.DoFn): def process(self, element): event_name = element['event_name'] event_date = element['event_date'] yield (event_name, event_date, element) def format_table_name(element): event_name, event_date, record = element sanitized_event_name = event_name.replace(' ', '_') sanitized_event_date = event_date.replace(' ', '_') table_name = f'PROJECT_ID:DATASET.{sanitized_event_name}_{sanitized_event_date}' return table_name, record def split_records(element): table_name, record = element json_record = [{ 'event_name': str(record.get('event_name', '')) if record.get('event_name') is not None else '', 'event_date': str(record.get('event_date', '')) if record.get('event_date') is not None else '', 'user_id': str(record.get('user_id', '')) if record.get('user_id') is not None else '', 'platform': str(record.get('platform', '')) if record.get('platform') is not None else '' }] yield (table_name,json_record) def print_record(record): logging.info(f"Record before WriteToBigQuery: {record}") return record def run(argv=None): options = PipelineOptions(argv) options.view_as(StandardOptions).streaming = True p = beam.Pipeline(options=options) # Define schema for BigQuery (this needs to match your record structure) schema = 'event_name:STRING, event_date:STRING, user_id:STRING, platform:STRING' # Read from BigQuery, apply windowing, and process records (p | 'ReadFromBigQuery' >> beam.io.ReadFromBigQuery(query=f''' SELECT * FROM `PROJECT_ID.DATASET.TABLE` WHERE _TABLE_SUFFIX = FORMAT_TIMESTAMP('%Y%m%d', CURRENT_TIMESTAMP()) ''', use_standard_sql=True) | 'ApplyWindowing' >> beam.WindowInto(FixedWindows(60)) # 60-second window | 'SplitByParameter' >> beam.ParDo(SplitByParameter()) # Split by event_name and event_date | 'FormatTableName' >> beam.Map(format_table_name) # Format the table name | 'LogBeforeFlatMap' >> beam.Map(lambda x: logging.info(f'Before FlatMap: {x}') or x) | 'SplitRecords' >> beam.FlatMap(split_records) # Convert record to desired format | 'LogBeforeWrite' >> beam.Map(log_before_write) | 'PrintRecord' >> beam.Map(print_record) # Print records before writing to BigQuery | 'WriteToBigQuery' >> beam.io.WriteToBigQuery( table=lambda x: x[0], # Table name is the first element of the tuple schema=schema, # Use the schema defined above write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND # Append data to existing tables ) ) p.run() if __name__ == '__main__': logging.getLogger().setLevel(logging.INFO) run()
解决方案
问题根源
报错的核心原因是传递给WriteToBigQuery的记录结构错误:当前代码中split_records函数返回的是表名 + 列表格式的记录,但WriteToBigQuery期望的是表名 + 单个字典格式的记录(动态指定表时,需用元组(表名, 单个记录)的格式)。
修复步骤
- 修改
split_records函数,去掉记录外层的列表包裹,直接返回单个字典 - 保留必要的空值处理逻辑,避免数据格式不匹配
修改后的关键代码
def split_records(element): table_name, record = element cleaned_record = { 'event_name': str(record.get('event_name', '')) if record.get('event_name') is not None else '', 'event_date': str(record.get('event_date', '')) if record.get('event_date') is not None else '', 'user_id': str(record.get('user_id', '')) if record.get('user_id') is not None else '', 'platform': str(record.get('platform', '')) if record.get('platform') is not None else '' } # 返回(表名, 单个字典记录) yield (table_name, cleaned_record)
完整优化后的管道代码
import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions from apache_beam.io.gcp.bigquery import WriteToBigQuery from apache_beam.transforms.window import FixedWindows import logging def log_before_write(element): table_name, record = element logging.info(f"About to write to BigQuery - Table: {table_name}, Record: {record}") return element class SplitByParameter(beam.DoFn): def process(self, element): event_name = element['event_name'] event_date = element['event_date'] yield (event_name, event_date, element) def format_table_name(element): event_name, event_date, record = element sanitized_event_name = event_name.replace(' ', '_') sanitized_event_date = event_date.replace(' ', '_') table_name = f'PROJECT_ID:DATASET.{sanitized_event_name}_{sanitized_event_date}' return table_name, record def clean_record(element): table_name, record = element cleaned_record = { 'event_name': str(record.get('event_name', '')) if record.get('event_name') is not None else '', 'event_date': str(record.get('event_date', '')) if record.get('event_date') is not None else '', 'user_id': str(record.get('user_id', '')) if record.get('user_id') is not None else '', 'platform': str(record.get('platform', '')) if record.get('platform') is not None else '' } return (table_name, cleaned_record) def run(argv=None): options = PipelineOptions(argv) options.view_as(StandardOptions).streaming = True p = beam.Pipeline(options=options) schema = 'event_name:STRING, event_date:STRING, user_id:STRING, platform:STRING' (p | 'ReadFromBigQuery' >> beam.io.ReadFromBigQuery(query=f''' SELECT * FROM `PROJECT_ID.DATASET.TABLE` WHERE _TABLE_SUFFIX = FORMAT_TIMESTAMP('%Y%m%d', CURRENT_TIMESTAMP()) ''', use_standard_sql=True) | 'ApplyWindowing' >> beam.WindowInto(FixedWindows(60)) | 'SplitByParameter' >> beam.ParDo(SplitByParameter()) | 'FormatTableName' >> beam.Map(format_table_name) | 'CleanRecord' >> beam.Map(clean_record) | 'LogBeforeWrite' >> beam.Map(log_before_write) | 'WriteToBigQuery' >> beam.io.WriteToBigQuery( table=lambda x: x[0], schema=schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND ) ) p.run() if __name__ == '__main__': logging.getLogger().setLevel(logging.INFO) run()
内容的提问来源于stack exchange,提问作者Rick Rickles
相关产品推荐
相关产品推荐

