如何去除PySpark DataFrame中array类型date字段的重复值得到预期输出
PySpark数组字段去重实现方案
你可以直接使用PySpark内置的array_distinct函数对date数组字段做去重,再提取去重后数组的唯一值即可,实现代码如下:
适配你场景的最优实现(保留原有分组统计结果,不合并记录)
不需要修改你原有的聚合逻辑,直接对聚合后的结果做字段处理即可:
from pyspark.sql.functions import count, array_distinct, col # 你原有的业务逻辑 timeseries_monthly = spark.read.options(header='True',inferschema='True',delimiter=',').parquet("url...") date = timeseries_monthly.select(timeseries_monthly["gps.date"]) grouped_df = date.groupBy('date').agg(count('date').alias('date_count')) # 新增date字段处理逻辑 result = grouped_df.withColumn("date", array_distinct(col("date"))[0]) # 输出结果 result.show(4, truncate=False)
说明
array_distinct是PySpark 2.4及以上版本的内置函数,会自动对数组内的重复元素做去重,你的场景中去重后每个数组仅包含1个唯一日期- 后续通过索引
[0]直接取出数组内的唯一值,即可得到单值格式的date字段,同时不会改动原有统计出来的date_count值,完全匹配你的预期输出
低版本PySpark兼容方案
如果你使用的PySpark版本低于2.4,没有内置array_distinct函数,可以自定义UDF实现相同效果:
from pyspark.sql.types import ArrayType, StringType from pyspark.sql.functions import udf # 自定义数组去重UDF def dedup_arr(arr): return list(set(arr)) dedup_udf = udf(dedup_arr, ArrayType(StringType())) # 处理字段 result = grouped_df.withColumn("date", dedup_udf(col("date"))[0]) result.show(4, truncate=False)
内容的提问来源于stack exchange,提问作者Nabih Bawazir
相关产品推荐
相关产品推荐

