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

Pandas UDF实现时间序列聚类未并行化问题排查求助

时间序列聚类最优K值计算的Spark并行化问题

问题背景

拥有20余万个类别的时间序列数据,需通过计算不同聚类数K对应的样本到质心平方距离和(inertia)确定最优K值。尝试基于K值用Pandas UDF实现分布式并行计算,但Spark UI仅显示单个任务持续运行,仅利用了45个可用工作节点中的2个。已尝试重分区操作,数据无倾斜,此前Pandas UDF并行化工作正常,怀疑问题与TimeSeriesKMeans函数相关。

问题代码

# Read data
stream = spark.read.format('delta').load('/mnt/datadrop-test/fw_fish/stream_bronze/clustering/cvr_cpa_clicks_4hourly_smoothed')

# Create an array per category
stream_array = stream.groupBy(['advertiser_id', 'campaign_name', 'keyword_text', 'keyword_id'])
                     .agg(F.collect_list('cvr_cpa_clicks').alias('cvr_cpa_clicks'))

# create a dataframe with the number of clusters to try
# these are the possible values of K 
n_clusters = spark.range(5, 500, 5).withColumnRenamed("id","n_cluster")

# cross join them
cross_joined1 = stream_array.crossJoin(F.broadcast(n_clusters)).repartition('n_cluster')


def get_sum_of_squared_distances(df):
    k = df.n_cluster.values[0]
    df_array = df.cvr_cpa_clicks.values 
    nrows = df_array.shape[0]
    df_array = np.concatenate(df_array).reshape(nrows, -1)
    km = TimeSeriesKMeans(n_clusters = k,
                      max_iter = 500,
                      n_init = 50,
                      metric = "euclidean",
                      n_jobs = -1,
                      verbose = False,
                      max_iter_barycenter = 500,
                      random_state = 0,
                      init = 'k-means++')

    km = km.fit(df_array)
    result_df = pd.DataFrame({'n_clusters': [k], 'sum_of_squared_distances': [km.inertia_]})
    return result_df
   
# Get the sum of squared distances for each k(number of clusters)
sum_of_squared_distances = cross_joined1.groupBy('n_cluster').applyInPandas(get_sum_of_squared_distances, schema = "n_clusters long, sum_of_squared_distances double")

display(sum_of_squared_distances)

问题原因

  1. TimeSeriesKMeans多线程与Spark资源冲突:n_jobs=-1会让每个Pandas UDF任务占用当前Executor的所有CPU核心,导致Spark无法调度更多任务到其他节点,只能同时运行少量任务(对应2个节点)。
  2. 分区策略未充分匹配K值数量:仅按n_cluster重分区,但未明确指定分区数等于K值的总数量,可能导致多个K值被分配到同一个分区,无法并行处理。
  3. TimeSeriesKMeans参数过于激进:n_init=50、max_iter=500等参数大幅增加了单任务计算耗时,进一步延缓了整体并行进度。

解决方案

1. 禁用TimeSeriesKMeans内部多线程

将n_jobs=1,让每个Pandas UDF任务仅占用单核心,Spark可同时调度更多任务到不同节点,充分利用集群资源。

2. 调整分区策略

明确设置分区数等于K值的总数量(从5到500步长5,共99个K值),确保每个K值对应独立分区,实现完全并行。

3. 优化TimeSeriesKMeans参数

降低初始化次数和迭代次数,在保证结果稳定性的前提下减少单任务计算时间。

修改后的代码

# 读取数据
stream = spark.read.format('delta').load('/mnt/datadrop-test/fw_fish/stream_bronze/clustering/cvr_cpa_clicks_4hourly_smoothed')

# 按类别聚合时间序列数组
stream_array = stream.groupBy(['advertiser_id', 'campaign_name', 'keyword_text', 'keyword_id'])
                     .agg(F.collect_list('cvr_cpa_clicks').alias('cvr_cpa_clicks'))

# 生成要测试的K值列表
n_clusters = spark.range(5, 500, 5).withColumnRenamed("id","n_cluster")
# 获取K值总数,用于设置分区数
k_total = n_clusters.count()

# 交叉连接并按K值重分区,确保每个K值对应独立分区
cross_joined1 = stream_array.crossJoin(F.broadcast(n_clusters)).repartition(k_total, 'n_cluster')


def get_sum_of_squared_distances(df):
    k = df.n_cluster.values[0]
    # 转换为NumPy数组(替代原concatenate方式,更简洁)
    df_array = np.array(df.cvr_cpa_clicks.tolist())
    # 初始化时间序列KMeans,禁用内部多线程并优化参数
    km = TimeSeriesKMeans(n_clusters=k,
                          max_iter=100,
                          n_init=10,
                          metric="euclidean",
                          n_jobs=1,
                          verbose=False,
                          max_iter_barycenter=100,
                          random_state=0,
                          init='k-means++')

    km.fit(df_array)
    return pd.DataFrame({'n_clusters': [k], 'sum_of_squared_distances': [km.inertia_]})
   
# 按K值分组并行计算平方距离和
sum_of_squared_distances = cross_joined1.groupBy('n_cluster').applyInPandas(
    get_sum_of_squared_distances, 
    schema="n_clusters long, sum_of_squared_distances double"
)

display(sum_of_squared_distances)

额外优化建议

  • 类别抽样:若全量20万类别计算耗时仍过高,可对stream_array做抽样(如sample(fraction=0.1)),用10%的类别计算最优K值,结果趋势与全量数据基本一致,能大幅缩短计算时间。
  • 资源配置检查:确保Spark Executor有足够内存,避免时间序列聚类加载样本时出现OOM错误。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 18:35:06