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

Dataflow加载Pub/Sub数据至BigQuery无数据写入问题求助

Dataflow写入BigQuery无数据问题排查与修复

核心问题1:数据格式不匹配(最可能原因)

你的parse_json_message返回的是JSON字符串(比如[{"KEY_DATA": "xxx", ...}]),但WriteToBigQuery期望接收的是Python字典对象,而非序列化后的字符串。这会导致BigQuery无法正确解析数据,要么被判定为无效行,要么直接跳过写入。

修复方案:

  • 修改parse_json_message,返回字典列表而非JSON字符串:
    def parse_json_message(message):
        row = json.loads(message)
        df = pd.json_normalize(row, max_level=4)
        # 重新赋值确保操作生效
        df = df.fillna('').astype(str)
        df.columns = [x.strip().replace('-','_').replace('.','_').upper() for x in df.columns]
        select_df = df[["KEY_DATA", "VALUE_DATA_KV_QUOTE_LASTPRICEDDATE", "VALUE_DATA_KV_QUOTE_REPRICEONOPEN", "VALUE_DATA_KV_QUOTE_LOCKED"]]
        # 返回字典列表,而非JSON字符串
        return select_df.to_dict('records')
    
  • 将beam.Map(parse_json_message)改为beam.FlatMap(parse_json_message),因为to_dict('records')返回的是列表,FlatMap会把列表中的每个字典作为单独元素输出,符合WriteToBigQuery的要求。

核心问题2:Pandas操作未生效

代码中df.fillna('').astype(str)没有重新赋值给df,Pandas的方法默认返回新对象,原DataFrame不会被修改,导致空值仍存在,可能后续写入时引发隐性错误。

修复:

df = df.fillna('').astype(str)

窗口触发延迟问题

你使用了60秒固定窗口,Dataflow默认的窗口触发策略是窗口结束后才触发输出。如果测试时刚启动任务,数据可能还在窗口缓存中,未到达触发时间,所以BigQuery暂时看不到数据。

临时测试方案:

  • 缩小窗口大小(比如改为10秒):
    | "Fixed-size-windows" >> beam.WindowInto(window.FixedWindows(10))
    
  • 或者添加提前触发策略,处理时间10秒后就触发:
    from apache_beam import window
    from apache_beam.transforms.trigger import AfterProcessingTime, AccumulationMode
    
    | "Fixed-size-windows" >> beam.WindowInto(
        window.FixedWindows(60),
        trigger=AfterProcessingTime(10),
        accumulation_mode=AccumulationMode.DISCARDING
    )
    

错误日志排查

检查你的错误日志表da-dwh-dev:test_dataset_for_all.error_log是否有数据:

  • 如果有错误记录,根据error_message字段直接定位问题(比如字段类型不匹配、必填字段缺失等)。
  • 如果没有错误记录,在数据流中添加日志输出验证数据:
    | "Log parsed data" >> beam.Map(lambda x: print(f"Parsed data: {x}"))
    

其他排查点

  • 确认Pub/Sub订阅确实有消息被推送:通过Pub/Sub控制台查看订阅的消息投递情况。
  • 确认BigQuery表的table_schema与输出字典的字段完全匹配(字段名、类型一致,大小写敏感)。
  • 检查Dataflow Worker日志:在GCP控制台的Dataflow任务详情中,查看是否有隐性异常未被捕获。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.16 20:23:18