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

时间序列组相似性计算:如何将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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.04 03:30:01