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

如何用Pandas/PySpark实现传感器时序数据的列化聚合汇总

传感器时序数据汇总表生成方案

问题背景

现有百万个单传感器小型时序数据表,需生成单传感器一行的汇总表(按sensor uuid分组),需解决两个核心问题:

  1. 将Python风格的复杂分组统计逻辑转换为Pandas/PySpark的简洁group.agg()形式,生成单传感器的一行汇总数据;
  2. 将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_colmin_colmax_colcol_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.5156320...00.260.230.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.21 22:18:29