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

Spark DataFrame中Double数组列的近似分位数计算方法

Solution for Per-Row Array Quantile Calculation in Spark

Since df.stat.approxQuantile is built to compute quantiles across an entire column (not per-row arrays), we need workarounds to calculate the 0.25th and 0.75th quantiles for each array in your amt_list column. Here are two reliable, practical methods:

Method 1: Use a User-Defined Function (UDF)

This is the most straightforward approach—we’ll create a UDF that takes an array, processes it, and returns the desired quantiles. We’ll use numpy for the quantile calculation (just make sure it’s installed on all your Spark workers).

Python Example:

from pyspark.sql import SparkSession
from pyspark.sql.functions import udf
from pyspark.sql.types import StructType, StructField, DoubleType
import numpy as np

# Initialize Spark session
spark = SparkSession.builder.appName("ArrayQuantiles").getOrCreate()

# Define schema for UDF output (holds both q25 and q75)
quantile_schema = StructType([
    StructField("q25", DoubleType(), nullable=True),
    StructField("q75", DoubleType(), nullable=True)
])

# UDF to compute quantiles for a single array
def calculate_array_quantiles(arr):
    # Handle empty or tiny arrays to avoid errors
    if not arr or len(arr) < 4:
        return (None, None)
    sorted_arr = sorted(arr)
    q25 = np.percentile(sorted_arr, 25, method="linear")
    q75 = np.percentile(sorted_arr, 75, method="linear")
    return (float(q25), float(q75))

array_quantile_udf = udf(calculate_array_quantiles, quantile_schema)

# Apply UDF and unpack the result into separate columns
result_df = df.withColumn("quantiles", array_quantile_udf(df.amt_list)) \
              .select("id", "amt_list", "ct_tran_amt", "quantiles.q25", "quantiles.q75")

result_df.show(truncate=False)

For Scala Users:

You can write a similar UDF using Scala’s built-in sorting and libraries like org.apache.commons.math3.stat.descriptive.rank.Percentile to compute quantiles without relying on external tools like numpy.

Method 2: Spark SQL Native Approach (No UDFs)

If you want to avoid UDFs (for better performance with large datasets), you can unnest the arrays, compute quantiles per id using window functions, then pivot the results back to a single row per id.

Python Example:

from pyspark.sql import Window
from pyspark.sql.functions import explode, percent_rank, first

# Step 1: Unnest the amt_list array into individual rows
exploded_df = df.select("id", "amt_list", "ct_tran_amt", explode(df.amt_list).alias("amt"))

# Step 2: Define a window partitioned by id, ordered by the transaction amount
window_spec = Window.partitionBy("id").orderBy("amt")

# Step 3: Calculate percent rank for each value in the array
ranked_df = exploded_df.withColumn("p_rank", percent_rank().over(window_spec))

# Step 4: Filter for rows near the 25th and 75th percentiles (use epsilon for edge cases)
quantile_rows = ranked_df.filter(
    (ranked_df.p_rank >= 0.24) & (ranked_df.p_rank <= 0.26) |
    (ranked_df.p_rank >= 0.74) & (ranked_df.p_rank <= 0.76)
)

# Step 5: Pivot results to get q25 and q75 as separate columns
result_df = quantile_rows.groupBy("id", "amt_list", "ct_tran_amt") \
                         .pivot("p_rank", [0.25, 0.75]) \
                         .agg(first("amt")) \
                         .withColumnRenamed("0.25", "q25") \
                         .withColumnRenamed("0.75", "q75")

result_df.show(truncate=False)

Notes:

  • The 0.01 epsilon in the filter handles cases where the exact percentile doesn’t align with a row’s percent_rank (common with small arrays). Adjust it based on your precision needs.
  • This method leverages Spark’s optimized native operations, making it faster for large-scale data compared to Python UDFs (which cross the JVM boundary).

Which Method Should You Pick?

  • Go with the UDF approach for simplicity, especially with small-to-medium arrays.
  • Use the Spark SQL native method for better performance with large datasets or if you want to avoid external dependencies.

内容的提问来源于stack exchange,提问作者user1124702

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 08:12:39