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

如何将基于循环的流程步骤追踪逻辑转换为可扩展PySpark实现

PySpark 无循环实现产品流程进度时间序列统计

核心思路

抛弃循环式逐批次计算,改用Spark原生的宽表转长表、日期序列生成、窗口函数累计、透视表组合方案,所有计算在单一分布式执行计划内完成,彻底解决集群扩展性问题。

实现步骤及代码

1. 初始化环境与测试数据

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

spark = SparkSession.builder.appName("ProcessStepTracking").getOrCreate()

# 测试输入数据(模拟产品各步骤到达日期)
test_data = [
    ("USA", "BrandA", "ModelX", "2024-01-01", "2024-01-03", "2024-01-05", "2024-01-07", "2024-01-09"),
    ("USA", "BrandA", "ModelX", "2024-01-02", "2024-01-04", "2024-01-06", "2024-01-08", "2024-01-10"),
    ("Canada", "BrandB", "ModelY", "2024-01-01", "2024-01-02", None, None, None),
    ("Canada", "BrandB", "ModelY", "2024-01-03", "2024-01-04", "2024-01-05", None, None)
]

schema = ["country", "brand", "model", "step1_date", "step2_date", "step3_date", "step4_date", "step5_date"]
df = spark.createDataFrame(test_data, schema=schema)

2. 宽表转长表

将每个产品的多步骤日期列转换为(步骤-日期)的长格式,统一后续计算维度:

melted_df = df.select(
    "country", "brand", "model",
    F.explode(
        F.array(
            *[F.struct(F.lit(f"step{i}").alias("step"), F.col(f"step{i}_date").alias("date")) for i in range(1,6)]
        )
    ).alias("step_info")
).select(
    "country", "brand", "model",
    F.col("step_info.step"),
    F.col("step_info.date").cast("date")
).filter(F.col("date").isNotNull())

3. 生成分组日期序列

为每个(国家、品牌、型号)分组生成完整的日期范围(从该组最早步骤日期到最晚步骤日期):

# 计算每个分组的日期边界
date_range_df = melted_df.groupBy("country", "brand", "model")\
    .agg(F.min("date").alias("min_date"), F.max("date").alias("max_date"))

# 生成连续日期序列
all_dates_df = date_range_df.select(
    "country", "brand", "model",
    F.explode(F.sequence(F.col("min_date"), F.col("max_date"), F.expr("interval 1 day"))).alias("date")
)

4. 计算每日累计步骤数量

用窗口函数按步骤分组累计,确保每日的步骤数量包含所有之前到达该步骤的产品:

window_spec = Window.partitionBy("country", "brand", "model", "step")\
    .orderBy("date")\
    .rowsBetween(Window.unboundedPreceding, Window.currentRow)

step_counts_df = melted_df.groupBy("country", "brand", "model", "step", "date")\
    .agg(F.count("*").alias("daily_add"))\
    .join(all_dates_df, on=["country", "brand", "model", "date"], how="right")\
    .fillna(0, subset=["daily_add"])\
    .withColumn("cumulative_count", F.sum("daily_add").over(window_spec))\
    .select("country", "brand", "model", "date", "step", "cumulative_count")

5. 透视回目标格式

将长格式数据转换为每日各步骤列的时间序列:

final_df = step_counts_df.groupBy("country", "brand", "model", "date")\
    .pivot("step")\
    .agg(F.first("cumulative_count"))\
    .fillna(0)\
    .orderBy("country", "brand", "model", "date")

# 查看结果
final_df.show()

方案优势

  • 完全基于Spark原生算子,无循环,避免多次作业提交的开销
  • 分布式执行计划自动优化,充分利用集群资源,支持大数据量扩展
  • 代码简洁可维护,逻辑清晰,便于后续调整步骤数量或日期范围

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 18:10:37