如何用SQL基于10分钟时间窗口统计各地点同时在岗唯一人员数量
业务数据说明
现有超过10万行结构如下的业务数据,包含字段:timestamp(时间戳)、Location(地点)、person(人员),样例数据如下:
timestamp Location person 2017-09-04 08:07:00 UTC A x 2017-09-04 08:08:00 UTC B y 2017-09-04 08:09:00 UTC A y 2017-09-04 08:07:00 UTC A x 2017-09-04 08:27:00 UTC B x
预期统计结果
需要输出按Location分组的各地点最大同时在岗人员数量,样例如下:
Location Nr_of_persons_working_at_the_same_time A 2 B 1
统计规则
同一地点内,不同人员的操作记录时间差不超过10分钟即视为同时在岗;时间差超过10分钟判定为换班,不计入同时在岗统计。规则标注参考:
timestamp Location person 2017-09-04 08:07:00 UTC A x <--- 人员x在A地点的第一条操作记录 2017-09-04 08:08:00 UTC B y <--- 人员y在B地点的第一条操作记录 2017-09-04 08:09:00 UTC A y <--- A地点的第二条操作记录,此时无法确认人员x是否已离岗 2017-09-04 08:07:00 UTC A x <--- 人员x仍有同时间段记录,因此A地点同时在岗人数为2 2017-09-04 08:27:00 UTC B x <--- 人员x在B地点的记录与上一条记录间隔20分钟,属于换班,因此B地点同时在岗人数仍为1
背景与已遇到的问题
- 数据已通过SQL查询获取,优先使用SQL实现统计逻辑;受限于10万行数据量级与云端计算资源,不推荐高开销方案
- 已尝试方案存在的问题:
- 按
location、timestamp分组统计存在时间硬切割问题,统计结果不符合预期 - 尝试使用窗口函数实现,但按
timestamp排序后无法避免不同地点的数据混淆
- 按
实现方案
SQL方案(优先推荐)
该方案通过事件拆分逻辑计算重叠在岗人数,用分区窗口隔离不同地点数据,完全解决现有问题,10万行数据计算效率极高,适配Hive、Spark SQL、PostgreSQL等多数支持标准SQL的引擎:
WITH person_valid_records AS ( -- 去重同一地点、同一人员、同一时间的重复记录,避免统计干扰 SELECT DISTINCT Location, person, timestamp AS start_time, timestamp + INTERVAL '10' MINUTE AS end_time FROM 你的原表名 ), event_points AS ( -- 将每个人员的10分钟在岗区间拆为进入(+1)、离开(-1)两个事件点 SELECT Location, start_time AS event_time, 1 AS delta FROM person_valid_records UNION ALL SELECT Location, end_time AS event_time, -1 AS delta FROM person_valid_records ), rolling_online_cnt AS ( -- 按地点分区、事件时间排序,滚动计算当前在岗人数 SELECT Location, SUM(delta) OVER ( PARTITION BY Location ORDER BY event_time ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS current_cnt FROM event_points ) -- 取每个地点的最大同时在岗人数 SELECT Location, MAX(current_cnt) AS Nr_of_persons_working_at_the_same_time FROM rolling_online_cnt GROUP BY Location ORDER BY Location;
注:如果使用的SQL引擎不支持
INTERVAL语法,可替换为对应时间计算函数,例如MySQL可用DATE_ADD(timestamp, INTERVAL 10 MINUTE)。
Python方案(备选)
如果确实需要用Python实现,可按地点分组后再计算,避免跨地点数据干扰,内存占用可控:
import pandas as pd # 读取数据,将timestamp转为datetime类型,此处替换为你的数据读取逻辑 df = pd.read_csv("你的数据文件路径") df["timestamp"] = pd.to_datetime(df["timestamp"]) result = [] # 按地点分组独立计算 for location, loc_group in df.groupby("Location"): # 去重同一人同一时间的重复记录 unique_records = loc_group.drop_duplicates(subset=["person", "timestamp"]) events = [] for _, row in unique_records.iterrows(): # 生成进入、离开事件 events.append((row["timestamp"], 1)) events.append((row["timestamp"] + pd.Timedelta(minutes=10), -1)) # 按时间排序事件 events.sort(key=lambda x: x[0]) max_cnt = current_cnt = 0 for _, delta in events: current_cnt += delta max_cnt = max(max_cnt, current_cnt) result.append({ "Location": location, "Nr_of_persons_working_at_the_same_time": max_cnt }) # 输出结果 result_df = pd.DataFrame(result) print(result_df)
内容的提问来源于stack exchange,提问作者Charles
相关产品推荐
相关产品推荐

