滞留を見てレプリカ数を決める
実装:
orchestration/scaletrigger// 実行:go test ./orchestration/scaletrigger/
待ち行列の長さを目標にすると必要な数がその場で出る。だが実物では、そもそも何が取れるかが相手の作りで決まる。読んでも消えないログなら書いた位置と読んだ位置の差が取れ、読んだら消えるキューは残りしか取れない。しかも見えている数だけで判断すると、全部が処理中のときに0個と答えてしまう。そして分割数を超えたレプリカは、立てても働かない。
この章で作るもの
カスタム指標とKEDAで、待ち行列の長さを目標にすると ceil(全体の量 ÷ 1 個あたりの目標) で必要な数がその場で出ることを見た。KEDA の Kafka スケーラも SQS スケーラも、この式そのものを使っている。
だが式より手前に問題がある。そもそも「待ち行列の長さ」として何が取れるのかは、相手の作りで変わる。ここを取り違えると、式が正しくても答えが狂う。
この章では配信元を 2 種類作って、取れる情報の違いと、そこから生まれる 3 つの落とし穴を数える。
ログ型(Kafka など) キュー型(SQS など)
書いた位置 ─────┐ ┌ 見えている(取り出せる)
[■■■■■■■■■■] [■■■■■■]
└ 読んだ位置 └ 処理中(取り出し済み・未確定)
ラグ = 書いた − 読んだ 残り = 見えている + 処理中
読んでも消えないので 取り出すと見えなくなるので
「何件流れたか」が残る 「何件流れたか」は残らない順に見ていく。
- ラグが取れる形と取れない形がある: 読んでも消えないなら位置の差が取れる。消えるなら残りしか取れない
- 見えている数だけでは足りない: 全部が処理中だと 0 と出る。そこで縮めると仕事が止まる
- 分割数がレプリカ数の上限になる: 式がいくつを要求しても、超えたぶんは働かない
① ラグが取れる形と取れない形がある
配信元を 2 つ作る。まずログ型で、読んでも消えないので位置を 2 つ持つ:
// Log は読んでも消えない配信元(Kafka のトピックなど)。
//
// 書いた位置と読んだ位置を別々に持つので、その差が「まだ処理していない量」になる。
type Log struct {
written int // 書き込まれた総数
committed int // 読み終えた位置
}
// Append は n 件を末尾に足す。
func (l *Log) Append(n int) { l.written += n }
// Commit は n 件を読み終えたことにする。書いた位置は動かない。
func (l *Log) Commit(n int) {
l.committed += n
if l.committed > l.written {
l.committed = l.written
}
}
// Written は書き込まれた総数。
func (l *Log) Written() int { return l.written }
// Committed は読み終えた位置。
func (l *Log) Committed() int { return l.committed }
// Lag は書いた位置と読んだ位置の差。ログ型でしか取れない。
//
// 読んでも消えないので、過去に何件流れたかと、どこまで読んだかが両方残る。
// この 2 つがあるから引き算ができる。
func (l *Log) Lag() int { return l.written - l.committed }次にキュー型。取り出すと見えなくなり、確定(ack)で消え、失敗(nack)で戻る:
// Queue は読んだら消える配信元(SQS などの分配キュー)。
//
// 取り出された件数は「処理中」に移り、確定(ack)で消え、失敗(nack)で戻る。
// 書いた位置という概念が無いので、ラグは定義できない。
type Queue struct {
visible int // まだ誰も取り出していない数
inFlight int // 取り出されたが、まだ確定していない数
}
// Enqueue は n 件を入れる。
func (q *Queue) Enqueue(n int) { q.visible += n }
// Receive は n 件を取り出す。取り出した分は見えなくなり、処理中へ移る。
func (q *Queue) Receive(n int) int {
if n > q.visible {
n = q.visible
}
q.visible -= n
q.inFlight += n
return n
}
// Ack は n 件を確定させる。ここで初めて消える。
func (q *Queue) Ack(n int) {
if n > q.inFlight {
n = q.inFlight
}
q.inFlight -= n
}
// Nack は n 件を戻す。処理に失敗したものは、また見えるようになる。
func (q *Queue) Nack(n int) {
if n > q.inFlight {
n = q.inFlight
}
q.inFlight -= n
q.visible += n
}
// Visible は取り出せる数。KEDA の SQS スケーラが既定で見るのはこちら。
func (q *Queue) Visible() int { return q.visible }
// InFlight は処理中の数。
func (q *Queue) InFlight() int { return q.inFlight }
// Outstanding は「まだ終わっていない総数」。見えている数と処理中の数の合計。
//
// 見えている数だけを見ると、処理中が多いときに「もう空いた」と誤読する。
func (q *Queue) Outstanding() int { return q.visible + q.inFlight }同じ「200 件溜まっている」状態から 50 件処理して、何が取れるかを比べた。
| 書いた総数 | 読んだ位置 | 残り | |
|---|---|---|---|
| ログ型 | 200 | 50 | 150 |
| キュー型 | 取れない | 取れない | 150 |
残りはどちらも 150 で一致する。テストでこの一致を固定した。違うのは、そこに至る位置を持っているかどうかになる。
ログ型は読んでも消えないので、「これまで何件流れたか」が残り続ける。だから引き算ができ、それがラグになる。キュー型は取り出した時点で消えるため、総数という情報がどこにも無い。Lag() という関数を定義しようがない。
この差は、メッセージキューとPub/Subで作った 2 つのモデルの違いがそのまま出たものになる。消えるか残るかという配送の設計が、監視で何を測れるかを決めている。
② 見えている数だけでは足りない
キュー型には落とし穴がある。取り出された分は見えなくなるので、全員が処理中のとき、見えている数は 0 になる。
200 件を全部取り出した直後で測るとこうなった。
| 見ているもの | 値 | 出る答え |
|---|---|---|
| 見えている数 | 0 | 0 個 |
| まだ終わっていない数(見えている + 処理中) | 200 | 20 個 |
見えている数で判断すると「もう仕事が無い」と読んで 0 個まで縮める。だが実際には 200 件が処理中で、20 個ぶんの仕事が残っている。テストでこの 2 つの答えの食い違いを固定した。
しかも処理に失敗すると、処理中だった分は見えている側へ戻ってくる。テストで、200 件を nack すると見えている数が 0 から 200 へ戻ることを固定した。縮めた直後に全部戻ってくるという最悪の順序があり得る。
KEDA の SQS スケーラが既定で見るのは見えている数(ApproximateNumberOfMessages)で、処理中(ApproximateNumberOfMessagesNotVisible)を足すかどうかは設定になっている。どちらが正しいかは処理時間で決まる。処理が一瞬なら処理中はほぼ 0 なので気にしなくてよく、処理が長いほどこの穴が開く。
③ 分割数がレプリカ数の上限になる
滞留から必要な数を出す式は、現在のレプリカ数に依存しない:
// Desired は滞留から必要なレプリカ数を返す。
//
// KEDA の lagThreshold と同じ式で、ceil(滞留 ÷ 1 レプリカが引き受ける量)。
// 現在のレプリカ数が式に出てこないので、何個で動いていても同じ滞留からは同じ答えが出る。
func Desired(backlog, perReplica int) int {
if perReplica <= 0 || backlog <= 0 {
return 0
}
return (backlog + perReplica - 1) / perReplica
}だから 2000 件溜まっていて 1 個が 10 件を引き受けるなら、答えは 200 個になる。ところがログ型では、1 つの分割を同時に読めるのは 1 レプリカだけという制約がある:
// Topic は分割された配信元。同時に読める数がパーティション数で決まる。
type Topic struct {
// Partitions は分割数。1 パーティションを同時に読めるのは 1 レプリカだけ。
Partitions int
// PerReplica は 1 レプリカが 1 単位時間に捌ける件数。
PerReplica int
}
// EffectiveReplicas は、実際に仕事をするレプリカ数を返す。
//
// パーティションを超えたぶんは、割り当てが無いので何もしない。
func (t Topic) EffectiveReplicas(replicas int) int {
if replicas > t.Partitions {
return t.Partitions
}
if replicas < 0 {
return 0
}
return replicas
}
// Idle は割り当てが無く遊ぶレプリカ数。
func (t Topic) Idle(replicas int) int { return replicas - t.EffectiveReplicas(replicas) }
// Throughput は 1 単位時間に捌ける件数。
func (t Topic) Throughput(replicas int) int {
return t.EffectiveReplicas(replicas) * t.PerReplica
}
// MaxUseful は、これ以上増やしても捌ける量が増えないレプリカ数。
func (t Topic) MaxUseful() int { return t.Partitions }分割 10、1 レプリカ 10 件で測った。
| レプリカ | 実際に働く | 遊ぶ | 捌ける件数 |
|---|---|---|---|
| 5 | 5 | 0 | 50 |
| 10 | 10 | 0 | 100 |
| 15 | 10 | 5 | 100 |
| 20 | 10 | 10 | 100 |
10 を超えると、立てても捌ける量が 1 件も増えない。テストで、20 個のときの処理量が 10 個のときと同じであることを固定した。
式が要求する数と噛み合わないと、こうなる。
滞留 2000 件 / 1 レプリカ 10 件 → 式は 200 個を要求
だが分割は 10 なので、働くのは 10 個。190 個は遊ぶ
捌け終わるまで 20 単位時間(200 個でも 10 個でも同じ)
分割を 50 に増やすと 4 単位時間190 個ぶんの資源を使って、速さは 1 ミリも変わらない。テストで、200 個と 10 個で所要時間が一致すること、分割を増やして初めて短くなることを固定した。
だから maxReplicaCount を分割数に合わせる。指標が上限を持たないこと(①で見た利点)と、実際に増やせる数に上限があることは別で、指標側に上限が無いからこそ、外から蓋をする必要がある。
動かす
下のデモは、ログ型とキュー型で取れる情報の違いを並べ、処理中の件数を動かすと見えている数だけの判断が破綻する様子を見る。分割数とレプリカ数を動かすと、上限を超えたところで捌ける量が伸びなくなる。
式は現在のレプリカ数に依存しないので、滞留が青天井なら要求も青天井になる。だが実際に働ける数には 物理的な上限があり、指標の側には現れない。上限の無い指標を使うときほど、外から蓋をする必要がある。
設計の観点
- 測れるものは配送の設計が決める: 読んでも消えないなら位置の差が取れ、消えるなら残りしか取れない。監視のために配送を選ぶことすらある
- 「見えている」は「残っている」ではない: 取り出し済みで未確定のものを数に入れるかどうかで、答えが 0 個と 20 個に割れる
- 処理時間の長さが穴の大きさを決める: 処理が一瞬なら処理中はほぼ 0。長いほど、見えている数だけの判断が外れる
- 失敗は見えている側へ戻る: 縮めた直後に全部戻ってくる順序があり得る。縮小を遅らせる理由がここにもある
- 上限の無い指標にこそ蓋が要る: 滞留は青天井なので、式は青天井の数を返す。物理的に働ける数で外から止める
- 遊ぶレプリカは害がある: 速くならないうえ、資源を占め、スケジューラの配置先も食う
対照と実例
| 配信元 | 取れる指標 | ラグ | 上限の要因 | KEDA での例 |
|---|---|---|---|---|
| ログ型(読んでも消えない) | 書いた位置・読んだ位置 | 取れる | 分割数 | Kafka、Pulsar |
| キュー型(読んだら消える) | 見えている数・処理中の数 | 取れない | 明示しなければ無し | SQS、RabbitMQ、Redis |
| 使用率 | 比率 | — | 100% で頭打ち | CPU、メモリ |
| 時刻 | 指標でない | — | — | cron |
上の 3 行が、カスタム指標で見た「上限のある指標とない指標」の続きになる。滞留は上限が無いので必要な数がその場で出るが、そのぶん物理的な上限を別に置くことになる。時刻は量ですらないので、カスタム指標の Activation と同じく「量でなく有無」の判断に属する。
裏どり:
- Kafka の 1 パーティション 1 コンシューマ: 同じコンシューマグループの中で、1 つのパーティションを同時に読めるのは 1 つだけという規約になっている。だからパーティション数がそのまま並列度の上限で、後から分割を増やすことはできても減らすことはできない
- SQS の 2 つの数:
ApproximateNumberOfMessagesが見えている数、ApproximateNumberOfMessagesNotVisibleが取り出し済みで未確定の数。どちらも近似値で、分散して保持されているため厳密な瞬間値ではない - 可視性タイムアウト: 取り出したメッセージは一定時間見えなくなり、その間に確定しなければ自動で戻る。処理がこの時間を超えると、まだ処理しているのに別のレプリカへ再配達される
- ラグは遅れの量であって、遅れの時間ではない: 10 万件のラグが 1 秒ぶんなのか 1 時間ぶんなのかは、流量が分からないと言えない。Cloud Monitoring で読み替えるで見た Pub/Sub の「件数と経過時間の 2 つ組」は、この曖昧さへの答えになっている
- KEDA 自身も定期的に見に行く: 外部システムを
pollingIntervalごとに問い合わせるので、メトリクスの集め方で測った pull の性質をそのまま持つ。監視対象が増えれば 1 周の予算に当たる
簡略化したこと
- 1 つの配信元だけ: 複数のトリガーを組み合わせたときの合成(最大を採る)は扱わない
- 時間は単位時間の整数: 可視性タイムアウトや再配達の時間は数えない
- 順序を持たない: パーティション内の順序保証や、キーによる割り当ては扱わない
- 消費側の失敗率なし: nack は手で起こす。実際の失敗率から必要数を出す話はしない
- 認証は扱わない: 各サービスへの接続情報の渡し方は範囲外
- トリガーの全一覧は載せない: 取れる情報の型に絞った。個々のスケーラの設定項目は公式ドキュメントに譲る
参考資料
- KEDA: Scalers — 各配信元から何を読むか
- KEDA: Apache Kafka scaler —
lagThresholdと分割数の関係 - KEDA: AWS SQS scaler — 見えている数と処理中の数のどちらを見るか
- Kafka: Consumer Groups — 1 パーティション 1 コンシューマの規約
- 実装: orchestration/scaletrigger