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

Spark(PySpark 1.6.3)中实现可变窗口大小的滚动求和问题

Hey there, let's break down your problem and find a working solution for PySpark 1.6.3!

First, why your original attempts failed

The core issue here is a limitation in PySpark 1.6's window function implementation:

  • The rangeBetween/rowsBetween methods only accept constant numeric values as window boundaries—you can't use column expressions (like rank_id +1) to define dynamic per-row windows. That's why both your PySpark API and SQL attempts threw errors.
  • For your final "simple SQL" test, the error happens because Spark 1.6's RANGE BETWEEN syntax expects either UNBOUNDED PRECEDING/UNBOUNDED FOLLOWING or specific numeric offsets, but it's also not designed for dynamic boundaries anyway.

Solution: Use a self-join + groupBy

Since we can't use window functions for dynamic per-row windows, we can achieve the same result by joining the DataFrame to itself, filtering for the desired rank range, then grouping to calculate the sum. This approach keeps all processing distributed (no collecting data to local) and works with large datasets.

Here's the step-by-step code:

from pyspark.mllib.random import RandomRDDs
import pyspark.sql.functions as psf
from pyspark.sql.window import Window

# 1. Generate your original data (same as your code)
data = RandomRDDs.uniformVectorRDD(sc, 15, 2)
df = data.map(lambda l: (float(l[0]), float(l[1]))).toDF()
df = df.selectExpr("_1 as x", "_2 as y")

# 2. Add rank_id (same as your code)
w = Window().orderBy("x")
df = df.withColumn("rank_id", psf.rowNumber().over(w)).sort("rank_id")

# 3. Self-join to get matching rows in the dynamic window
# Join the DataFrame to itself, where b's rank_id falls in [a.rank_id+1, 2*a.rank_id-1]
df_joined = df.alias("a").join(
    df.alias("b"),
    (psf.col("b.rank_id") >= psf.col("a.rank_id") + 1) & 
    (psf.col("b.rank_id") <= 2 * psf.col("a.rank_id") - 1),
    how="left"  # Keep all rows from the original DataFrame
)

# 4. Calculate the rolling sum, handling cases with no matching rows (sum = 0)
df_result = df_joined.groupBy("a.rank_id", "a.x", "a.y") \
    .agg(psf.coalesce(psf.sum("b.y"), psf.lit(0)).alias("roll_var")) \
    .orderBy("rank_id")

# Check the result
df_result.show()

How this works:

  1. Self-join: We alias the original DataFrame as a (the base rows) and b (the rows we want to include in the window for each a row). The join condition ensures we only pull b rows where rank_id is within your desired range for each a row.
  2. Left join: This guarantees we keep every row from the original DataFrame, even if there are no matching b rows (e.g., when rank_id=1, 2*1-1=1 which is less than 1+1=2—no rows qualify, so sum is null).
  3. Coalesce: We use coalesce to replace null sums with 0, which gives a clean result for rows with no matching window data.
  4. GroupBy: Grouping by all columns from a ensures we retain the original row data while calculating the sum of b.y for each row's window.

Performance note for large datasets

Since rank_id is a sequential integer (generated via rowNumber()), the join condition is efficient—each a row only joins to a small, contiguous range of b rows. Spark can optimize this range-based join, so it won't create an expensive Cartesian product even for large datasets.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 07:55:32