You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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期望的是表名 + 单个字典格式的记录(动态指定表时,需用元组(表名, 单个记录)的格式)。

修复步骤

  1. 修改split_records函数,去掉记录外层的列表包裹,直接返回单个字典
  2. 保留必要的空值处理逻辑,避免数据格式不匹配

修改后的关键代码

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.16 03:37:02