Skip to content

JobとCronJob

実装: orchestration/job/ / 実行: go test ./orchestration/job/

この編で扱ってきたワークロードは、どれも動き続けるものだった。だが終わる仕事もある。バッチ集計、移行、バックアップ。終わることが正常な結末なので、扱いが逆になる。数えるのは ready な数でなく成功した数で、失敗という結末が加わるぶん、何回まで試すかを決めておく必要が出る。周期実行には、前の実行が終わらないときどうするかという固有の問題がある。

この章で作るもの

この編で扱ってきたワークロードは、どれも動き続けるものだった。Pod は起動して、リクエストを受けて、止められるまで生きている。ヘルスチェックが見ていたのも「まだ受けられるか」で、終わることは想定していなかった。調整ループが「3個であれ」と言うのも、3個が居続けることを意味していた。

だが終わる仕事もある。夜間のバッチ集計、データの移行、バックアップ。これらは走って、終わって、消える。終わることが正常な結末なので、扱いが逆になる。動き続けるものは「止まったら異常」だが、こちらは「終わったら正常」で、「終わらないほうが異常」になる。

数え方も変わる。ready な数でなく、成功した数を数える。そして失敗という結末が加わるので、何回まで試すかを決めておく必要が出てくる。決めておかないと、直らない失敗を永久に再試行し続ける。

  動き続けるもの                終わるもの

  数えるもの: ready な数        数えるもの: 成功した数
  目標に達したら: 保ち続ける     目標に達したら: 終わる
  Pod が消えたら: 作り直す       Pod が成功したら: それでよい
  失敗は: 異常                  失敗は: 想定内(ただし上限が要る)

  completions 6 / parallelism 2 なら
   ■■ ■■ ■■   ← 2つずつ、3周期で6個
   量(いくつ要るか)と速さ(同時にいくつ)は別の軸
動き続けるものと終わるもの。何を数えるかと、目標に達したあとどうするかが逆になる

順に見ていく。

  1. 成功した数を数える: 目標に達したら終わり。保ち続けるのではない
  2. 量と速さは別の軸: いくつ要るか(completions)と、同時にいくつ走らせるか(parallelism)
  3. 諦める点を決める: 決めておかないと、直らない失敗を永久に試し続ける

① 成功を数えて、終わる

まず Job 本体を作る:

go

// Phase は Job の結末。
type Phase int

const (
	Running  Phase = iota // まだ走っている
	Complete              // 必要な数だけ成功した
	Failed                // 再試行の上限を超えた
)

func (p Phase) String() string {
	switch p {
	case Running:
		return "Running"
	case Complete:
		return "Complete"
	case Failed:
		return "Failed"
	}
	return "Unknown"
}

// Config は Job の設定。3つの数が別々の軸を決める。
type Config struct {
	// Completions は何個成功すれば終わりか。仕事の量。
	Completions int
	// Parallelism は同時に何個走らせるか。速さ。
	Parallelism int
	// BackoffLimit は失敗を何回まで許すか。これを超えると Job ごと失敗にする。
	BackoffLimit int
	// FailFirst は最初の何回の試行が失敗するか(台本。乱数を避けるため)。
	FailFirst int
}

// Job は「必要な数だけ成功したら終わる」ワークロード。
type Job struct {
	cfg    Config
	active int // 今走っている数

	Succeeded int
	Failed    int
	Attempts  int
	Phase     Phase
	Log       []string
}

// New は設定 cfg の Job を作る。
func New(cfg Config) *Job {
	if cfg.Parallelism < 1 {
		cfg.Parallelism = 1
	}
	if cfg.Completions < 1 {
		cfg.Completions = 1
	}
	return &Job{cfg: cfg}
}

// Active は今走っている数を返す。
func (j *Job) Active() int { return j.active }

// Done は結末が決まったかを返す。
func (j *Job) Done() bool { return j.Phase != Running }

// Step は1周期進める。走っているものに結末をつけ、足りなければ起動する。
//
// 数えるのは ready な数ではなく、成功した数になる。動き続けるものとは
// 数え方が逆で、目標に達したら終わりであって、目標を保ち続けるのではない。
func (j *Job) Step() {
	if j.Done() {
		return
	}

	// ① 走っているものに結末がつく。台本の回数までは失敗する。
	for ; j.active > 0; j.active-- {
		j.Attempts++
		if j.Attempts <= j.cfg.FailFirst {
			j.Failed++
			j.logf("試行 " + itoa(j.Attempts) + " が失敗(通算 " + itoa(j.Failed) + " 回目)")
			continue
		}
		j.Succeeded++
		j.logf("試行 " + itoa(j.Attempts) + " が成功(通算 " + itoa(j.Succeeded) + " 個)")
	}

	// ② 失敗が上限を超えたら、Job ごと失敗にする。
	// 上限を決めておかないと、直らない失敗を永久に再試行し続ける。
	if j.Failed > j.cfg.BackoffLimit {
		j.Phase = Failed
		j.logf("失敗が上限 " + itoa(j.cfg.BackoffLimit) + " を超えた。Job を失敗として終える")
		return
	}

	// ③ 必要な数だけ成功したら終わり。
	if j.Succeeded >= j.cfg.Completions {
		j.Phase = Complete
		j.logf("必要な " + itoa(j.cfg.Completions) + " 個が成功した")
		return
	}

	// ④ 残りを、同時実行の上限まで起動する。
	remaining := j.cfg.Completions - j.Succeeded
	n := j.cfg.Parallelism
	if remaining < n {
		n = remaining // 要る数より多くは走らせない
	}
	j.active = n
}

