PySpark新手求助:如何从JSON中筛选超龄记录及获取最大年龄
嘿,我来帮你搞定这个问题!你遇到的错误根源有两个:一是你的employees字段是数组类型,直接用df['employees.age']其实拿到的是一个数组列表,没法直接和整数22比较;二是你的age字段存的是字符串类型,就算数组问题解决了,字符串和整数比较也会有类型不匹配的问题。下面一步步给你解决:
第一步:解析JSON并处理嵌套结构
首先,你需要把嵌套的employees数组展开成单独的员工记录,同时把age字段从字符串转换成整数类型。假设你已经把JSON数据加载到了DataFramedf,先查看下原始数据的结构:
df.printSchema()
输出应该类似这样:
root |-- employees: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- age: string (nullable = true) | | |-- firstName: string (nullable = true) | | |-- lastName: string (nullable = true)
接下来用explode函数展开数组,并转换age的类型:
from pyspark.sql.functions import explode, col # 展开employees数组,将每个员工转为单独一行 exploded_df = df.select(explode(col("employees")).alias("employee")) # 拆分结构体字段,同时把age转为整数 clean_df = exploded_df.select( col("employee.age").cast("int").alias("age"), col("employee.firstName").alias("firstName"), col("employee.lastName").alias("lastName") ) # 查看处理后的结果 clean_df.show()
这时候clean_df里的每条记录都是一个独立的员工,age也变成了整数类型,后续操作就不会有类型不匹配的问题了。
第二步:筛选年龄大于指定值的记录
现在可以正常使用filter来筛选符合条件的记录了:
# 筛选年龄大于22的员工 filtered_df = clean_df.filter(col("age") > 22) filtered_df.show()
第三步:获取最大年龄
用agg函数配合max函数就能轻松拿到最大年龄:
from pyspark.sql.functions import max # 计算并获取最大年龄 max_age_result = clean_df.agg(max(col("age"))).collect()[0][0] print(f"所有员工中的最大年龄是:{max_age_result}")
可选:链式调用简化代码
如果你不想分步写变量,也可以把所有操作链式调用起来,更简洁:
from pyspark.sql.functions import explode, col, max # 链式筛选年龄>22的记录 df.select(explode(col("employees")).alias("employee")) \ .select(col("employee.age").cast("int").alias("age"), col("employee.firstName"), col("employee.lastName")) \ .filter(col("age") > 22) \ .show() # 链式计算最大年龄 max_age = df.select(explode(col("employees")).alias("employee")) \ .select(col("employee.age").cast("int").alias("age")) \ .agg(max(col("age"))) \ .collect()[0][0] print(f"最大年龄:{max_age}")
内容的提问来源于stack exchange,提问作者Prashant Patel
相关产品推荐
相关产品推荐

