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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.27 06:00:10