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

基于PySpark从Spark DataFrame分组结果生成时间序列

按Party和CounterParty生成时间序列并分析的完整流程

Hey there! 既然你已经搞定了按Party和CounterParty分组的操作,那接下来的时间序列生成、可视化和深度分析其实可以分成几个清晰的步骤来走,我给你一步步拆解:

1. 补全完整的时间序列(处理缺失时间点)

原始数据里大概率存在某些时间点没有记录的情况,这会直接影响后续分析的准确性。我们需要为每一对分组生成连续的时间轴,步骤如下:

  • 先确定全局的时间范围(比如从数据中最早的日期到最晚的日期,按天/小时/分钟等业务需要的粒度)
  • 用Spark的sequence函数生成这个时间范围内的所有时间点
  • 将生成的时间序列和所有唯一的(Party, CounterParty)对做交叉连接,得到每一对分组的完整时间轴
  • 左连接原始数据,用合理的方式填充缺失值(比如0、前向值、均值等,依业务场景而定)

示例Scala代码:

import org.apache.spark.sql.functions._
import org.apache.spark.sql.types.TimestampType

// 假设你的原始DataFrame为df,包含timestamp、Party、CounterParty、value字段
val minTs = df.select(min("timestamp")).first().getAs[Timestamp](0)
val maxTs = df.select(max("timestamp")).first().getAs[Timestamp](0)

// 生成按天粒度的时间序列(可替换为hour/minute调整粒度)
val timeSeries = spark.sql(s"""SELECT sequence(to_timestamp('${minTs}'), to_timestamp('${maxTs}'), interval 1 day) as date""")
  .withColumn("timestamp", explode(col("date")))
  .drop("date")

// 获取所有唯一的(Party, CounterParty)组合
val uniquePairs = df.select("Party", "CounterParty").distinct()

// 交叉连接得到每对组合的完整时间轴
val fullTimeSeries = uniquePairs.crossJoin(timeSeries)

// 左连接原始数据,用0填充缺失的value(也可用窗口函数last()做前向填充)
val filledDf = fullTimeSeries.join(df, Seq("Party", "CounterParty", "timestamp"), "left")
  .na.fill(0, Seq("value"))

2. 数据预处理(保障分析质量)

在正式分析前,先做些基础预处理:

  • 确认timestamp列是TimestampType,如果不是就用to_timestamp()转换
  • 处理异常值:比如用四分位数法过滤极端值,或用中位数替换异常点
  • 若后续要用到机器学习模型,可对数值列做标准化/归一化处理

3. 时间序列可视化

Spark本身不支持直接绘图,我们可以把分组后的数据转成Pandas DataFrame,再用Matplotlib/Seaborn/Plotly来实现可视化:

  • 可以生成分面图,每个子图对应一对(Party, CounterParty)的时间序列趋势
  • 也可以聚焦重点组合单独绘制,方便对比差异

示例Python代码(Matplotlib):

import pandas as pd
import matplotlib.pyplot as plt

# 将Spark DataFrame转换为Pandas DataFrame
pd_df = filledDf.toPandas()

# 设置绘图风格
plt.style.use('seaborn-v0_8')

# 获取所有唯一的(Party, CounterParty)组合
unique_pairs = pd_df[['Party', 'CounterParty']].drop_duplicates().values.tolist()

# 生成子图布局,每行显示3个图
n_rows = (len(unique_pairs) + 2) // 3
fig, axes = plt.subplots(n_rows, 3, figsize=(18, n_rows*5))
axes = axes.flatten()

for idx, (party, counter_party) in enumerate(unique_pairs):
    # 筛选当前组合的数据
    subset = pd_df[(pd_df['Party'] == party) & (pd_df['CounterParty'] == counter_party)]
    # 绘制时间序列
    axes[idx].plot(subset['timestamp'], subset['value'], marker='o', linestyle='-', linewidth=1)
    axes[idx].set_title(f'Party: {party} | CounterParty: {counter_party}')
    axes[idx].tick_params(axis='x', rotation=45)

# 隐藏多余的空白子图
for idx in range(len(unique_pairs), len(axes)):
    axes[idx].axis('off')

plt.tight_layout()
plt.show()

4. 深度时间序列分析

完成可视化后,可以开展这些深度分析工作:

  • 趋势分析:用滑动窗口计算均值/标准差,观察长期趋势;或用线性回归拟合趋势线,判断增长/下降态势
  • 季节性分析:使用STL分解(季节-趋势-残差分解)识别周期性波动;或用傅里叶变换检测主要周期
  • 异常检测:用Spark MLlib的孤立森林(Isolation Forest)、LOF算法检测异常点;也可用3σ原则快速识别极端值
  • 相关性分析:计算不同(Party, CounterParty)组合时间序列的相关性,找出关联度高的配对
  • 预测:若需要预测未来值,可使用ARIMA、Prophet(需转Pandas)或Spark MLlib的时序预测模型

实用小贴士

  • 如果数据量极大,建议先对分组进行采样,先分析样本确定流程,再批量处理所有组合
  • 时间粒度的选择要贴合业务:日交易数据用天粒度,实时数据用小时/分钟粒度
  • 缺失值填充要结合业务场景:比如交易金额缺失可用0填充,累计类指标适合用前向值填充

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 10:19:00