// Run は結末がつくまで最大 max 周期まわす。
func (j *Job) Run(max int) {
	for i := 0; i < max && !j.Done(); i++ {
		j.Step()
	}
}

Step が数えているのは Succeeded で、Ready ではない。そして Succeeded >= Completions になったら Complete にして、そこで止まる。動き続けるものとの違いがここに出ている。あちらは目標の数を保ち続けるので、Pod が消えたら作り直す。こちらは目標に達したら、もう何も作らない。

CompletionsParallelism が別の設定であることにも意味がある。6個の仕事を同時2つで処理するなら3周期、同時3つなら2周期。仕事の量は変わらず、終わるまでの時間だけが変わる。テストで、並列を増やすと同じ量が短い周期で終わることを固定した。

残りより多く起動しないようにもしてある。残り1個のときに同時実行の設定が5でも、1個しか起動しない。余分に走らせても、成功が必要数を超えるだけで意味がない。テストで、残り分しか起動しないことを固定した。

② 諦める点を決める

失敗の扱いが、動き続けるものと最も違うところになる。

動き続ける Pod が落ちたら、調整ループは無条件で作り直す。何度落ちても作り直す。それでよい。目標は「3個居ること」で、落ちた回数は目標に関係しないからだ。だが Job は違う。目標は「6個成功すること」なので、失敗し続けるなら永久に達成しない。作り直し続ければ、永久に走り続けることになる。

BackoffLimit がその歯止めになる。失敗がこの回数を超えたら、Job ごと失敗として終える。テストで、上限の内なら再試行して最後には完了し、超えたら諦めることを固定した。

上限を0にすると、1回の失敗で諦める。大きくすると粘る。どちらが正しいということはなく、失敗の性質で決まる。一時的な障害なら再試行が効くが、入力データが壊れているなら何度試しても直らない。後者に大きな上限を設定すると、直らない失敗のために資源を使い続ける。実物の Kubernetes では、この再試行の間隔が指数的に延びるようになっていて、無駄を減らす工夫が入っている。

③ 前の実行が終わらないとき

周期実行には、固有の問題がある:

go

// Policy は前の実行がまだ終わっていないときの振る舞い。
type Policy int

const (
	// Allow は重ねて走らせる。処理が重複しても構わない仕事向け。
	Allow Policy = iota
	// Forbid はこの回を飛ばす。重複が許されない仕事向け。
	Forbid
	// Replace は前の実行を止めて置き換える。最新だけが要る仕事向け。
	Replace
)

func (p Policy) String() string {
	switch p {
	case Allow:
		return "Allow"
	case Forbid:
		return "Forbid"
	case Replace:
		return "Replace"
	}
	return "Unknown"
}

// CronConfig は周期実行の設定。
type CronConfig struct {
	// Every は何周期ごとに起動するか。
	Every int
	// Policy は前の実行が残っているときの扱い。
	Policy Policy
	// Job は起動する Job の設定。
	Job Config
}

// CronJob は決まった周期で Job を起動する。
type CronJob struct {
	cfg  CronConfig
	now  int
	runs []*Job

	Started  int // 起動した数
	Skipped  int // 飛ばした数
	Replaced int // 置き換えた数
	Log      []string
}

// NewCron は設定 cfg の周期実行を作る。
func NewCron(cfg CronConfig) *CronJob {
	if cfg.Every < 1 {
		cfg.Every = 1
	}
	return &CronJob{cfg: cfg}
}

// Runs はこれまでに起動した Job を起動順に返す。
func (c *CronJob) Runs() []*Job { return c.runs }

// Active はまだ終わっていない Job の数を返す。
func (c *CronJob) Active() int {
	n := 0
	for _, j := range c.runs {
		if !j.Done() {
			n++
		}
	}
	return n
}

// Tick は時刻を1つ進める。走っている Job を進め、起動の時刻なら方針に従う。
func (c *CronJob) Tick() {
	c.now++
	for _, j := range c.runs {
		j.Step()
	}
	if c.now%c.cfg.Every != 0 {
		return
	}

	// 起動の時刻。前の実行が残っているかで振る舞いが変わる。
	if c.Active() > 0 {
		switch c.cfg.Policy {
		case Forbid:
			c.Skipped++
			c.logf("前の実行が終わっていないので、この回は飛ばす")
			return
		case Replace:
			for _, j := range c.runs {
				if !j.Done() {
					j.Phase = Failed
					j.logf("次の実行に置き換えられた")
				}
			}
			c.Replaced++
			c.logf("前の実行を止めて置き換える")
		}
	}
	j := New(c.cfg.Job)
	c.runs = append(c.runs, j)
	c.Started++
	c.logf("実行を起動(" + itoa(c.Started) + " 回目)")
}

