You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在PySpark中将数组列转换为单调递减序列?

PySpark实现数组转单调递减序列

可以实现这个需求,下面提供两种可行方案:

方案一:使用自定义UDF

通过编写用户自定义函数(UDF)遍历数组,维护单调递减的序列逻辑:

  1. 创建示例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"])
  1. 定义并注册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()))
  1. 应用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 18:13:21