【发布时间】: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 脚本:
- 查询 BQ 以获取一个 ID 的所有数据
- 按时间戳对数据进行排序
- 使用“位置”上的 diff 函数来确定它何时发生变化
- 使用 cumsum() 为每组具有相同值的 ping 标记所有项目。
- 在 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 语句要求所有被排序的数据都驻留在一个节点上。
问题
- 是否可以使用上述 SQL 方法,但将所有逻辑嵌套在一个分区语句中?基本上,按 user_id 分区并在该分区内执行所有排序、cumsum 等?
- 有没有更好的方法来解决这个问题?
我感谢任何和所有输入。我完全不知道解决这个问题的最佳方法,并且感觉超出了我的深度。
【问题讨论】:
标签: python sql google-bigquery