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
相关产品推荐
相关产品推荐

