PySpark中按ArrayType子元素统计distinct出现次数
解决PySpark中ArrayType子元素的distinct计数问题
我来帮你搞定这个统计需求!根据你的描述,虽然ArrayType的键看起来被扁平化了,但printSchema()能看到完整结构,说明数组里的元素本质是StructType,咱们只需要先把数组展开,再提取目标子元素统计就行,具体步骤如下:
步骤1:导入必要的PySpark函数
首先得导入用来展开数组的explode和用来统计次数的count函数:
from pyspark.sql.functions import explode, count
步骤2:展开数组列
假设你的数组类型列名叫array_col(你可以从df.printSchema()的输出里找到实际列名),用explode把数组拆成单行,每个数组元素对应一行:
# 把数组列展开,并重命名为elements方便后续访问子元素 exploded_df = df.select(explode("array_col").alias("elements"))
步骤3:分组统计element_y的出现次数
现在可以直接通过elements.element_y访问到目标子元素,然后分组计数:
# 按element_y分组,统计每组的行数(也就是出现次数) count_result = exploded_df.groupBy("elements.element_y").agg(count("*").alias("occurrences"))
如果需要排除element_y为null的无效数据,可以在分组前加个过滤:
filtered_exploded = exploded_df.filter("elements.element_y IS NOT NULL") count_result = filtered_exploded.groupBy("elements.element_y").agg(count("*").alias("occurrences"))
步骤4:输出成你想要的格式
最后把统计结果转换成你预期的22x: 2, 23x: 1这种格式:
# 把结果收集到本地,然后遍历输出 result_dict = {row["element_y"]: row["occurrences"] for row in count_result.collect()} for key, value in result_dict.items(): print(f"{key}: {value}", end=", ")
执行后就能得到和你预期一致的输出啦。
举个实际的例子,如果你的DataFrame schema是这样的:
root |-- id: integer (nullable = true) |-- array_col: array (nullable = true) | |-- element: struct (containsNull = true) | | |-- element_x: string (nullable = true) | | |-- element_y: string (nullable = true)
那上面的代码完全适用,只需要把array_col替换成你实际的数组列名就行。
内容的提问来源于stack exchange,提问作者Ivan Bilan
相关产品推荐
相关产品推荐

