Apache Beam流式写入BigQuery时additional_bq_parameters可调用对象不生效问题
问题结论
这个不是你的配置错误,属于Apache Beam当前版本的特性差异:批量FILE_LOADS写入模式已经支持additional_bq_parameters传入可调用对象实现动态参数,但STREAMING_INSERTS流式写入模式暂未实现该逻辑。
根因说明
你排查源码的结果是正确的:
- 批量写入对应的
BigQueryBatchFileLoads类中,已经兼容了可调用类型的additional_bq_parameters,会自动传入目标表参数执行函数拿到实际配置 - 流式写入对应的实现逻辑中,没有添加对
additional_bq_parameters的可调用判断,会直接把传入的函数对象当做参数进行**解包,所以抛出_MessageClass object argument after ** must be a mapping, not function报错
临时适配方案
你可以根据自己的业务场景选择对应的适配方式:
- 单表写入场景:直接传入参数映射即可,不需要传函数,示例:
# 提前调用函数拿到对应表的参数 bq_params = additional_bq_parameters_fn("你的目标表完整名称") (processed_rows | "Write to Bigquery Streaming" >> apache_beam.io.WriteToBigQuery( table_fn, schema=schema_fn, additional_bq_parameters=bq_params, write_disposition=apache_beam.io.BigQueryDisposition.WRITE_APPEND, create_disposition=apache_beam.io.BigQueryDisposition.CREATE_IF_NEEDED, method=apache_beam.io.WriteToBigQuery.Method.STREAMING_INSERTS) )
- 多表动态路由写入场景:可以先把数据按目标表分组,每组调用函数拿到对应参数后再分别执行写入操作。
你可以到Apache Beam Jira提交对应的特性需求工单,补全流式写入和批量写入的参数兼容逻辑。
内容的提问来源于stack exchange,提问作者Djuls
相关产品推荐
相关产品推荐

