1、設定空閑狀態保留時間,Flink SQL 可以指定空閑狀態(即未更新的狀態)被保留的最小時間,當狀態中某個 key對應的狀態未更新的時間達到閾值時,該條狀態被自動清理:
#引數指定
configuration.setString("table.exec.state.ttl", "1 h");
2、開啟 MiniBatch
MiniBatch Aggregation,思路是記憶體快取 batch 資料再進行聚合,減少狀態訪問次數, 從而提升吞吐并減少資料的輸出量,MiniBatch 主要依靠在每個 Task 上注冊的 Timer 執行緒 來觸發微批,需要消耗一定的執行緒調度性能, MiniBatch 默認關閉,開啟方
如下:
// 設定引數,開啟 miniBatch
configuration.setString("table.exec.mini-batch.enabled", "true");
// 批量輸出的間隔時間
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
// 防止 OOM 設定每個批次最多快取資料的條數,可以設為 2 萬條
configuration.setString("table.exec.mini-batch.size", "20000");
適用場景:
??微批處理通過增加延遲換取高吞吐,如果有超低延遲的要求,不建議開啟微批處理,通 常對于聚合的場景,微批處理可以顯著的提升系統性能,建議開啟, 1.12 之前的版本有 bug,開啟 miniBatch,不會清理過期狀態,也就是說如果設定 狀態的 TTL,法清理過期狀態,1.12 版本才修復這個問題,
3、開啟 LocalGlobal
3.1 原理概述 Local Global Aggregation,思路是聚合操作拆分為兩階段, Local 階段預聚合減少資料條數,Global 解決全域聚合,即 MapReduce 模型中的Combine+Reduce 處理模式,本質上能夠靠 LocalAgg 的聚合篩除部分傾斜資料,從而降低 GlobalAgg 的熱點,提升性能,
??第一階段在上游節點本地攢一批資料 進行聚合(localAgg),并輸出這次微批的增量值(Accumulator),
??第二階段再將收到 的 Accumulator 合并(Merge),得到最終的結果(GlobalAgg),
LocalGlobal 開啟方式:
??1.LocalGlobal 優化需要先開啟 MiniBatch,依賴于MiniBatch 的引數,
??2.table.optimizer.agg-phase-strategy: 聚合策略,默認 AUTO,支持引數 AUTO、 TWO_PHASE(使用 LocalGlobal 兩階段聚合)、ONE_PHASE(僅使用 Global 一階段聚合),
// 設定引數:
// 開啟 miniBatch
configuration.setString("table.exec.mini-batch.enabled", "true");
// 批量輸出的間隔時間
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
// 防止 OOM 設定每個批次最多快取資料的條數,可以設為 2 萬條
configuration.setString("table.exec.mini-batch.size", "20000");
// 開啟 LocalGlobal
configuration.setString("table.optimizer.agg-phase-strategy", "TWO_PHASE");
注意事項:
??1.需要先開啟 MiniBatch
??2.開啟 LocalGlobal 需要 UDAF 實作 Merge 方法
4、開啟 Split Distinct(Split Distinct Aggregation,思路是針對 count distinct 場景, 對分組 key 先分桶預聚合, 再對分桶結果全域聚合)LocalGlobal 優化針對普通聚合(例如 SUM、COUNT、MAX、MIN 和 AVG)有較好的效果,對于 DISTINCT 的聚合(如 COUNT DISTINCT)收效不明顯,因為 COUNT DISTINCT 在 Local 聚合時,對于 DISTINCT KEY 的去重率不高,導致在 Global 節點仍然存在熱點,
4.1 原理概述 之前為了解決 COUNT DISTINCT 的熱點問題,通常需要手動改寫為兩層聚合(增加按 Distinct Key 取模的打散層),從 Flink1.9.0 版本開始 , 提供了 COUNT DISTINCT 自動打散功能 , 通過 HASH_CODE(distinct_key) %BUCKET_NUM 打散,不需要手動重寫,Split Distinct 和 LocalGlobal 的原理對比參見下圖,
--Distinct 舉例:
SELECT a, COUNT(DISTINCT b)
FROM TB?
GROUP BY a
--手動打散舉例:
SELECT a, SUM(cnt)
FROM (
SELECT a, COUNT(DISTINCT b) as cnt
FROM TB
GROUP BY a, MOD(HASH_CODE(b), 1024)
) GROUP BY a
Split Distinct 開啟方式默認不開啟,使用引數顯式開啟:
??table.optimizer.distinct-agg.split.enabled: true,默認 false,
??table.optimizer.distinct-agg.split.bucket-num: Split Distinct 優化在第一層聚合中,被打散的 bucket 數目,默認 1024,
// 設定引數:(要結合 minibatch 一起使用)
// 開啟 Split Distinct
configuration.setString("table.optimizer.distinct-agg.split.enabled", "true");
// 第一層打散的 bucket 數目
configuration.setString("table.optimizer.distinct-agg.split.bucket-num", "1024");
注意事項:
??1.目前不能在包含 UDAF 的 Flink SQL 中使用 Split Distinct 優化方法,
??2.拆分出來的兩個 GROUP 聚合還可參與 LocalGlobal 優化,
??3.該功能在 Flink1.9.0 版本及以上版本才支持,
5.多維 DISTINCT 使用 Filter
5.1 原理概述 在某些場景下,可能需要從不同維度來統計 count(distinct)的結果(比如統計 uv、app 端的 uv、web 端的 uv),可能會使用如下 CASE WHEN 語法,
SELECTB
a,
COUNT(DISTINCT b) AS total_b,
COUNT(DISTINCT ** CASE WHEN ** c IN ('A', 'B') THEN b ELSE NULL END) AS AB_b,
COUNT(DISTINCT ** CASE WHEN ** c IN ('C', 'D') THEN b ELSE NULL END) AS CD_b
FROM TB
GROUP BY a
在這種情況下,建議使用 FILTER 語法, 目前的 Flink SQL 優化器可以識別同一唯一鍵 上的不同 FILTER 引數,如,在上面的示例中,三個 COUNT DISTINCT 都作用在 b 列上, 此時,經過優化器識別后,Flink 可以只使用一個共享狀態實體,而不是三個狀態實體,可 減少狀態的大小和對狀態的訪問,將上邊的 CASE WHEN 替換成 FILTER 后,如下所示:
SELECTB
a,
COUNT(DISTINCT b) AS total_b,
COUNT(DISTINCT b) FILTER (WHERE c IN ('A', 'B')) AS AB_b,
COUNT(DISTINCT b) FILTER (WHERE c IN ('C', 'D')) AS CD_b
FROM TB
GROUP BY a
6、總結以上的調優引數,代碼如下:
// 設定引數:
// 開啟 miniBatch
configuration.setString("table.exec.mini-batch.enabled", "true");
// 批量輸出的間隔時間
configuration.setString("table.exec.mini-batch.allow-latency", "5 s");
// 防止 OOM 設定每個批次最多快取資料的條數,可以設為 2 萬條
configuration.setString("table.exec.mini-batch.size", "20000");
// 開啟 LocalGlobal
configuration.setString("table.optimizer.agg-phase-strategy", "TWO_PHASE");
// 開啟 Split Distinct
configuration.setString("table.optimizer.distinct-agg.split.enabled", "true");
// 第一層打散的 bucket 數目
configuration.setString("table.optimizer.distinct-agg.split.bucket-num", "1024");
// 指定時區
configuration.setString("table.local-time-zone", "Asia/Shanghai");
附錄:FlinkSQL 官網配置引數: https://nightlies.apache.org/flink/flink-docs-release-1.13/docs/dev/table/config/
腦子是空的不要緊,主要是不要進水······轉載請註明出處,本文鏈接:https://www.uj5u.com/shujuku/547977.html
標籤:大數據
上一篇:交易系統之資料庫弱依賴解決方案
下一篇:初識Kafka
