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/rowsBetweenmethods only accept constant numeric values as window boundaries—you can't use column expressions (likerank_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 BETWEENsyntax expects eitherUNBOUNDED PRECEDING/UNBOUNDED FOLLOWINGor 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:
- Self-join: We alias the original DataFrame as
a(the base rows) andb(the rows we want to include in the window for eacharow). The join condition ensures we only pullbrows whererank_idis within your desired range for eacharow. - Left join: This guarantees we keep every row from the original DataFrame, even if there are no matching
brows (e.g., whenrank_id=1,2*1-1=1which is less than1+1=2—no rows qualify, so sum is null). - Coalesce: We use
coalesceto replace null sums with 0, which gives a clean result for rows with no matching window data. - GroupBy: Grouping by all columns from
aensures we retain the original row data while calculating the sum ofb.yfor 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

