如何用Pandas/PySpark实现传感器时序数据的列化聚合汇总
传感器时序数据汇总表生成方案
问题背景
现有百万个单传感器小型时序数据表,需生成单传感器一行的汇总表(按sensor uuid分组),需解决两个核心问题:
- 将Python风格的复杂分组统计逻辑转换为Pandas/PySpark的简洁
group.agg()形式,生成单传感器的一行汇总数据; - 将
df.groupby("b").agg(avg=("y", "mean"), cnt=("b", "count")).reset_index()的行式聚合结果转换为列式结构(如col_0(cnt)、col_0(avg)),替代Python循环的非优雅实现。
原始示例数据表
y a b 1 0 1 0 2 0 1 0 2 1 0 2 0 4 0 1 0 1 0 6 0 0 6 0 0 6 0 1 0 1 0 8 0
现有Python统计函数
def get_stats(df): g = df.groupby("a").size() c = g.count() a = df["y"].sum() b = df["y"].count() avg_col = (b - a) / (c - 1) min_col = g.min() max_col = g.max() return avg_col, min_col, max_col
期望汇总表结构
| avg_col | min_col | max_col | col_0(cnt) | col_1(cnt) | col_2(cnt) | col_3(cnt) | ... | col_n(cnt) | col_0(avg) | col_1(avg) | col_2(avg) | ... | col_n(avg) |
|---|---|---|---|---|---|---|---|---|---|---|---|---|---|
| 1.5 | 1 | 5 | 6 | 3 | 2 | 0 | ... | 0 | 0.26 | 0.23 | 0.69 | ... | 0 |
解决方案
1. 把统计逻辑改成Pandas/PySpark的agg简洁写法
Pandas版本
直接通过一次agg计算所有中间统计值,再派生目标字段,避免分步循环:
import pandas as pd # 示例数据 df = pd.DataFrame([ [1,0,1], [0,2,0], [1,0,2], [1,0,2], [0,4,0], [1,0,1], [0,6,0], [0,6,0], [0,6,0], [1,0,1], [0,8,0] ], columns=["y", "a", "b"]) # 单传感器一行汇总 sensor_summary = ( df .agg( y_sum=("y", "sum"), y_count=("y", "count"), g_min=("a", lambda x: x.groupby(x).size().min()), g_max=("a", lambda x: x.groupby(x).size().max()), g_count=("a", lambda x: x.groupby(x).size().count()) ) .assign(avg_col=lambda x: (x["y_count"] - x["y_sum"]) / (x["g_count"] - 1)) [["avg_col", "g_min", "g_max"]] .rename(columns={"g_min": "min_col", "g_max": "max_col"}) )
输出为单个传感器的一行统计结果:
avg_col 1.5 min_col 1.0 max_col 5.0 dtype: float64
PySpark版本
PySpark分两步聚合再合并,逻辑更清晰:
from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder.getOrCreate() # 示例数据 df = spark.createDataFrame([ (1,0,1), (0,2,0), (1,0,2), (1,0,2), (0,4,0), (1,0,1), (0,6,0), (0,6,0), (0,6,0), (1,0,1), (0,8,0) ], ["y", "a", "b"]) # 统计a分组后的size的min/max/count a_group_stats = df.groupBy("a").agg(F.count("*").alias("a_size")).agg( F.min("a_size").alias("min_col"), F.max("a_size").alias("max_col"), F.count("a_size").alias("g_count") ) # 统计y的全局sum和count y_global_stats = df.agg(F.sum("y").alias("y_sum"), F.count("y").alias("y_count")) # 合并计算avg_col sensor_summary = y_global_stats.crossJoin(a_group_stats).withColumn( "avg_col", (F.col("y_count") - F.col("y_sum")) / (F.col("g_count") - 1) ).select("avg_col", "min_col", "max_col")
输出结果:
+-------+-------+-------+ |avg_col|min_col|max_col| +-------+-------+-------+ | 1.5| 1| 5| +-------+-------+-------+
2. 行式聚合转列式结构(替代循环)
针对df.groupby("b").agg(avg=("y", "mean"), cnt=("b", "count")).reset_index()的结果,用透视表直接转成目标列式结构:
Pandas版本
# 先做b的分组聚合 b_grouped = df.groupby("b").agg(avg=("y", "mean"), cnt=("b", "count")).reset_index() # 转成列式结构 pivot_result = ( b_grouped .melt(id_vars="b", var_name="stat_type", value_name="value") .assign(col_name=lambda x: f"col_{x['b']}({x['stat_type']})") .pivot(columns="col_name", values="value") .fillna(0) ) # 和之前的传感器汇总合并 final_summary = pd.concat([sensor_summary, pivot_result], axis=1)
输出结果(对应示例数据):
avg_col min_col max_col col_0(cnt) col_0(avg) col_1(cnt) col_1(avg) col_2(cnt) col_2(avg) 0 1.5 1.0 5.0 6.0 0.000000 3.0 1.000000 2.0 1.000000
PySpark版本
# 先做b的分组聚合 b_grouped = df.groupBy("b").agg( F.mean("y").alias("avg"), F.count("b").alias("cnt") ) # 转成列式结构 pivot_result = ( b_grouped .select("b", F.explode(F.array( F.struct(F.lit("cnt").alias("stat_type"), F.col("cnt").alias("value")), F.struct(F.lit("avg").alias("stat_type"), F.col("avg").alias("value")) )).alias("stats") .select("b", "stats.stat_type", "stats.value") .withColumn("col_name", F.concat(F.lit("col_"), F.col("b"), F.lit("("), F.col("stat_type"), F.lit(")"))) .groupBy() .pivot("col_name") .agg(F.first("value")) .na.fill(0) ) # 合并传感器汇总 final_summary = sensor_summary.crossJoin(pivot_result)
输出结果:
+-------+-------+-------+----------+----------+----------+----------+----------+----------+ |avg_col|min_col|max_col|col_0(cnt)|col_0(avg)|col_1(cnt)|col_1(avg)|col_2(cnt)|col_2(avg)| +-------+-------+-------+----------+----------+----------+----------+----------+----------+ | 1.5| 1| 5| 6| 0.0| 3| 1.0| 2| 1.0| +-------+-------+-------+----------+----------+----------+----------+----------+----------+
内容的提问来源于stack exchange,提问作者Yiffany
相关产品推荐
相关产品推荐

