如何在PySpark中从含多字典的字符串列提取指定数据生成新列
PySpark数据转换:提取指定条件的JSON数组元素
问题背景
现有PySpark DataFrame,schema如下:
id: string类型features: string类型,存储包含多个字典的JSON数组
输入示例:
| id | features |
|---|---|
| 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'的字典,转换为如下结构化格式:
| id | source | bus_name | type |
|---|---|---|---|
| a | abc | business_342 | tpA |
| b | abc | business_574 | tpB |
| c | abc | business_112 | tpC |
实现代码
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()
代码说明
- 解析JSON字符串:用
from_json将features列的字符串内容转换为符合预定义schema的数组结构,让非结构化的JSON字符串可被Spark操作。 - 展开数组:通过
explode把数组中的每个字典拆分为单独行,实现对单个元素的精准过滤。 - 条件过滤:筛选出
source字段为abc的记录,保留目标数据。 - 字段提取:从展开后的字典中提取所需字段并设置别名,最终得到符合要求的结构化DataFrame。
内容的提问来源于stack exchange,提问作者suraj jadhav
相关产品推荐
相关产品推荐

