BigQuery优化方案:计算各时间点有库存卖家的产品最低价格
基于BigQuery CDC数据计算各时间点产品最低价格的优化方案
我们有来自产品列表表变更捕获(CDC)的变更数据,包含时间点(dt)、产品ID(product_id)、卖家(seller)、库存(stock)、价格(price)字段。需求是计算每个时间点有库存卖家中的产品最低价格及对应卖家,当前使用交叉连接的方案性能低下、资源占用高,需要更高效的实现方式。
示例数据
WITH a AS ( SELECT 1 AS dt, 'p' as product_id, 'A' AS seller, 10 as stock, 100 as price UNION ALL SELECT 2 AS dt, 'p' as product_id, 'B' AS seller, 10 as stock, 120 as price UNION ALL SELECT 3 AS dt, 'p' as product_id, 'C' AS seller, 10 as stock, 150 as price UNION ALL SELECT 4 AS dt, 'p' as product_id, 'D' AS seller, 10 as stock, 300 as price UNION ALL SELECT 5 AS dt, 'p' as product_id, 'E' AS seller, 10 as stock, 400 as price UNION ALL SELECT 6 AS dt, 'p' as product_id, 'F' AS seller, 10 as stock, 500 as price UNION ALL SELECT 7 AS dt, 'p' as product_id, 'G' AS seller, 10 as stock, 600 as price UNION ALL SELECT 8 AS dt, 'p' as product_id, 'A' AS seller, 0 as stock, 100 as price UNION ALL SELECT 9 AS dt, 'p' as product_id, 'B' AS seller, 10 as stock, 110 as price UNION ALL SELECT 10 AS dt, 'p' as product_id, 'B' AS seller, 10 as stock, 190 as price UNION ALL SELECT 11 AS dt, 'p' as product_id, 'G' AS seller, 10 as stock, 800 as price UNION ALL SELECT 12 AS dt, 'p' as product_id, 'G' AS seller, 10 as stock, 100 as price ) SELECT * FROM a
预期输出
| dt | product_id | minimum_price | seller_with_minimum_price |
|---|---|---|---|
| 1 | p | 100 | A |
| 2 | p | 100 | A |
| 3 | p | 100 | A |
| 4 | p | 100 | A |
| 5 | p | 100 | A |
| 6 | p | 100 | A |
| 7 | p | 100 | A |
| 8 | p | 120 | B |
| 9 | p | 110 | B |
| 10 | p | 150 | C |
| 11 | p | 150 | C |
| 12 | p | 100 | G |
优化实现方案
核心思路是先为每个卖家生成其状态(库存、价格)的生效时间段,再将所有时间点与对应生效的卖家状态关联,最后聚合计算最低价格,彻底避免交叉连接的高开销。
WITH -- 1. 为每个卖家的变更记录标记下一次变更时间,确定当前状态的生效区间 seller_state_changes AS ( SELECT product_id, seller, dt, stock, price, -- 获取该卖家下一次变更的时间,无后续变更则设为极大值 LEAD(dt) OVER (PARTITION BY product_id, seller ORDER BY dt) AS next_dt FROM a ), -- 2. 提取所有需要计算的时间点 all_dts AS ( SELECT DISTINCT dt FROM a ), -- 3. 关联时间点与卖家的有效状态:仅保留时间点处于状态生效区间且库存>0的记录 active_sellers_at_dt AS ( SELECT ad.dt, ssc.product_id, ssc.seller, ssc.price FROM all_dts ad JOIN seller_state_changes ssc ON ad.dt >= ssc.dt AND (ad.dt < ssc.next_dt OR ssc.next_dt IS NULL) WHERE ssc.stock > 0 ), -- 4. 对每个时间点的有效卖家价格排序,标记最低价格的记录 ranked_prices AS ( SELECT dt, product_id, price, seller, ROW_NUMBER() OVER (PARTITION BY dt, product_id ORDER BY price ASC) AS rn FROM active_sellers_at_dt ) -- 5. 筛选出每个时间点的最低价格及对应卖家 SELECT dt, product_id, price AS minimum_price, seller AS seller_with_minimum_price FROM ranked_prices WHERE rn = 1 ORDER BY dt
方案说明
- seller_state_changes:通过
LEAD窗口函数明确每个卖家状态的生效时间段,避免重复计算历史状态。 - all_dts:提取所有CDC记录的时间点,确保覆盖所有需要计算的时刻。
- active_sellers_at_dt:通过区间关联替代交叉连接,精准匹配每个时间点的有效卖家状态,大幅减少无效数据量。
- ranked_prices:用
ROW_NUMBER对每个时间点的价格排序,快速定位最低价格记录。 - 最后筛选排名第一的记录,得到目标结果。
该方案利用BigQuery的窗口函数和高效关联能力,在大数据量场景下性能远优于交叉连接方案。
内容的提问来源于stack exchange,提问作者armamut
相关产品推荐
相关产品推荐

