【问题标题】:BigQuery - nest operations within partition in order to aggregate consecutive records from a table with 17 billion recordsBigQuery - 在分区内嵌套操作,以便从具有 170 亿条记录的表中聚合连续记录
【发布时间】:2019-05-13 20:44:03
【问题描述】:

我对 SQL 和 BigQuery 还很陌生,一周来一直在努力寻找可行的解决方案来解决这个问题。我的两种解决方案无法扩展。

背景

拥有一个包含 170 亿条记录的 BigQuery 表。每条记录代表一个设备 ping。每条记录都包含一个时间戳、一个用于识别用户的 ID 以及接收 ping 的位置名称。

获取此数据表,按 ID 对其进行分区并按时间戳排序。然后你有一组按时间顺序排列的 ping。用户可能对位置 A 有 1 个 ping,然后对位置 B 有 7 个 ping,然后对位置 C 有 2 个 ping,再对 A 有 2 个。

ID        timestamp             Location
ABC123    2017-10-12 10:20:37   A
ABC123    2017-10-12 11:15:21   B
ABC123    2017-10-12 11:21:47   B
ABC123    2017-10-12 11:25:05   B
ABC123    2017-10-12 11:32:12   B
ABC123    2017-10-12 11:36:24   B
ABC123    2017-10-12 11:47:13   B
ABC123    2017-10-12 11:59:08   B
ABC123    2017-10-12 12:04:42   C
ABC123    2017-10-12 17:04:52   C
ABC123    2017-10-12 19:15:37   A
ABC123    2017-10-12 19:18:37   A

我想做的是拿这张桌子并制作一个新的桌子,每次“旅行”一行。其中一次行程是一组连续的 ping,具有“first_ping”和“last_ping”列。如果行程包含 1 个 ping,则该时间戳既是第一个 ping,也是最后一个 ping。

ID        first_ping            last_ping             Location
ABC123    2017-10-12 10:20:37   2017-10-12 10:20:37   A
ABC123    2017-10-12 11:15:21   2017-10-12 11:59:08   B
ABC123    2017-10-12 12:04:42   2017-10-12 17:04:52   C
ABC123    2017-10-12 19:15:37   2017-10-12 19:18:37   A

解决方案的尝试

Python

我从来没有使用过这么大的数据,而且我一直使用 Python。所以我第一次尝试的解决方案是一个 Python 脚本:

  1. 查询 BQ 以获取一个 ID 的所有数据
  2. 按时间戳对数据进行排序
  3. 使用“位置”上的 diff 函数来确定它何时发生变化
  4. 使用 cumsum() 为每组具有相同值的 ping 标记所有项目。
  5. 在 cumsum() 上使用 df.groupby() 来获取每条记录的一行,并使用 first() 和 last() 来获取 first_ping 和 last_ping 值。

这个解决方案产生了我需要的输出,但对于 170 亿条记录和 6900 万个唯一 ID 来说是不可行的。每个 ID 大约需要 10 秒,也就是大约 19 万小时的运行时间。

SQL

WITH visitWithIsChange AS 
(select
   *,
   LAG(location,1,'') OVER (PARTITION BY user_id ORDER BY timestamp) previous,
    CASE 
     WHEN (LAG(location,1,'') 
           OVER (PARTITION BY user_id ORDER BY timestamp)) = location
           THEN 0
           ELSE 1
     END ischange
 FROM `ping_table` ORDER BY user_id, timestamp),
 visitsWithcumsum AS (
   SELECT 
      t1.*,
      SUM(t2.ischange) AS cumulativeSum 
   FROM visitWithIsChange t1
        INNER JOIN
             visitWithIsChange t2
               ON
                 t1.local_timestamp >=t2.local_timestamp
                 AND
                 t1.user_id=t2.user_id
   GROUP BY 
     t1.local_timestamp,
     t1.user_id,
     t1.chain_id,
     t1.previous,
     t1.isChange
   ORDER BY user_id, timestamp
)
SELECT 
  MIN(timestamp) AS first_ping,
  MAX(local_timestamp) AS last_ping,
  user_id,
  chain_id,
FROM visitsWithcumsum
GROUP BY
  user_id,
  cumulativeSum,
  chain_id,
ORDER BY user_id, first_ping

我知道 SQL 语句的问题是在分区之外使用 ORDER BY。每次对超过几十万行调用 ORDER BY 时,BigQuery 都会引发资源超出错误。我的理解是,这是因为 ORDER BY 语句要求所有被排序的数据都驻留在一个节点上。

问题

  1. 是否可以使用上述 SQL 方法,但将所有逻辑嵌套在一个分区语句中?基本上,按 user_id 分区并在该分区内执行所有排序、cumsum 等?
  2. 有没有更好的方法来解决这个问题?

我感谢任何和所有输入。我完全不知道解决这个问题的最佳方法,并且感觉超出了我的深度。

【问题讨论】:

    标签: python sql google-bigquery


    【解决方案1】:

    cumulativeSum 应该使用累积和而不是非等连接来计算:

    WITH visitWithIsChange AS 
    (select
       *,
        CASE 
         WHEN (LAG(location,1,'') 
               OVER (PARTITION BY user_id ORDER BY timestamp)) = location
               THEN 0
               ELSE 1
         END ischange
     FROM `ping_table`
     -- I don't now about BigQuery, but why do you need this?
     --ORDER BY user_id, timestamp
     ),
     visitsWithcumsum AS (
       SELECT 
          *,
          SUM(ischange)
          OVER (PARTITION BY user_id
                ORDER BY timestamp
                ROWS UNBOUNDED PREDECING) AS cumulativeSum 
       FROM visitWithIsChange  
    )
    SELECT 
      MIN(timestamp) AS first_ping,
      MAX(local_timestamp) AS last_ping,
      user_id,
      chain_id,
    FROM visitsWithcumsum
    GROUP BY
      user_id,
      cumulativeSum,
      chain_id,
    ORDER BY user_id, first_ping
    

    【讨论】:

    • 嘿@dnoeth,非常感谢您的回答。我能够在离开办公室之前启动查询,并且运行成功。没有资源超出错误。您注释掉的 ORDER BY 确实是无关的,而不是由于 BigQuery。这是我第一次使用 SQL(也是第一次使用大数据)。所以我犯了错误,但试图了解更多。我只需要明天早上到办公室时对这个查询的结果进行 QA。
    【解决方案2】:

    尝试以下版本(BigQuery 标准 SQL)

    #standardSQL
    SELECT 
      id, 
      MIN(timestamp) AS first_ping, 
      MAX(timestamp) AS last_ping, 
      ANY_VALUE(location) AS location
    FROM (
      SELECT id, timestamp, location,
        COUNTIF(flag) OVER(PARTITION BY id ORDER BY timestamp) grp
      FROM (
        SELECT *, 
          location != LAG(location) OVER(PARTITION BY id ORDER BY timestamp) flag
        FROM `project.dataset.ping_table`
      )
    )
    GROUP BY id, grp
    

    【讨论】:

    • 嘿@Mikhail-Berlyant,谢谢你的回答!早上到电脑前我会试一试。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-09-29
    • 2021-08-02
    • 1970-01-01
    • 2021-02-02
    • 1970-01-01
    • 1970-01-01
    • 2023-01-25
    相关资源
    最近更新 更多