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

从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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:00:55