BigQuery使用LAST_VALUE() OVER()遇资源超限及前向填充需求
解决BigQuery大表前向填充时内存溢出的问题
你的问题核心是全局窗口排序导致内存耗尽——原查询中OVER(ORDER BY datetime)没有指定分区,BigQuery需要对整个1.126亿行数据进行全局排序,这显然超出了内存限制。好在你的表已经按天分区且以datetime聚类,我们可以利用这个特性大幅优化查询。
优化方案:利用分区缩小窗口范围
因为你的表是按天分区的,前向填充的需求通常可以限制在单个日期内(如果需要跨天填充,后面会补充方案)。通过在窗口函数中添加PARTITION BY DATE(datetime),让每个窗口只处理一天的数据,内存压力会骤降。同时明确窗口框架,确保LAST_VALUE取到当前行之前最近的非空值。
优化后的查询代码
INSERT INTO project.dataset.table (datetime, col1, col2, col3, col4, col5, col6, col7) WITH filled_data AS ( SELECT datetime, -- 对每一列应用带分区的LAST_VALUE前向填充 LAST_VALUE(col1 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col1, LAST_VALUE(col2 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col2, LAST_VALUE(col3 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col3, LAST_VALUE(col4 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col4, LAST_VALUE(col5 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col5, LAST_VALUE(col6 IGNORE NULLS) OVER ( PARTITION BY DATE(datetime) ORDER BY datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW ) AS col6, col7 -- col7无空值,直接保留 FROM project.dataset.origin_table ) SELECT * FROM filled_data
关键优化点说明
PARTITION BY DATE(datetime):将窗口限制在单个日期的分区内,每个分区的数据量远小于全局表,避免了全局排序的内存开销。同时因为表是按datetime聚类的,BigQuery可以直接利用聚类后的有序数据,不需要额外排序,进一步提升性能。- 明确窗口框架:
ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW确保窗口范围是从当前分区的第一条数据到当前行,这样LAST_VALUE(IGNORE NULLS)会准确取到当前行之前最近的非空值,完美实现前向填充的需求。
可选:跨天前向填充方案
如果需要处理跨天的空值(比如某天的第一条数据为空,需要用上一天最后一个非空值填充),可以先计算每个日期的各列最后非空值,再关联到下一天的数据中:
WITH daily_last_values AS ( SELECT DATE(datetime) AS date, LAST_VALUE(col1 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col1, LAST_VALUE(col2 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col2, LAST_VALUE(col3 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col3, LAST_VALUE(col4 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col4, LAST_VALUE(col5 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col5, LAST_VALUE(col6 IGNORE NULLS) OVER (PARTITION BY DATE(datetime) ORDER BY datetime) AS last_col6, ROW_NUMBER() OVER (PARTITION BY DATE(datetime) ORDER BY datetime DESC) AS rn FROM project.dataset.origin_table QUALIFY rn = 1 -- 取每个日期的最后一条数据(即该日期的最后非空值) ), filled_data AS ( SELECT o.datetime, -- 先尝试当天的前向填充,为空则用上一天的最后非空值 COALESCE( LAST_VALUE(o.col1 IGNORE NULLS) OVER (PARTITION BY DATE(o.datetime) ORDER BY o.datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW), dl.last_col1 ) AS col1, COALESCE( LAST_VALUE(o.col2 IGNORE NULLS) OVER (PARTITION BY DATE(o.datetime) ORDER BY o.datetime ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW), dl.last_col2 ) AS col2, -- 其他列同理 o.col7 FROM project.dataset.origin_table o LEFT JOIN daily_last_values dl ON DATE(o.datetime) = DATE_ADD(dl.date, INTERVAL 1 DAY) ) SELECT * FROM filled_data
内容的提问来源于stack exchange,提问作者Pedro Pablo Severin Honorato
相关产品推荐
相关产品推荐

