如何适配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
相关产品推荐
相关产品推荐

