如何在PySpark中将数组列转换为单调递减序列?
PySpark实现数组转单调递减序列
可以实现这个需求,下面提供两种可行方案:
方案一:使用自定义UDF
通过编写用户自定义函数(UDF)遍历数组,维护单调递减的序列逻辑:
- 创建示例DataFrame
from pyspark.sql import SparkSession from pyspark.sql import functions as F from pyspark.sql.types import ArrayType, IntegerType spark = SparkSession.builder.appName("monotonic_decrease").getOrCreate() data = [(1, [50,10,5,20,2]), (2, [42,10,15,5,3])] df = spark.createDataFrame(data, ["ID", "Pct"])
- 定义并注册UDF
def make_monotonic_decrease(arr): if not arr: return arr result = [arr[0]] for num in arr[1:]: # 取当前元素与前一个结果的较小值,保证序列递减 result.append(min(num, result[-1])) return result monotonic_decrease_udf = F.udf(make_monotonic_decrease, ArrayType(IntegerType()))
- 应用UDF生成目标列
df_result = df.withColumn("Mon_dec_pct", monotonic_decrease_udf(F.col("Pct"))) df_result.show(truncate=False)
方案二:使用PySpark内置高阶函数(推荐)
利用aggregate高阶函数实现逻辑,避免UDF的序列化开销,性能更优:
df_result = df.withColumn( "Mon_dec_pct", F.expr(""" aggregate( Pct, cast(array() as array<int>), (acc, x) -> concat(acc, array(if(size(acc) == 0, x, least(x, acc[size(acc)-1])))), acc -> acc ) """) ) df_result.show(truncate=False)
两种方案执行后都能得到预期结果:
+---+-------------------+-------------------+ |ID |Pct |Mon_dec_pct | +---+-------------------+-------------------+ |1 |[50, 10, 5, 20, 2] |[50, 10, 5, 5, 2] | |2 |[42, 10, 15, 5, 3] |[42, 10, 10, 5, 3] | +---+-------------------+-------------------+
内容的提问来源于stack exchange,提问作者Peter Yang
相关产品推荐
相关产品推荐

