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)
问题原因
- TimeSeriesKMeans多线程与Spark资源冲突:
n_jobs=-1会让每个Pandas UDF任务占用当前Executor的所有CPU核心,导致Spark无法调度更多任务到其他节点,只能同时运行少量任务(对应2个节点)。 - 分区策略未充分匹配K值数量:仅按
n_cluster重分区,但未明确指定分区数等于K值的总数量,可能导致多个K值被分配到同一个分区,无法并行处理。 - 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
相关产品推荐
相关产品推荐

