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

如何在PySpark中转换含字符串格式字典列表的流DataFrame?

实现PySpark流式DataFrame结构转换

要将你的原始流式DataFrame转换成目标结构,按以下步骤操作:

1. 定义JSON Schema

先明确JSON数组的结构,定义对应的Spark Schema:

from pyspark.sql.types import StructType, StructField, IntegerType, ArrayType

# 定义数组内每个元素的结构体Schema
json_schema = ArrayType(
    StructType([
        StructField("num", IntegerType(), nullable=False),
        StructField("cor", IntegerType(), nullable=False)
    ])
)

2. 处理流式DataFrame

针对原始DataFrame执行数据转换:

from pyspark.sql.functions import col, from_json, explode, regexp_replace

# 假设原始流式DataFrame名为 raw_stream_df
processed_df = raw_stream_df \
    # 1. 去除JSON字符串首尾的双引号,并重命名字段方便后续操作
    .withColumn("json_str", regexp_replace(col("`stringdecode(value, UTF-8)`"), '^"|"$', '')) \
    # 2. 将清洗后的字符串解析为Spark数组类型
    .withColumn("json_array", from_json(col("json_str"), json_schema)) \
    # 3. 将数组拆分为多行,实现一对多展开
    .withColumn("json_item", explode(col("json_array"))) \
    # 4. 提取目标字段,保留原始时间戳和偏移量
    .select(
        "timestamp",
        "offset",
        col("json_item.num").alias("num"),
        col("json_item.cor").alias("cor")
    )

关键细节说明

  • 字段引用:原始JSON字段名包含特殊字符(括号、空格),必须用反引号``包裹才能正确引用。
  • 清洗JSON字符串:原始数据中的JSON被双引号包裹(如"[{...}]"),用regexp_replace去除首尾引号,确保from_json能正常解析。
  • 数组拆分:explode函数会将数组中的每个元素拆成单独的行,对应原始每条数据生成多行结果。

执行以上代码后,就能得到你需要的目标DataFrame结构。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 23:08:15