时间序列组相似性计算:如何将VBA逻辑迁移为pandas/pyspark实现
实现方案选型与落地代码
适用场景选型
- 数据量<100万行、单节点内存可容纳:优先选pandas,开发成本低、逻辑调整灵活
- 数据量>100万行、分布式存储:选pyspark,支持水平扩展,适配大数据量计算需求
注意:以下代码默认日度数值字段命名为
day_1到day_N,如果是200+天的跨度,仅需修改day_cols为所有数值字段的列名列表即可,计算逻辑不需要调整,完全兼容你的数据结构。
1. Pandas实现方案
逻辑说明
Excel的CORREL函数本质是皮尔逊相关系数,计算时自动忽略双端均为有效值之外的缺失值,和你的需求完全匹配;得分并列时取最新日期的规则可以通过排序后去重实现。
完整代码
import pandas as pd import numpy as np # 1. 数据读取,history_df为全量历史序列,test_df为待检测序列 day_cols = [f'day_{i}' for i in range(1, 8)] # 200+天则替换为所有数值字段列名 common_cols = ['GROUP', '序列日期'] + day_cols # 2. 按分组关联待检测和历史序列,避免跨组匹配 merged_df = test_df[common_cols].merge( history_df[common_cols], on='GROUP', suffixes=('', '_history') ) # 3. 逐行计算相关系数,仅取双端非空的有效值 def calc_correl(row): test_vals = [row[col] for col in day_cols] hist_vals = [row[f'{col}_history'] for col in day_cols] # 过滤双端都非空的数值对 valid_pairs = [(t, h) for t, h in zip(test_vals, hist_vals) if pd.notna(t) and pd.notna(h)] if len(valid_pairs) < 2: # 有效值不足2个无法计算相关系数 return np.nan t_list, h_list = zip(*valid_pairs) return np.corrcoef(t_list, h_list)[0][1] merged_df['Similarity_Score'] = merged_df.apply(calc_correl, axis=1) # 4. 按得分降序、历史日期降序排序,取每个待检测序列的Top1匹配 result_df = merged_df.sort_values( by=['Similarity_Score', '序列日期_history'], ascending=[False, False] ).drop_duplicates(subset=['GROUP', '序列日期'], keep='first') # 5. 格式化输出指定字段 output_df = result_df.rename(columns={ '序列日期_history': 'Similarity_Week', 'GROUP': 'Group' }) output_df = output_df[['序列日期'] + day_cols + ['Similarity_Score', 'Similarity_Week', 'Group']]
2. PySpark实现方案
逻辑说明
大数据量场景下通过广播小表(历史序列表)减少shuffle开销,用窗口函数取Top1匹配,性能远高于pandas。
完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, array, udf, desc, row_number from pyspark.sql.window import Window import numpy as np spark = SparkSession.builder.appName("TimeSeriesMatch").getOrCreate() # 1. 读取数据 history_df = spark.read.parquet("历史数据存储路径") test_df = spark.read.parquet("待检测数据存储路径") day_cols = [f'day_{i}' for i in range(1, 8)] # 200+天则替换为所有数值字段列名 # 2. 广播历史表,按分组关联 joined_df = test_df.join(history_df.hint("broadcast"), on='GROUP', how='inner') # 3. 自定义UDF计算相关系数,逻辑对齐Excel CORREL @udf(returnType="double") def calc_correl_udf(test_arr, hist_arr): valid_pairs = [(t, h) for t, h in zip(test_arr, hist_arr) if t is not None and h is not None] if len(valid_pairs) < 2: return None t_list, h_list = zip(*valid_pairs) return float(np.corrcoef(t_list, h_list)[0][1]) joined_df = joined_df.withColumn( "Similarity_Score", calc_correl_udf(array(*[col(c) for c in day_cols]), array(*[col(c) for c in day_cols])) ) # 4. 窗口函数取Top1匹配 window_spec = Window.partitionBy("GROUP", "序列日期").orderBy(desc("Similarity_Score"), desc("序列日期_history")) result_df = joined_df.withColumn("rn", row_number().over(window_spec)).filter(col("rn") == 1).drop("rn") # 5. 格式化输出 output_df = result_df.selectExpr( "序列日期", *day_cols, "Similarity_Score", "序列日期_history as Similarity_Week", "GROUP as Group" )
其他适用的相似性度量方法
- 余弦相似度:适合你当前的占比类字段,对数值绝对值差异不敏感,仅匹配趋势方向,运算速度快
- 动态时间规整(DTW):如果允许日度序列存在小范围偏移匹配(比如周1的特征和周2对齐也可判定为相似),DTW比皮尔逊相关系数的适配性更强
- 曼哈顿距离归一化:如果需要同时考虑数值大小的绝对差异,可以将曼哈顿距离归一化后转为0-1区间的相似性得分使用
内容的提问来源于stack exchange,提问作者Aussie_Stats
相关产品推荐
相关产品推荐

