Spark中如何提取嵌套struct列字段并解决parallelize序列化报错?
问题说明
当前Spark作业存在两个待解决问题:
- 生成的DataFrame所有业务字段均嵌套在名为
table的STRUCT类型单列中,需要将列内所有字段提取为顶层独立列 - 执行
Spark.parallelize(table)操作时抛出报错:
TypeError: cannot pickle '_thread.RLock' object
现有相关代码与结构信息如下:
table.printSchema()输出结构:
root |-- table: struct (nullable = true) | |-- name: string (nullable = true) | |-- school: date (nullable = true) | |-- studentid: string (nullable = true) | |-- class: integer (nullable = true) | |-- grade: string (nullable = true) | |-- age: string (nullable = true)
- 定义解析schema的代码:
schema1 = StructType([ StructField("name", StringType(), True), StructField("school", DateType(), True), StructField("StudentID", StringType(), True), StructField("class", StringType(), True), StructField("grade", IntegerType(), True), StructField("age", StringType(), True), ])
- 解析JSON生成
table列的代码:
table = df.select(from_json(col("body").cast("string"), schema1).alias("table")) table.printSchema()
解决方案
提取STRUCT嵌套列为独立顶层列
两种实现方式,按需选择:
- 全量快速展开:直接用
.*语法选中struct下所有字段,无需逐个枚举
# 打平table列所有字段为独立列 flatten_df = table.select("table.*") # 验证输出结构 flatten_df.printSchema()
- 自定义选择/重命名字段:如果只需要部分字段,或者需要给字段起别名,逐个指定字段路径即可
flatten_df = table.select( col("table.name"), col("table.school"), col("table.studentid").alias("student_id"), col("table.class"), col("table.grade"), col("table.age") )
注意:你当前定义的schema1存在两处不匹配问题,会导致字段解析为null:一是字段名大小写不一致,schema中写的是StudentID,实际JSON字段和输出schema均为小写studentid;二是字段类型不匹配,schema中class定义为StringType、grade定义为IntegerType,实际解析结果中class为integer类型、grade为string类型,请根据原始JSON的实际格式对齐schema定义。
修复pickle线程锁报错
这个报错的根因是用法错误:SparkContext.parallelize()的入参只能是本地Python集合(如list、dict、tuple等原生Python对象),你传入的table是Spark DataFrame分布式对象,内部持有SparkContext的线程锁实例,本身不支持pickle序列化,因此无法传入parallelize。
不需要用parallelize处理现有DataFrame:
- 如果只是要打平嵌套列,直接使用上述select方法处理DataFrame即可,无需转RDD
- 如果确实需要将DataFrame转为RDD,直接调用DataFrame自带的
.rdd属性即可,不要使用parallelize:
# DataFrame转RDD的正确写法 table_rdd = table.rdd # 打平后的DataFrame转RDD flatten_rdd = flatten_df.rdd
内容的提问来源于stack exchange,提问作者Justin Rigger
相关产品推荐
相关产品推荐

