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

如何用PySpark Streaming重构Kafka JSON数据流

PySpark Streaming 重构Kafka JSON数据流为嵌套格式解决方案

问题分析

输入为扁平结构的JSON数据流,需按timestamp聚合,将同时间戳下不同type(oxygen/temp)的指标分别嵌套为对象,每个对象内以name为键、value为值。原代码存在冗余操作、流处理API误用、聚合逻辑错误的问题,无法生成目标嵌套格式。

修正后的完整代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, map_from_arrays, collect_list
from pyspark.sql.types import StructType, StringType, LongType

# 初始化SparkSession
spark = SparkSession.builder.appName("KafkaJsonNestedTransform").getOrCreate()

# 读取Kafka流数据
df = (
    spark.readStream.format("kafka")
    .option("kafka.bootstrap.servers", "broker:29092")
    .option("subscribe", "testing")
    .load()
    .selectExpr("CAST(value AS STRING)")  # 仅需一次转换为字符串
)

# 定义JSON数据的Schema(StructType更规范,避免字符串解析的潜在问题)
json_schema = StructType() \
    .add("name", StringType()) \
    .add("value", StringType()) \
    .add("timestamp", LongType()) \
    .add("type", StringType())

# 解析JSON字符串为结构化数据(流处理专用from_json,替代批处理的read.json)
parsed_df = df.select(from_json(col("value"), json_schema).alias("data")).select("data.*")

# 第一步聚合:按timestamp和type分组,收集name和value的列表并转成键值对Map
type_grouped_df = parsed_df.groupBy("timestamp", "type") \
    .agg(
        map_from_arrays(
            collect_list(col("name")),
            collect_list(col("value"))
        ).alias("metrics")
    )

# 第二步聚合:按timestamp分组,将不同type的metrics转为顶层嵌套字段
result_df = type_grouped_df.groupBy("timestamp") \
    .pivot("type") \
    .agg(col("metrics")) \
    .select(
        col("timestamp"),
        col("oxygen"),
        col("temp")
    )

# 输出到控制台(流处理模式可按需调整)
output_query = (
    result_df.writeStream
    .outputMode("update")  # 输出更新后的数据;若需仅输出新增批次用append,全量输出用complete
    .format("console")
    .option("truncate", False)  # 避免控制台截断长JSON内容
    .start()
)

output_query.awaitTermination()

代码关键说明

  • 流处理适配:使用from_json替代批处理的spark.read.json,支持持续处理Kafka流数据
  • Map类型聚合:通过map_from_arrays直接将收集到的name和value列表转为键值对Map,完美匹配目标嵌套结构
  • 分层聚合逻辑:先按timestamp+type聚合生成各类型的指标Map,再按timestamp聚合并通过pivot将type转为顶层字段,最终得到期望的嵌套格式

输出验证

运行代码后,控制台输出将与目标格式一致:

+-------------+-----------------------------+-----------------------------------+
|timestamp    |oxygen                       |temp                               |
+-------------+-----------------------------+-----------------------------------+
|1699258974900|{position_y -> 1.9, position_H -> 3.6}|{position_y -> 1.9}|
|169925900    |{position_x -> 5.1}          |{position_z -> 3.999, position_x -> 5.1}|
+-------------+-----------------------------+-----------------------------------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 04:57:44