从BigQuery导入的分组DataFrame中Pandas时间窗口滚动计算结果异常
问题描述
致歉:本问题无法复现——将DataFrame转为字典再转回DataFrame后,异常消失。
BigQuery查询语句
SELECT published_at, from_author_id, text FROM `project.message.message`
数据转DataFrame代码
client = bigquery.Client(location="europe-west1", project="project") df = client.query(sql).to_dataframe()
出错的滚动计数代码及错误结果
执行以下代码会得到错误结果:
import pandas as pd #df['published_at'] = pd.to_datetime(df['published_at']) df = df.sort_values(by=['from_author_id', 'published_at']) df.groupby('from_author_id').rolling('3s', on='published_at')['text'].count()
使用pd.to_datetime()对滚动函数的结果无影响。错误输出示例:
from_author_id published_at 0001fcf4-94f5-4e42-8444-0cb6c2870bdc 2024-08-19 18:28:50.197000+00:00 1.0 2024-08-19 18:33:26.837000+00:00 2.0 2024-08-19 18:33:42.960000+00:00 3.0 2024-08-19 18:33:57.083000+00:00 4.0 2024-08-19 18:34:18.863000+00:00 5.0 ... fff7a574-a2fe-4eac-b7c6-d5de8dc5ff0c 2024-08-19 16:26:24.252000+00:00 6.0 2024-08-19 16:32:40.697000+00:00 7.0 2024-08-19 16:32:42.013000+00:00 8.0 2024-08-19 18:09:03.469000+00:00 1.0 2024-08-19 18:09:04.979000+00:00 2.0
可见第一个作者的每条消息间隔均超过3秒,滚动计数应返回1,但结果却累加。
可得到正确结果的自定义函数
以下自定义函数可得到预期结果:
def compute_correct_rolling_count(df, window_seconds=3): msg_counts = [] for _, group_df in df.groupby('from_author_id'): count_list = [] for i in range(len(group_df)): start_time = group_df.iloc[i]['published_at'] - pd.Timedelta(seconds=window_seconds) count = group_df[(group_df['published_at'] > start_time) & (group_df['published_at'] <= group_df.iloc[i]['published_at'])].shape[0] count_list.append(count) msg_counts.extend(count_list) return msg_counts # 计算每个作者3秒窗口内的滚动消息数 df['msg_count_last_3secs'] = compute_correct_rolling_count(df, window_seconds=3)
表结构与额外信息
BigQuery表project.message.message的Schema:
- published_at (TIMESTAMP)
- from_author_id (STRING)
- text (STRING)
- 其他字段
额外信息:
- 默认舍入模式:ROUNDING_MODE_UNSPECIFIED
- 按DAY分区
- 分区字段:published_at
- 使用的google-cloud-bigquery版本:3.25.0
内容的提问来源于stack exchange,提问作者user1627466
相关产品推荐
相关产品推荐

