天天看點

Flink的重新開機政策Flink的重新開機政策概覽執行個體重新開機政策

Flink的重新開機政策

Flink支援不同的重新開機政策,這些重新開機政策控制着job失敗後如何重新開機。叢集可以通過預設的重新開機政策來重新開機,這個預設的重新開機政策通常在未指定重新開機政策的情況下使用,而如果Job送出的時候指定了重新開機政策,這個重新開機政策就會覆寫掉叢集的預設重新開機政策。

概覽

叢集在啟動時會伴随一個預設的重新開機政策,在沒有定義具體重新開機政策時,會使用該預設重新開機政策,如果在工作送出時指定了一個重新開機政策,那麼該政策會覆寫叢集的預設政策。

預設的重新開機政策可以通過Flink的配置檔案

flink-conf.yaml

指定,配置參數

restart-strategy

定義了哪個政策被使用。

常用的重新開機政策:

  • 固定間隔(Fixed delay)
  • 失敗率(Failure rate)
  • 無重新開機(No restart)
  1. 如果checkpoint未啟動,就會采用no restart政策。
  2. 如果啟動了checkpoint機制,但是未指定重新開機政策的話,就會采用fixed-delay政策,重試Integer.MAX_VALUE次。

請參考下面的可用重新開機政策來了解哪些值是支援的。

每個重新開機政策都有自己的參數來控制它的行為,這些值也可以在配置檔案中設定,每個重新開機政策的描述都包含着各自的配置值資訊。

重新開機政策 重新開機政策值
Fixed delay fixed-delay
Failure rate failure-rate
No restart None

除了定義一個預設的重新開機政策之外,你還可以為每一個Job指定它自己的重新開機政策,這個重新開機政策可以在ExecutionEnvironment中調用setRestartStrategy()方法來程式化地調用,主意這種方式同樣适用于StreamExecutionEnvironment。

執行個體

下面的例子展示了我們如何為我們的Job設定一個固定延遲重新開機政策,一旦有失敗,系統就會嘗試每10秒重新開機一次,重新開機3次。

java方式

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
  3, // 嘗試重新開機次數
  Time.of(10, TimeUnit.SECONDS) // 延遲時間間隔
));
           

scala方式

val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
  3, // 重新開機次數
  Time.of(10, TimeUnit.SECONDS) // 延遲時間間隔
))
           

重新開機政策

下面部分描述了重新開機政策特定的配置項

固定延遲重新開機政策(Fixed Delay Restart Strategy)

固定延遲重新開機政策會嘗試一個給定的次數來重新開機Job,如果超過了最大的重新開機次數,Job最終将失敗。在連續的兩次重新開機嘗試之間,重新開機政策會等待一個固定的時間。

重新開機政策可以配置flink-conf.yaml的下面配置參數來啟用,作為預設的重新開機政策:

restart-strategy: fixed-delay
           
配置參數 描述 預設值
restart-strategy.fixed-delay.attempts 在Job最終宣告失敗之前,Flink嘗試執行的次數 1,如果啟用checkpoint的話是Integer.MAX_VALUE
restart-strategy.fixed-delay.delay 延遲重新開機意味着一個執行失敗之後,并不會立即重新開機,而是要等待一段時間。 akka.ask.timeout,如果啟用checkpoint的話是1s

例子:

restart-strategy.fixed-delay.attempts: 3
restart-strategy.fixed-delay.delay: 10 s
           

固定延遲重新開機也可以在程式中設定:

java:

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
  3, // 重新開機次數
  Time.of(10, TimeUnit.SECONDS) // 重新開機時間間隔
));
           

scala:

val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.fixedDelayRestart(
  3, // 重新開機次數
  Time.of(10, TimeUnit.SECONDS) // 重新開機時間間隔
))
           

失敗率重新開機政策(Failure rate)

失敗率重新開機政策在Job失敗後會重新開機,但是超過失敗率後,Job會最終被認定失敗。在兩個連續的重新開機嘗試之間,重新開機政策會等待一個固定的時間。

失敗率重新開機政策可以在flink-conf.yaml中設定下面的配置參數來啟用:

restart-strategy:failure-rate
           
配置參數 描述 預設值
restart-strategy.failure-rate.max-failures-per-interval 在一個Job認定為失敗之前,最大的重新開機次數 1
restart-strategy.failure-rate.failure-rate-interval 計算失敗率的時間間隔 1分鐘
restart-strategy.failure-rate.delay 兩次連續重新開機嘗試之間的時間間隔 akka.ask.timeout

例子:

restart-strategy.failure-rate.max-failures-per-interval: 3
restart-strategy.failure-rate.failure-rate-interval: 5 min
restart-strategy.failure-rate.delay: 10 s
           

失敗率重新開機政策也可以在程式中設定:

Java代碼:

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
env.setRestartStrategy(RestartStrategies.failureRateRestart(
  3, // 每個測量時間間隔最大失敗次數
  Time.of(5, TimeUnit.MINUTES), //失敗率測量的時間間隔
  Time.of(10, TimeUnit.SECONDS) // 兩次連續重新開機嘗試的時間間隔
));
           

Scala代碼::

val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.failureRateRestart(
  3, // 每個測量時間間隔最大失敗次數
  Time.of(5, TimeUnit.MINUTES), //失敗率測量的時間間隔
  Time.of(10, TimeUnit.SECONDS) // 兩次連續重新開機嘗試的時間間隔
))
           

無重新開機政策

Job直接失敗,不會嘗試進行重新開機。如果沒有啟動checkpoint,則預設情況下就是無重新開機

restart-strategy: none
           

無重新開機政策也可以在程式中設定

Java代碼:

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();
env.setRestartStrategy(RestartStrategies.noRestart());
           

Scala代碼:

val env = ExecutionEnvironment.getExecutionEnvironment()
env.setRestartStrategy(RestartStrategies.noRestart())
           

繼續閱讀