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

如何适配Beam ReadFromBigQuery输出字典至beam_nuggets写入数据库

问题:Dataflow从BigQuery读取数据写入PostgreSQL失败(beam_nuggets)

问题背景

在Python 3.9中通过FlexTemplate运行Dataflow管道,从BigQuery读取非嵌套/非重复记录,尝试使用beam_nuggets写入PostgreSQL测试数据库。BigQuery输出符合预期的Python字典,但写入数据库时管道失败。

BigQuery输出的测试数据(写入GCS内容)

{'order_id': 'CM-2011-110', 'order_date': '4/10/2011', 'ship_date': '10/10/2011', 'ship_mode': 'Standard Class', 'customer_name': 'Alejandro Grove', 'segment': 'Consumer', 'state': 'Est', 'country': 'Cameroon', 'market': 'Africa', 'region': 'Africa', 'product_id': 'OFF-CAR-10002031', 'category': 'Office Supplies', 'sub_category': 'Binders', 'product_name': 'Cardinal 3-Hole Punch, Durable', 'sales': 30, 'quantity': 1, 'discount': 0.0, 'profit': 13.92, 'shipping_cost': 2.57, 'order_priority': 'Medium', 'year': 2011}
{'order_id': 'CM-2011-110', 'order_date': '4/10/2011', 'ship_date': '10/10/2011', 'ship_mode': 'Standard Class', 'customer_name': 'Alejandro Grove', 'segment': 'Consumer', 'state': 'Est', 'country': 'Cameroon', 'market': 'Africa', 'region': 'Africa', 'product_id': 'TEC-CAN-10002879', 'category': 'Technology', 'sub_category': 'Copiers', 'product_name': 'Canon Copy Machine, High-Speed', 'sales': 521, 'quantity': 2, 'discount': 0.0, 'profit': 93.72, 'shipping_cost': 30.83, 'order_priority': 'Medium', 'year': 2011}

原代码

def run(save_main_session=True):
    beam_options = PipelineOptions()
    args = beam_options.view_as(MyOptions)
    with beam.Pipeline(options=beam_options) as p:
        db_source = db_reader(args.bqprojectid, args.dataset, args.table, args.limit)
        db = db_writer(args.destinationip, args.port, args.destinationusername, args.destinationpassword, args.destinationtable, args.database_name)
        result = (
            p | beam.io.ReadFromBigQuery(use_standard_sql=True, query=db_source.sql_query()
                )
               # 调试用步骤,后续计划移除
              |'Write to GCS' >> WriteToText('SOMEURI')
              |'Write to Database' >> relational_db.Write(
                  source_config = (db.sink_config()),
                  table_config = (db.table_config())
              ))

错误信息

File "/usr/local/lib/python3.9/site-packages/beam_nuggets/io/relational_db.py", line 181, in process
    assert isinstance(element, dict)

验证与尝试

  • 确认数据库凭证有效:使用beam_nuggets文档中的静态PCollection可以成功上传数据:
    months = p | "Reading month records" >> beam.Create([
            {'name': 'Jan', 'num': 1},
            {'name': 'Feb', 'num': 2},
        ])
    
  • 曾考虑用正则转换BigQuery输出,但认为不是最佳实践,希望找到更简便的方式将输出转换成静态示例中的字典结构。

解决方案

通过添加映射函数,将BigQuery输出转换为符合beam_nuggets要求的字典结构,修改后的完整代码如下:

映射转换函数

def map_to_beam_nuggets_data(element):
    return {
        'order_id': element['order_id'],
        'order_date': element['order_date'],
        'ship_date': element['ship_date'],
        'ship_mode': element['ship_mode'],
        'customer_name': element['customer_name'],
        'segment': element['segment'],
        'state': element['state'],
        'country': element['country'],
        'market': element['market'],
        'region': element['region'],
        'product_id': element['product_id'],
        'category': element['category'],
        'sub_category': element['sub_category'],
        'product_name': element['product_name'],
        'sales': element['sales'],
        'quantity': element['quantity'],
        'discount': element['discount'],
        'profit': element['profit'],
        'shipping_cost': element['shipping_cost'],
        'order_priority': element['order_priority'],
    }

修改后的主运行函数

# 控制器函数
def run(save_main_session=True):
    beam_options = PipelineOptions()
    args = beam_options.view_as(MyOptions)
    with beam.Pipeline(options=beam_options) as p:
        db_source = db_reader(args.bqprojectid, args.dataset, args.table, args.limit)
        db = db_writer(args.destinationip, args.port, args.destinationusername, args.destinationpassword, args.destinationtable, args.database_name)
              # 从BigQuery读取结果
        result = (p | beam.io.ReadFromBigQuery(use_standard_sql=True, query=db_source.sql_query()))

        # 写入GCS用于调试
        result | 'Write to GCS' >> WriteToText('SOMEURI')

        # 转换数据结构后写入数据库
        (result
         | 'Map to beam nuggets data' >> beam.Map(map_to_beam_nuggets_data)
         | 'Write to Database' >> relational_db.Write(
                    source_config=(db.sink_config()),
                    table_config=(db.table_config())
                ))


if __name__ == '__main__':
    logging.getLogger().setLevel(logging.INFO)
    run()

内容的提问来源于stack exchange,提问作者infiniteloop314

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 19:35:25