如何在PySpark DataFrame中将字符串列转为字典列表(不使用StructType)
问题描述
给定如下数据集:
data = [ { "employee_id": 873547690, "employee_name":"abc", "emp_deptname": "sales", "emp_endpoint": "abc@abc01storage", "emp_folder": "messages", "emp_address": [ { "type": "Home", "city": "Delhi", "apartment_number": "H.Number 124, B block, XYZ Towers", "pincode": "123abc" }, { "type": "Postbox", "city": "Delhi", "apartment_number": "Post Office 12A, Sector 22", "pincode": "456xyz" } ], 'event_timestamp': '1995-07-16 13:10:43' }, ]
只能使用以下字符串类型Schema创建PySpark DataFrame:
schema = "employee_id long, employee_name string, emp_deptname string, emp_endpoint string, emp_folder string, emp_address string, event_timestamp string"
创建后emp_address字段会被转为字符串,丢失原有的字典列表结构,直接转list会按字符拆分,无法满足后续操作需求。需要在不使用StructType/StructField的前提下,将该字符串字段转回字典列表。
解决方案
可以使用PySpark的from_json函数,配合DDL格式的字符串schema来解析emp_address字段,无需手动定义StructType:
- 先创建初始DataFrame:
from pyspark.sql import SparkSession spark = SparkSession.builder.appName("ParseEmpAddress").getOrCreate() df = spark.createDataFrame(data, schema=schema)
- 定义
emp_address对应的DDL字符串schema(描述嵌套结构):
emp_address_ddl = "array<struct<type:string, city:string, apartment_number:string, pincode:string>>"
- 使用
from_json解析字符串字段,替换原字段或新增字段:
from pyspark.sql.functions import from_json # 替换原emp_address字段 df_parsed = df.withColumn("emp_address", from_json(df.emp_address, emp_address_ddl)) # 验证结构 df_parsed.printSchema() df_parsed.select("emp_address").show(truncate=False)
处理后emp_address会恢复为包含结构体的数组类型,后续可直接进行扁平化、字段提取等操作:
# 示例:扁平化地址字段 from pyspark.sql.functions import explode df_flat = df_parsed.withColumn("address", explode("emp_address")) \ .select("employee_id", "employee_name", "address.*") df_flat.show(truncate=False)
内容的提问来源于stack exchange,提问作者LearneR
相关产品推荐
相关产品推荐

