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

PySpark结构化流:拆分Kafka Value嵌套JSON并写入Delta Lake

拆分Kafka JSON Value并写入Delta Lake的实现方案

核心步骤说明

针对你提供的嵌套JSON结构(包含嵌套对象与数组),需通过解析JSON、展平嵌套/数组结构、写入Delta Lake三个核心步骤完成,以下是基于PySpark的具体实现:


1. 读取Kafka数据并转换Value列

Kafka的value列默认是二进制类型,需先转换为字符串格式,才能进行JSON解析:

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, from_json, explode
from pyspark.sql.types import StructType, StructField, DoubleType, IntegerType, LongType

# 初始化SparkSession(需确保Delta Lake依赖已加载)
spark = SparkSession.builder \
    .appName("KafkaToDelta") \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .config("spark.sql.catalog.spark_catalog", "org.apache.spark.sql.delta.catalog.DeltaCatalog") \
    .getOrCreate()

# 读取Kafka数据(批处理示例,流处理替换为readStream)
kafka_df = spark.read \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "your-kafka-broker:9092") \
    .option("subscribe", "your-topic-name") \
    .load()

# 将二进制value转换为字符串
json_str_df = kafka_df.select(col("value").cast("string").alias("json_value"))

2. 定义JSON Schema并解析嵌套结构

根据你提供的JSON示例,手动定义Schema(避免自动推断的性能问题):

# 定义完整的JSON Schema
components_schema = StructType([
    StructField("co", DoubleType()),
    StructField("no", DoubleType()),
    StructField("no2", DoubleType()),
    StructField("o3", DoubleType()),
    StructField("so2", DoubleType()),
    StructField("pm2_5", DoubleType()),
    StructField("pm10", DoubleType()),
    StructField("nh3", DoubleType())
])

main_schema = StructType([
    StructField("aqi", IntegerType())
])

list_item_schema = StructType([
    StructField("main", main_schema),
    StructField("components", components_schema),
    StructField("dt", LongType())
])

root_schema = StructType([
    StructField("coord", StructType([
        StructField("lon", DoubleType()),
        StructField("lat", DoubleType())
    ])),
    StructField("list", list_item_schema)
])

# 解析JSON字符串为结构化DataFrame
parsed_df = json_str_df.select(from_json(col("json_value"), root_schema).alias("data"))

3. 展平嵌套与数组结构

由于list是数组类型,需用explode拆分为单行记录,同时将嵌套字段提取为顶级列,方便后续SQL查询与MLlib操作:

# 展平coord字段 + 拆分list数组
flattened_df = parsed_df.select(
    col("data.coord.lon").alias("lon"),
    col("data.coord.lat").alias("lat"),
    explode(col("data.list")).alias("list_item")
)

# 提取list_item内的所有字段
final_df = flattened_df.select(
    "lon",
    "lat",
    col("list_item.main.aqi").alias("aqi"),
    col("list_item.components.co").alias("co"),
    col("list_item.components.no").alias("no"),
    col("list_item.components.no2").alias("no2"),
    col("list_item.components.o3").alias("o3"),
    col("list_item.components.so2").alias("so2"),
    col("list_item.components.pm2_5").alias("pm2_5"),
    col("list_item.components.pm10").alias("pm10"),
    col("list_item.components.nh3").alias("nh3"),
    col("list_item.dt").alias("timestamp")
)

# 可选:将时间戳转换为可读格式
from pyspark.sql.functions import from_unixtime
final_df = final_df.withColumn("readable_timestamp", from_unixtime(col("timestamp")))

4. 写入Delta Lake

支持批处理与流处理两种模式,写入后可直接通过SQL查询或MLlib加载数据:

批处理写入

final_df.write \
    .format("delta") \
    .mode("append") \
    .save("/path/to/delta-table")  # 或使用saveAsTable("database.table_name")

流处理写入(实时场景)

query = final_df.writeStream \
    .format("delta") \
    .option("checkpointLocation", "/path/to/checkpoint") \
    .mode("append") \
    .start("/path/to/delta-table")

query.awaitTermination()

后续操作验证

写入完成后,可通过SQL查询验证结果:

SELECT * FROM delta.`/path/to/delta-table` LIMIT 10;

也可直接加载为DataFrame用于MLlib操作:

delta_df = spark.read.format("delta").load("/path/to/delta-table")
# 后续可直接用于特征工程、模型训练等MLlib流程

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 01:02:42