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

PySpark中基于窗口的线性回归计算时间序列斜率报错求助

问题分析

你遇到的AttributeError是因为lr.fit(vector_df).coefficients[0]返回的是一个全局训练得到的单一数值(numpy.float64类型),它不是Spark的列表达式,无法通过.over(sliding_window)实现窗口级别的计算。你当前的代码是对全量数据训练了一个线性回归模型,而不是针对每个ID的滑动窗口分别训练。

解决方案

针对每个ID的滑动窗口计算时间序列斜率,有两种高效的实现方式:

方法一:利用单变量线性回归的解析公式(推荐,性能更高)

单变量线性回归的斜率有直接的数学公式,不需要调用Spark ML模型,用窗口聚合函数就能实现:
斜率公式:

slope = (n*sum(xy) - sum(x)*sum(y)) / (n*sum(x²) - (sum(x))²)

其中:

  • n:窗口内的样本数量
  • x:时间戳(date_integer)
  • y:目标序列值(series)

代码实现

from pyspark.sql import functions as F
from pyspark.sql.window import Window

# 定义滑动窗口:按ID分组,按日期排序,取当前行及前12行(注意:如果要严格回溯12个月,建议用rangeBetween按时间戳范围,而非行数)
sliding_window = Window.partitionBy('ID').orderBy('date').rowsBetween(-12, 0)

# 转换日期为时间戳整数
df = df.withColumn("date_integer", F.unix_timestamp(df['date']))

# 窗口聚合计算所需的统计量
df_stats = df.withColumn("xy", F.col("date_integer") * F.col("series")) \
             .withColumn("x_squared", F.col("date_integer") ** 2) \
             .withColumn("n", F.count("series").over(sliding_window)) \
             .withColumn("sum_x", F.sum("date_integer").over(sliding_window)) \
             .withColumn("sum_y", F.sum("series").over(sliding_window)) \
             .withColumn("sum_xy", F.sum("xy").over(sliding_window)) \
             .withColumn("sum_x2", F.sum("x_squared").over(sliding_window))

# 计算斜率,处理分母为0的情况(比如窗口内数据不足)
df_result = df_stats.withColumn(
    "slope_window",
    F.when(
        F.col("n") >= 2,  # 至少需要2个样本才能计算斜率
        (F.col("n") * F.col("sum_xy") - F.col("sum_x") * F.col("sum_y")) / 
        (F.col("n") * F.col("sum_x2") - F.col("sum_x") ** 2)
    ).otherwise(None)  # 样本不足时返回空
)

# 可选:清理中间列
df_result = df_result.drop("date_integer", "xy", "x_squared", "n", "sum_x", "sum_y", "sum_xy", "sum_x2")

注意事项

如果你的需求是严格回溯12个月(而非固定12行),需要把窗口改成按时间戳范围计算:

# 计算12个月对应的秒数(近似值,忽略闰年)
twelve_months_seconds = 12 * 30 * 24 * 3600
sliding_window = Window.partitionBy('ID').orderBy(F.unix_timestamp('date')) \
                       .rangeBetween(-twelve_months_seconds, 0)

方法二:用Pandas UDF实现窗口级线性回归(适合多变量场景)

如果需要处理多变量线性回归,或者更习惯用ML模型,可以用Pandas UDF对每个ID的滑动窗口应用回归:

from pyspark.sql import functions as F
from pyspark.sql.types import DoubleType
from sklearn.linear_model import LinearRegression
import pandas as pd

# 定义滑动窗口
sliding_window = Window.partitionBy('ID').orderBy('date').rowsBetween(-12, 0)

# 转换日期为时间戳
df = df.withColumn("date_integer", F.unix_timestamp(df['date']))

# 定义Pandas UDF,对每个窗口的x和y计算斜率
@F.pandas_udf(DoubleType())
def calculate_slope(x: pd.Series, y: pd.Series) -> pd.Series:
    if len(x) < 2:
        return pd.Series([None])
    lr = LinearRegression()
    lr.fit(x.values.reshape(-1, 1), y.values)
    return pd.Series([lr.coef_[0]])

# 应用UDF到窗口
df_result = df.withColumn(
    "slope_window",
    calculate_slope(F.col("date_integer"), F.col("series")).over(sliding_window)
)
关键提醒
  • 不要直接用Spark ML的fit方法全局训练模型,再试图应用到窗口:fit返回的是训练好的模型实例,其系数是全局值,无法和窗口函数结合。
  • 优先选择方法一,因为基于公式的计算比调用ML模型性能高很多,尤其在大数据量场景下。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 06:15:14