如何在BigQuery中并行处理单国家下多Ngram以优化脚本执行效率
Optimize BigQuery Script: Process All Ngrams Per Country Instead of Per Combination
Got it, let's tackle this performance issue head-on. The core problem with your original script is that it loops over every single country+ngram pair, which means it's scanning dataset_1.table_22 hundreds/thousands of times—way too inefficient. Here's how to refactor it to process all ngrams per country in one go, cutting down on repeated table scans and execution time drastically:
DECLARE result ARRAY<STRUCT<country STRING, Metric STRING, consecutiveCount INT64, startDate TIMESTAMP, endDate TIMESTAMP>> DEFAULT []; -- Loop over distinct countries instead of country+ngram pairs FOR country_record IN (SELECT DISTINCT country FROM `dataset_1.table_11`) DO DECLARE current_country STRING DEFAULT country_record.country; -- Get all ngrams associated with the current country DECLARE country_ngrams ARRAY<STRING> DEFAULT ARRAY( SELECT ngram FROM `dataset_1.table_11` WHERE country = current_country ); SET result = ARRAY_CONCAT(result, ARRAY( SELECT STRUCT(current_country AS country, Metric, consecutiveCount, startDate, endDate) FROM ( WITH term_t AS ( -- Single scan of table_22 per country, join with all ngrams for the country SELECT t.trending_at, current_country AS country, ng AS Metric, COUNTIF(CONTAINS_SUBSTR(LOWER(t.title), ng)) AS Value FROM `dataset_1.table_22` t CROSS JOIN UNNEST(country_ngrams) AS ng WHERE t.country = current_country GROUP BY t.trending_at, current_country, ng ), -- Reuse your original consecutive count logic unchanged StockRow AS ( SELECT Metric, Value, trending_at, ROW_NUMBER() OVER(PARTITION BY Metric ORDER BY trending_at) rn FROM term_t ), RunGroup AS ( SELECT Base.Metric, Base.trending_at, MAX(Restart.rn) OVER(PARTITION BY Base.Metric ORDER BY Base.trending_at) groupingId FROM StockRow Base LEFT JOIN StockRow Restart ON Restart.Metric = Base.Metric AND Restart.rn = Base.rn - 1 AND Restart.Value >= Base.Value ), F AS ( SELECT Metric, COUNT(*) AS consecutiveCount, MIN(trending_at) AS startDate, MAX(trending_at) AS endDate FROM RunGroup GROUP BY Metric, groupingId HAVING COUNT(*) >= 3 ORDER BY Metric, startDate ) SELECT * FROM F ) )); END FOR; INSERT INTO `dataset_1.table_33` SELECT * FROM UNNEST(result);
Key Optimizations Breakdown
- Loop at the Country Level: Instead of iterating every
country+ngrampair, we loop once per unique country. This cuts the number of iterations from potentially thousands to just the number of countries in your data. - Single Table Scan Per Country: Using
CROSS JOIN UNNEST(country_ngrams)lets us scantable_22once per country to calculate values for all ngrams in that country. Your original approach would scantable_22once per ngram—this is the biggest performance and cost-saving win. - Unchanged Core Logic: The consecutive counting logic (
StockRow,RunGroup,F) stays exactly the same because we're still grouping byMetric(each ngram) and processing time-series data the same way.
Why This Works for Your Sample Data
Running this script on your sample data:
- For the US, it processes both "lorem" and "ipsum" in one scan of
table_22. - The consecutive count logic correctly identifies that "lorem" has 3 consecutive days of increasing counts (1 on 2021-03-04, 2 on 2021-03-05, 3 on 2021-03-06), so it's included in the final result.
- Other ngrams don't meet the
COUNT(*) >=3criteria, so they're excluded—matching your expected output perfectly.
内容的提问来源于stack exchange,提问作者Ammar Alyousfi
相关产品推荐
相关产品推荐

