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

如何在PySpark中从含多字典的字符串列提取指定数据生成新列

PySpark数据转换:提取指定条件的JSON数组元素

问题背景

现有PySpark DataFrame,schema如下:

  • id: string类型
  • features: string类型,存储包含多个字典的JSON数组

输入示例:

idfeatures
a[{"source":"abc","bus_name":"business_342","type":"tpA"},{"source":"pqr","bus_name":"busness_342","type":"tpP"}]
b[{"source":"abc","bus_name":"business_574","type":"tpB"},{"source":"pqr","bus_name":"busness_896","type":"tpQ"}]
c[{"source":"abc","bus_name":"business_112","type":"tpC"},{"source":"pqr","bus_name":"busness_312","type":"tpR"}]

需要提取features数组中source='abc'的字典,转换为如下结构化格式:

idsourcebus_nametype
aabcbusiness_342tpA
babcbusiness_574tpB
cabcbusiness_112tpC

实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import from_json, explode, col
from pyspark.sql.types import StructType, StructField, StringType

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

# 模拟输入数据
data = [
    ("a", '[{"source":"abc","bus_name":"business_342","type":"tpA"},{"source":"pqr","bus_name":"busness_342","type":"tpP"}]'),
    ("b", '[{"source":"abc","bus_name":"business_574","type":"tpB"},{"source":"pqr","bus_name":"busness_896","type":"tpQ"}]'),
    ("c", '[{"source":"abc","bus_name":"business_112","type":"tpC"},{"source":"pqr","bus_name":"busness_312","type":"tpR"}]')
]

# 创建初始DataFrame
df = spark.createDataFrame(data, schema=["id", "features"])

# 定义features数组中单个字典的schema
feature_schema = StructType([
    StructField("source", StringType(), nullable=True),
    StructField("bus_name", StringType(), nullable=True),
    StructField("type", StringType(), nullable=True)
])

# 解析JSON字符串 -> 展开数组 -> 过滤目标记录 -> 提取字段
result_df = df.withColumn("parsed_features", from_json(col("features"), feature_schema)) \
              .withColumn("exploded_features", explode(col("parsed_features"))) \
              .filter(col("exploded_features.source") == "abc") \
              .select(
                  col("id"),
                  col("exploded_features.source").alias("source"),
                  col("exploded_features.bus_name").alias("bus_name"),
                  col("exploded_features.type").alias("type")
              )

# 查看结果
result_df.show()

代码说明

  1. 解析JSON字符串:用from_json将features列的字符串内容转换为符合预定义schema的数组结构,让非结构化的JSON字符串可被Spark操作。
  2. 展开数组:通过explode把数组中的每个字典拆分为单独行,实现对单个元素的精准过滤。
  3. 条件过滤:筛选出source字段为abc的记录,保留目标数据。
  4. 字段提取:从展开后的字典中提取所需字段并设置别名,最终得到符合要求的结构化DataFrame。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 03:32:39