使用Apache Beam将Pub/Sub数据写入BigQuery时出现表未找到警告
问题场景
我正在处理从Pub/Sub获取数据的业务,使用Apache Beam的ReadFromPubSub方法读取数据,Pub/Sub发布的消息格式如下:
b"('A', 'Stream2', 10)" b"('B', 'Stream1', 14)" b"('D', 'Stream3', 16)"
需求是将Stream1的数据写入BigQuery的dflow_stream1表,Stream2、3的数据写入dflow_stream23表,因此使用side_outputs实现分流。以下是实现代码:
import json import os import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.options.pipeline_options import GoogleCloudOptions from apache_beam.io.gcp.internal.clients import bigquery table_spec1 = bigquery.TableReference( projectId=<PROJECT_ID>, datasetId='training', tableId='dflow_stream1') table_spec23 = bigquery.TableReference( projectId=<PROJECT_ID>, datasetId='training', tableId='dflow_stream23') SCHEMA = { "fields": [ { "name": 'name', "type": "STRING", }, { "name": 'stream', "type": "STRING" }, { "name": 'salary', "type": "INT64", "mode": "NULLABLE" } ] } pipeline_options = PipelineOptions( streaming=True) class ProcessWords(beam.DoFn): def process(self, ele): Name,Stream,Salary=eval(ele) if Stream=="Stream1": yield {"Name":Name,"Stream":Stream,"Salary":Salary} else: yield beam.pvalue.TaggedOutput('Stream23', {"Name":Name,"Stream":Stream,"Salary":Salary}) class word_split(beam.DoFn): def process(selff,ele): Name,Stream,Salary=eval(ele) yield {"Name":Name,"Stream":Stream,"Salary":Salary} with beam.Pipeline(options=pipeline_options) as p: out= ( p | "Read from Pub/Sub subscription" >> beam.io.ReadFromPubSub(subscription="projects/<PROJECT_ID>/subscriptions/Test-sub") | "Decode and parse Json" >> beam.Map(lambda element: element.decode("utf-8")) |"Formatting " >> beam.ParDo(ProcessWords()).with_outputs("Stream23",main="Stream1") ) s1=out.Stream1 s23=out.Stream23 s1 | "Table1" >> beam.io.WriteToBigQuery( table=table_spec1, dataset='training', schema=SCHEMA, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, triggering_frequency=5, with_auto_sharding=True ) s23 |"Table23" >> beam.io.WriteToBigQuery( table=table_spec23, dataset='training', schema=SCHEMA, create_disposition=beam.io.BigQueryDisposition.CREATE_IF_NEEDED, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, triggering_frequency=5, with_auto_sharding=True ) p.run()
目前流水线已实现预期功能:从Pub/Sub读取数据、自动在数据集创建对应表并将数据追加到相应表中,但运行时会抛出如下警告:
WARNING:apache_beam.io.gcp.bigquery:There were errors inserting to BigQuery. Will retry. Errors were [{'index': 0, 'errors': [{'message': 'POST https://bigquery.googleapis.com/bigquery/v2/projects/<PROJECT_ID>/datasets/training/tables/dflow_stream23/insertAll?prettyPrint=false: Table 431017404487:training.dflow_stream23 not found.', 'reason': 'Not Found'}]}]
原因分析
- 表创建与写入的时序竞争:当流水线首次运行时,
WriteToBigQuery的CREATE_IF_NEEDED配置会先检查表是否存在,不存在则触发创建。但BigQuery的表创建属于异步操作,需要时间完成元数据同步和表初始化。而写入请求可能在表完全就绪前就发起,导致出现"表不存在"的错误。 - 内置重试机制的正常触发:Apache Beam的BigQuery写入组件自带重试逻辑,遇到这类临时的"表不存在"错误时会自动重试,因此最终数据能成功写入,只是会先抛出警告。
- 分流逻辑的影响:如果Stream2/3的数据先到达写入节点,对应表
dflow_stream23还在创建过程中,就会比Stream1的写入更早触发该警告;反之如果Stream1的数据先到,dflow_stream1先完成创建,后续写入就不会触发警告。
内容的提问来源于stack exchange,提问作者Ajay S Pal
相关产品推荐
相关产品推荐