// Completed は成功して終わった Job の数を返す。
func (c *CronJob) Completed() int {
	n := 0
	for _, j := range c.runs {
		if j.Phase == Complete {
			n++
		}
	}
	return n
}

Tick が起動の時刻に達したとき、前の実行がまだ終わっていることがある。1時間ごとに走る集計が、データ量が増えて1時間半かかるようになった、というのはよくある。このとき何をするか。

Policy の3つが、その答えになる。Allow は重ねて走らせる。Forbid はこの回を飛ばす。Replace は前を止めて置き換える。テストで、同じ状況で3つの方針が別々の結果になることを固定した。

どれを選んでも何かを失う。重ねれば同じデータを二重に処理するかもしれない。飛ばせば、その回の処理は永久に行われない。置き換えれば、途中まで進んだ処理が捨てられる。仕事の性質で決めるしかない。集計をやり直しても結果が同じなら重ねてよいし、最新の状態だけが要るなら置き換えでよい。

Forbid について1つ注意が要る。飛ばされた回は、後で埋め合わされない。「1時間ごとに必ず走る」ことは保証されない。テストで、起動された数と飛ばされた数の合計が予定回数になることを固定した。予定どおりの回数だけ走ったかを確かめたいなら、走った記録のほうを数える必要がある。

動かす

下のデモは、Job と CronJob を切り替えて見る。Job 側では、並列を上げると終わるまでの周期が減り、失敗の設定を上げると諦める側に倒れる。CronJob 側では、1回の仕事が起動の周期より長い状況で、3つの方針が別々の帯を描く。

デモJobとCronJobFailed ・ 成功 1 / 失敗 3 ・ 3 周期
Job(1回の仕事)CronJob(周期実行)
completions(仕事の量)369
parallelism(同時に走らせる数)123
backoffLimit(諦めるまでの失敗数)025
最初の何回が失敗するか038
成功1 / 6
失敗3(上限 2)
失敗が上限 2 を超えたので Job ごと失敗にした。上限を決めていなければ、ここで永久に再試行し続けていた

completions は仕事の量、parallelism は速さで、別の軸になっている。並列を上げると同じ量が少ない周期で終わる。 backoffLimit は諦める点で、これを 0 にすると1回の失敗で Job ごと失敗になる。逆に決めておかないと、 直らない失敗を永久に試し続けることになる。「最初の何回が失敗するか」を 8 にすると、どの上限でも諦める側に倒れる。

設計の観点

  • 終わることが正常な結末: 動き続けるものと、目標の扱いが逆になる。同じ「3」でも、保つ3と達する3は別のもの
  • 諦める点は必ず決める: 決めないと、直らない失敗のために資源を使い続ける。何回試せば直る見込みがあるかは、失敗の性質で決まる
  • 重複を前提に設計する: 再試行がある以上、同じ処理が2回走ることは起こりうる。何度実行しても結果が変わらないように作っておくと、方針の選択が楽になる。メッセージキューの at-least-once と同じ話
  • 周期より長い仕事は方針が要る: どれを選んでも何かを失う。仕事が周期より長くなったこと自体を検知して、周期のほうを見直すのが本筋
  • 飛ばした回は埋まらない: Forbid は静かに実行を落とす。予定どおり走ったかは、別に数える必要がある
  • Operatorとの関係: 終わる仕事の管理も、宣言と差の判定として書かれている。「6個成功した状態であれ」という宣言に、この編の一貫した形が出ている

対照と実例

動き続けるワークロードJob
数えるものready な数成功した数
目標に達したら保ち続ける終わる
Pod が消えたら作り直す成功していれば何もしない
失敗したら無条件に作り直す上限まで再試行し、超えたら諦める
正常な終わり方無い(止められるまで)終わること

裏どり:

  • Job: completionsparallelism が別の設定で、backoffLimit が諦める点になる
  • 再試行の間隔: 実物は失敗のたびに待ち時間が指数的に延びる。上限に達するまで無駄に速く試さない
  • CronJob の concurrencyPolicy: Allow / Forbid / Replace の3つ。既定は Allow
  • 飛ばされた実行: Forbid で飛ばされた回は後から埋め合わされない。startingDeadlineSeconds を過ぎた回も同様。ただし1つ例外があって、Forbid で前の実行が終わったとき、startingDeadlineSeconds の中に収まっていれば実行が起きることがある

簡略化したこと

  • Pod を持たない: 実物は Job が Pod を作る。ここでは走っている数だけを数える
  • バックオフの間隔なし: 再試行のたびに待ち時間が延びる仕組みは入れていない
  • completionMode は既定のみ: 序数つきの Indexed は扱わない
  • 時刻はカレンダーでない: 実物は cron 式で書く。ここは何周期ごとという単純な形
  • 後片付けなし: 終わった Job を自動で消す仕組みは扱わない

参考資料

  • Jobs — completions / parallelism / backoffLimit の関係
  • CronJob — 3つの concurrencyPolicy と、飛ばされた実行の扱い
  • 実装: orchestration/job