使用Apache Beam从PubSub读取数据时遇JSON解析错误求助
问题解决:Pub/Sub到BigQuery的JSON解析错误
问题背景
作为Apache Beam新手,实操流程为:从JSON文件读取数据发送至Pub/Sub Topic,再用Beam从Pub/Sub读取数据写入BigQuery,但执行时出现JSON解析错误:
RuntimeError: json.decoder.JSONDecodeError: Expecting property name enclosed in double quotes: line 1 column 2 (char 1) [while running 'CustomParse']
错误根源分析
问题出在publishEvents.py的消息序列化步骤:
- 代码中先通过
json.loads(line)把JSON转成Python字典,然后用str(traffic)将字典转为字符串,这个字符串是Python语法的字典(键用单引号),不是标准JSON格式(键必须用双引号)。 - 消费端尝试用
json.loads()解析非标准JSON字符串,直接触发解析错误,后续的ast.literal_eval也无法补救这个前置错误。
修正方案
1. 修正发布端代码(publishEvents.py)
将消息序列化为标准JSON字符串,而不是Python字典的字符串表示:
import json from google.cloud import pubsub_v1 from concurrent import futures from typing import Callable project='ApacheBeam' topic='incomingData' publisher = pubsub_v1.PublisherClient() topic_path = publisher.topic_path(project, topic) publish_futures = [] def get_callback( publish_future: pubsub_v1.publisher.futures.Future, data: str ) -> Callable[[pubsub_v1.publisher.futures.Future], None]: def callback(publish_future: pubsub_v1.publisher.futures.Future) -> None: try: print(publish_future.result(timeout=60)) except futures.TimeoutError: print(f"Publishing {data} timed out.") return callback # 修正消息序列化逻辑 with open('/home/test/Resources/Input/Input_Data.json', 'r') as f: for line in f: # 方式1:直接用原JSON行(如果文件每行都是标准JSON) data = line.strip() # 方式2:如果需要修改字典后再序列化,用json.dumps # traffic = json.loads(line) # data = json.dumps(traffic) publish_future = publisher.publish(topic_path, data.encode("utf-8")) publish_future.add_done_callback(get_callback(publish_future, data)) publish_futures.append(publish_future) futures.wait(publish_futures, return_when=futures.ALL_COMPLETED) print(f"Published messages with error handler to {topic_path}.")
关键修改:用json.dumps()生成标准JSON字符串,或者直接使用文件中已有的标准JSON行,避免将Python字典转为非标准字符串。
2. 修正消费端代码(readFromPubSub.py)
简化解析逻辑,因为现在接收的是标准JSON,无需ast.literal_eval,同时修复json.dumps()后无法添加键的错误(json.dumps()返回字符串,不能直接加键):
import argparse import json import os import logging import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions, StandardOptions logging.basicConfig(level=logging.INFO) logging.getLogger().setLevel(logging.INFO) os.environ["GOOGLE_APPLICATION_CREDENTIALS"] = "home/test/Resources/Config/gcpkey.json" INPUT_SUBSCRIPTION = "projects/ApacheBeam/subscriptions/ReadData" BIGQUERY_TABLE = "ApacheBeam:JSONData.StoreEvents" BIGQUERY_SCHEMA = "user_id:STRING,event_ts:STRING,timestamp:STRING,event_id:STRING,ifa:STRING,ifv:STRING,country:STRING,chip_balance:STRING,game:STRING,user_group:STRING,user_condition:STRING,device_type:STRING,device_model:STRING,user_name:STRING,fb_connect:BOOLEAN,is_active_event:BOOLEAN,event_payload:STRING" class CustomParsing(beam.DoFn): def to_runner_api_parameter(self, unused_context): return "beam:transforms:custom_parsing:custom_v0", None def process(self, element: bytes, timestamp=beam.DoFn.TimestampParam): # 解析标准JSON字符串为字典 json_data = json.loads(element.decode("utf-8")) # 添加timestamp字段 json_data["timestamp"] = timestamp.to_rfc3339() # 返回字典,BigQuery Write支持直接接收字典 yield json_data def run(): parser = argparse.ArgumentParser() parser.add_argument( "--input_subscription", help='Input PubSub subscription of the form "projects/<PROJECT>/subscriptions/<SUBSCRIPTION>."', default=INPUT_SUBSCRIPTION, ) parser.add_argument( "--output_table", help="Output BigQuery Table", default=BIGQUERY_TABLE ) parser.add_argument( "--output_schema", help="Output BigQuery Schema in text format", default=BIGQUERY_SCHEMA, ) known_args, pipeline_args = parser.parse_known_args() pipeline_options = PipelineOptions(pipeline_args) pipeline_options.view_as(StandardOptions).streaming = True with beam.Pipeline(options=pipeline_options) as p: ( p | "ReadFromPubSub" >> beam.io.gcp.pubsub.ReadFromPubSub( subscription=known_args.input_subscription, timestamp_attribute=None ) | "CustomParse" >> beam.ParDo(CustomParsing()) | "WriteToBigQuery" >> beam.io.WriteToBigQuery( known_args.output_table, schema=known_args.output_schema, write_disposition=beam.io.BigQueryDisposition.WRITE_APPEND, ) ) if __name__ == "__main__": run()
关键修改:
- 移除多余的
ast.literal_eval,直接用json.loads()解析标准JSON - 修复
json.dumps()后无法添加字段的问题:直接操作字典,再返回字典(BigQuery的WriteToBigQuery支持接收字典类型) - 更新BigQuery Schema,将原
create_ts改为timestamp(对应添加的字段)
验证步骤
- 重新运行修正后的
publishEvents.py,确保消息以标准JSON格式发送到Pub/Sub - 启动Beam消费管道
readFromPubSub.py,验证数据能正常写入BigQuery
内容的提问来源于stack exchange,提问作者Data writer
相关产品推荐
相关产品推荐

