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

使用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(对应添加的字段)

验证步骤

  1. 重新运行修正后的publishEvents.py,确保消息以标准JSON格式发送到Pub/Sub
  2. 启动Beam消费管道readFromPubSub.py,验证数据能正常写入BigQuery

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 12:54:59