๐ง ๋ค์ด๊ฐ๊ธฐ ์
์ปจ์๋จธ๋ฅผ ํ๋ ๋ ๋์ ์ ๋ฟ์ธ๋ฐ ๊ทธ๋ฃน ์ ์ฒด์ ์๋น๊ฐ ๋ช ์ด์ฉ ๋ฉ์ถ๋ ๊ฒฝํ์ ํํ๋ค. ๊ทธ๋ ๋ค๊ณ ์ด๊ฒ ์นดํ์นด์ ๋ฒ๊ทธ๋ ์๋๊ณ , classic ํ๋กํ ์ฝ์ด ํฉ์๋ฅผ ๋ง๋ค์ด ๋ด๋ ๋ฐฉ์์์ ๊ตฌ์กฐ์ ์ผ๋ก ๋ฐ๋ผ ๋์ค๋ ๊ฒฐ๊ณผ๋ค. ๊ทธ๋์ "๋ฆฌ๋ฐธ๋ฐ์ค๊ฐ ๋๋ฆฌ๋ค"๋ฅผ ํ๋ ๋ฌธ์ ๋ก๋ง ๋ณด๋ฉด ์ด๋๋ฅผ ๊ฑด๋๋ ค์ผ ํ๋์ง๊ฐ ๋๋ด ๋ณด์ด์ง ์๋๋ค.
๊ฐ์ ๋ฆฌ๋ฐธ๋ฐ์ค๋ผ๋ ์ด๋ฆ ์๋์์๋ ๋ฉ์ถ๋ ๋ฐฉ์์ ๋ค๋ฅด๋ค. classic ์ eager rebalancing ์ ๋ณด์ ํํฐ์ ์ ์ ๋ถ ๋๊ณ ๊ทธ๋ฃน ์ ์ฒด๊ฐ JoinGroup ๋ฐฐ๋ฆฌ์ด์ ๋ชจ์ธ๋ค. incremental cooperative rebalancing ์ ์ด๋ํ ํํฐ์ ๋ง ๋์ง๋ง, JoinGroup ๊ณผ SyncGroup ์ผ๋ก ์๋ณตํ๋ classic ์ ๊ณจ๊ฒฉ์ ๊ทธ๋๋ก ์ฌ์ฉํ๋ค. KIP-848 ์ consumer ํ๋กํ ์ฝ์ ๊ทธ ๊ณจ๊ฒฉ ์์ฒด๋ฅผ ์์ ๊ณ ๋ฉค๋ฒ๋ณ heartbeat ๋ก ๋ชฉํ ํ ๋น์ ์๋ ดํ๋ค.
์ฐจ์ด๋ฅผ ๊ฐ๋ฅด๋ ์ถ์ ๋ ๊ฐ์ง๋ค. assignor ๊ฐ ์ด๋์์ ์คํ๋๋๊ฐ, ๊ทธ๋ฆฌ๊ณ ๋๊ฐ ๋๊ตฌ๋ฅผ ๊ธฐ๋ค๋ ค์ผ ํ๋๊ฐ๋ค. leader ์ปจ์๋จธ์ JVM ๊ณผ ๊ทธ๋ฃน ๋จ์ ๋ฐฐ๋ฆฌ์ด๊ฐ ๊ฒฐํฉํ๋ฉด ๊ฐ์ฅ ๋๋ฆฐ ๋ฉค๋ฒ๊ฐ ์ ์ฒด๋ฅผ ๋ถ์ก๋๋ค. ๊ณ์ฐ์ ๋ธ๋ก์ปค๋ก ์ฎ๊ธฐ๊ณ ๋ฉค๋ฒ๋ณ reconciliation ์ผ๋ก ๋ฐ๊พธ๋ฉด ์ํฅ์ ๋ฐ์ง ์๋ ํํฐ์ ์ ๊ณ์ ์๋นํ ์ ์๋ค. ์ด ๋ ์ถ์ ๊ธฐ์ค์ผ๋ก classic ์ eager/cooperative ์ KIP-848 consumer ํ๋กํ ์ฝ์ ์์๋๋ก ๋ฐ๋ผ๊ฐ๋ค.
1. ๋ฆฌ๋ฐธ๋ฐ์ค์ ๊ตฌ์ฑ์์
๋ฆฌ๋ฐธ๋ฐ์ค์ ๋จ๊ณ์ ๋ค์ด๊ฐ๊ธฐ ์ ์ ๊ตฌ์ฑ์์์ ๊ทธ๋ฃน์ ์ํ ์ ์ด๋ถํฐ ์ ๋ฆฌํด๋ณด๊ณ ์ ํ๋ค.
1-1. ์ธ ๊ฐ์ง ๊ตฌ์ฑ์์
| ์ญํ | ํ๋ ์ผ |
|---|---|
| GroupCoordinator | ํน์ ๋ธ๋ก์ปค๊ฐ ๊ทธ๋ฃน ํ๋๋ฅผ ๋ด๋นํ๋ค. ๋ฉค๋ฒ ๋ช ๋จ, ๊ทธ๋ฃน ์ํ, ์คํ์ ์ ๋ค๊ณ ์๋ค |
| ๋ฉค๋ฒ(์ปจ์๋จธ) | ์ฃผ๊ธฐ์ ์ผ๋ก heartbeat ๋ฅผ ๋ณด๋ด ์ด์์์์ ์๋ฆฌ๊ณ , ํ ๋น๋ฐ์ ํํฐ์ ์ ์๋นํ๋ค |
| assignor | ์ค์ ํํฐ์ ๋ถ๋ฐฐ๋ฅผ ๊ณ์ฐํ๋ ์๊ณ ๋ฆฌ์ฆ. ์ด๊ฒ ์ด๋์ ์คํ๋๋์ง๊ฐ ๋ ํ๋กํ ์ฝ์ ๊ฒฐ์ ์ ์ฐจ์ด๋ค |
assignor๊ฐ ์คํ๋๋ ์์น๊ฐ ์ด ๊ธ ์ ์ฒด์ ์ถ์ด๋ค. assignor ๊ฐ ์ปจ์๋จธ์ JVM ์์ ๋๋๊ฐ ๋ธ๋ก์ปค์์ ๋๋๊ฐ์ ๋ฐ๋ผ, ๋ฆฌ๋ฐธ๋ฐ์ค๊ฐ ๊ทธ๋ฃน ๋จ์ ์๋ณต์ด ๋๊ธฐ๋ ํ๊ณ ๋ฉค๋ฒ๋ณ ์๋ ด ๋ฃจํ๊ฐ ๋๊ธฐ๋ ํ๋ค.
1-2. classic ๊ทธ๋ฃน์ ์ํ ์ ์ด
classic ํ๋กํ ์ฝ์์ ๋ฆฌ๋ฐธ๋ฐ์ค๋ Stable ์ ๋ฒ์ด๋ ํ ๋ฐํด ๋์์ค๋ ์ฌ์ ์ด๋ค.
stateDiagram-v2
[*] --> Empty
Empty --> PreparingRebalance: ์ฒซ JoinGroup ๋์ฐฉ
Stable --> PreparingRebalance: ๋ฉค๋ฒ ๋ณํ ๊ฐ์ง
PreparingRebalance --> CompletingRebalance: ์ ์ JoinGroup ๋์ฐฉ = join phase ์ข
๋ฃ
CompletingRebalance --> Stable: leader ์ SyncGroup ํ ๋น ์์ = sync phase ์ข
๋ฃ
PreparingRebalance --> Empty: ์ ์ ์ดํ
[!IMPORTANT]
classic eager ์์ ์ฝ๋๋ค์ดํฐ๊ฐPreparingRebalance๋ก ์ ํํ๋ค๊ณ ์ฆ์ ๋ชจ๋ ์๋น๊ฐ ๋ฉ์ถ๋ ๊ฒ์ ์๋๋ค. ๊ธฐ์กด ๋ฉค๋ฒ๋ ๋ฆฌ๋ฐธ๋ฐ์ค๋ฅผ ์ธ์งํ๊ณ poll ๊ฒฝ๋ก์์ ํํฐ์ ์ revoke ํ๊ธฐ ์ ๊น์ง ์๋นํ ์ ์๋ค. ๊ฐ ๋ฉค๋ฒ๊ฐ ์ ์ฒด ํํฐ์ ์ revoke ํ ์์ ๋ถํฐ ์ ํ ๋น์ ๋ฐ์ ๋๊น์ง๊ฐ ๊ทธ ๋ฉค๋ฒ์ ์ค์ ์๋น ์ ์ง ๊ตฌ๊ฐ์ด๋ค.
2. classic ํ๋กํ ์ฝ
2-1. ๋ฆฌ๋ฐธ๋ฐ์ค ๋ผ์ด๋์ ๋จ๊ณ๋ณ ์ ๊ฐ
ํ ๋ผ์ด๋๋ ํฌ๊ฒ ์๋์ ๊ฐ์ ๋ค ๋จ๊ณ๋ก ์ด๋ฃจ์ด์ง๋ค.
| ๋จ๊ณ | ๋จ๊ณ๋ช | ๋ํ ์์ฒญ๊ณผ ์ํ |
|---|---|---|
| 1๋จ๊ณ | rebalance trigger | PreparingRebalance ์ ํ, ์๋ต ์๋ฌ ์ฝ๋ REBALANCE_IN_PROGRESS |
| 2๋จ๊ณ | join phase | JoinGroup ์์ฒญ, JoinGroup ์ rebalance timeout ๋๊ธฐ |
| 3๋จ๊ณ | client-side assignment | leader ์ปจ์๋จธ๊ฐ ConsumerPartitionAssignor ์คํ |
| 4๋จ๊ณ | sync phase | SyncGroup ์์ฒญ๊ณผ ์๋ต, CompletingRebalance โ Stable |
[!NOTE]
Kafka ์ ํด๋ผ์ด์ธํธ์ธก ํ ๋น ์ ์ ๋ฌธ์๋ ๊ทธ๋ฃน ๋ฉค๋ฒ์ญ ํ๋กํ ์ฝ์ "join phase ์ sync phase ๋ ๋จ๊ณ"๋ก ์ค๋ช ํ๋ค.
flowchart LR
S1[1๋จ๊ณ rebalance trigger<br/>PreparingRebalance ์ ํ] --> S2[2๋จ๊ณ join phase<br/>JoinGroup ๋ฐฐ๋ฆฌ์ด]
S2 --> S3[3๋จ๊ณ client-side assignment<br/>leader ์ง๋ช
๊ณผ ํ ๋น ๊ณ์ฐ]
S3 --> S4[4๋จ๊ณ sync phase<br/>SyncGroup ์๋ต]
S4 --> ST[Stable<br/>ํ์ฌ generation ์ ํ ๋น ํ์ ]์ ์ง ์๊ฐ์ด ๊ธธ์ด์ง๋ ์ง์ ์ ๋ฏธ๋ฆฌ ๋งํด ๋๋ฉด ๊ทธ๋ฆผ์ ์ฝ๊ธฐ๊ฐ ์ฝ๋ค. join phase ๋ ์๋ ค์ง ๋ฉค๋ฒ๋ค์ด ๋ค์ ํฉ๋ฅํ ๋๊น์ง ๊ธฐ๋ค๋ฆฌ๋ฏ๋ก, ๊ฐ์ฅ ๋ฆ๊ฒ JoinGroup ์ ๋ณด๋ด๋ ๋ฉค๋ฒ ํ ๋ช ์ด ์ด ๊ตฌ๊ฐ์ ๊ธธ์ด๋ฅผ ๊ฒฐ์ ํ ์ ์๋ค.
2-1-1. 1๋จ๊ณ: rebalance trigger
์ฝ๋๋ค์ดํฐ์ ์ปจ์๋จธ๋ ๋ฉค๋ฒ์ญ์ด๋ ๊ตฌ๋ ๋ฉํ๋ฐ์ดํฐ์ ๋ณํ๋ฅผ ๊ฐ์งํ๋ฉด ์ฌํฉ๋ฅ๋ฅผ ์์ํ๋ค. Kafka 4.3 ConsumerRebalanceListener Javadoc์ ๊ธฐ์ค์ผ๋ก ๋ํ์ ์ธ rebalance trigger ๋ ๋ค์๊ณผ ๊ฐ๋ค.
- ์ ๊ท ๋ฉค๋ฒ์ JoinGroup ๋์ฐฉ
- ๋ฉค๋ฒ ์ดํ, ์ฆ LeaveGroup ์์
- ๊ตฌ๋ ํ ํฝ ๋ณ๊ฒฝ
- session timeout ๋ง๋ฃ, ์ฆ heartbeat ๊ฐ ๋๊ธด ๋ฉค๋ฒ
- ๊ตฌ๋ ์ค์ธ ํ ํฝ์ ํํฐ์ ์ ๋ณ๊ฒฝ
- ์ ๊ท์ ๊ตฌ๋ ์ ์๋ก ๋งค์นญ๋๋ ํ ํฝ์ ์์ฑ ๋๋ ๊ตฌ๋ ํ ํฝ์ ์ญ์
์ด ๋จ๊ณ์์ ์ค๊ฐ๋ ๋ฉ์์ง๋ง ๋ผ์ด ๋ณด๋ฉด ์๋์ ๊ฐ๋ค.
sequenceDiagram
participant A as Consumer A
participant B as Consumer B
participant GC as GroupCoordinator
participant N as Consumer C - ์ ๊ท
Note over GC: Stable
N->>GC: JoinGroup ์ ๊ท ํฉ๋ฅ
Note over GC: PreparingRebalance ๋ก ์ ํ
A->>GC: Heartbeat ์ฃผ๊ธฐ์ ํธ์ถ
GC-->>A: REBALANCE_IN_PROGRESS
B->>GC: Heartbeat ์ฃผ๊ธฐ์ ํธ์ถ
GC-->>B: REBALANCE_IN_PROGRESS
Note over A,B: ๊ฐ์์ ๋ค์ heartbeat ๋ ๋น๋ก์ ์ธ์ง์ฃผ๋ชฉํ ๊ณณ์ ์ฝ๋๋ค์ดํฐ์์ A ์ B ๋ก ๋๊ฐ๋ ํ์ดํ๊ฐ ์ ๋ถ ์๋ต ํ์ดํ๋ผ๋ ์ ์ด๋ค. ์ฝ๋๋ค์ดํฐ๊ฐ ๋ฉค๋ฒ์๊ฒ ๋ฅ๋์ ์ผ๋ก ํธ์ํ๋ ์ ์ ๊ทธ๋ฆผ ์ด๋์๋ ์๋ค. ๊ฐ ๋ฉค๋ฒ๊ฐ ๋ณด๋ด๋ ์ฃผ๊ธฐ์ heartbeat ์ ์๋ต์ REBALANCE_IN_PROGRESS ์๋ฌ ์ฝ๋๋ฅผ ์ค์ด์ ๊ฐ์ ์ ์ผ๋ก ์๋ฆฐ๋ค. ์นดํ์นด์ ์์ฒญ๊ณผ ์๋ต ๋ชจ๋ธ์ด ํด๋ผ์ด์ธํธ ์ฃผ๋(pull)์ด๊ธฐ ๋๋ฌธ์ ์ด๋ ๊ฒ ์ค๊ณ๋ ๊ฒ์ด๋ค.
classic ์ heartbeat.interval.ms ๊ธฐ๋ณธ๊ฐ์ 3์ด์ด๋ฏ๋ก, ์ ์์ ์ธ ๋คํธ์ํฌ ์กฐ๊ฑด์์๋ trigger ๋ฐ์ ํ ๊ธฐ์กด ๋ฉค๋ฒ๊ฐ ๋ค์ heartbeat ์๋ต์ ๋ฐ์ ๋๊น์ง 0~ํ ์ค์ ์ฃผ๊ธฐ ์ ๋์ ์ธ์ง ์ง์ฐ์ด ์๊ธธ ์ ์๋ค. ๋ค๋ง ์ด ์๊ฐ ์ ์ฒด๋ฅผ ์๋น ์ ์ง๋ก ๊ณ์ฐํ๋ฉด ์ ๋๋ค. KIP-62์ ๋ฐ๋ผ heartbeat ๋ ๋ฐฑ๊ทธ๋ผ์ด๋์์ ๊ณ์๋๊ณ , ์ค์ ์ฌํฉ๋ฅ์ revoke ์ฝ๋ฐฑ์ ์ ํ๋ฆฌ์ผ์ด์
์ค๋ ๋๊ฐ poll ๋ก ๋์์์ ๋ ์งํ๋๋ค. ๋ฉค๋ฒ๋ revoke ์ ๊น์ง ๊ธฐ์กด ํ ๋น์ ์๋นํ ์ ์๋ค.
2-1-2. 2๋จ๊ณ: join phase
ํต๋ณด๋ฐ์ ๋ฉค๋ฒ๋ค์ ์๊ธฐ ๊ตฌ๋ ์ ๋ณด, ์ฆ ์ํ๋ ํ ํฝ๊ณผ ์ง์ํ๋ assignor ๋ชฉ๋ก์ ๋ด์ JoinGroup ์ ๋ค์ ๋ณด๋ธ๋ค. eager assignor ๋ผ๋ฉด ์ด ์ง์ ์ ๋ณด์ ํํฐ์ ์ ๋ถ๋ฅผ revoke ํ๋ฏ๋ก, ์ฌ๊ธฐ์๋ถํฐ ๊ทธ๋ฃน ์ ์ฒด๊ฐ ์๋น ์ ์ง ์ํ๋ค.
[!TIP]
assignor ์ ์ข ๋ฅ:partition.assignment.strategy์ ๋ฃ์ ์ ์๋ ๊ธฐ๋ณธ ๊ตฌํ์ ๋ค ๊ฐ์ง์ด๊ณ , "๋๊ตฌ์๊ฒ ๋ฌด์์ ์ฃผ๋๊ฐ"(๋ถ๋ฐฐ ๊ท์น)์ "์ธ์ ํํฐ์ ์ ๋๋๊ฐ"(rebalance protocol) ๋ ์ถ์ผ๋ก ๊ฐ๋ฆฐ๋ค.
assignor ๋ถ๋ฐฐ ๊ท์น ํ๋กํ ์ฝ ์ฑ๊ฒฉ RangeAssignorํ ํฝ๋ณ๋ก ํํฐ์ ์ ์ ๋ ฌํด ์ฐ์ ๊ตฌ๊ฐ์ ๋๋ ์ค๋ค eager ๊ธฐ๋ณธ ๋ชฉ๋ก์ ์ฒซ ํญ๋ชฉ์ด๋ผ ๊ธฐ๋ณธ ์ ํ๋๋ค. ํ ํฝ ์ฌ๋ฌ ๊ฐ๋ฅผ ๊ฐ์ด ๊ตฌ๋ ํ๋ฉด ์ ๋ฒํธ ์ปจ์๋จธ์ ์ ๋ฆด ์ ์๋ค RoundRobinAssignor์ ์ฒด ํํฐ์ ์ ํ ์ค๋ก ์ธ์ ๋ฉค๋ฒ์๊ฒ ๋ฒ๊ฐ์ ์ค๋ค eager ๊ท ๋ฑํ์ง๋ง ๋ฆฌ๋ฐธ๋ฐ์ค๋ง๋ค ๋ฐฐ์น๊ฐ ํต์งธ๋ก ํ๋ค๋ฆฐ๋ค StickyAssignor์ต๋ํ ๊ท ๋ฑํ๊ฒ ๋๋๋ ์ง์ ํ ๋น์ ์ต๋ํ ์ ์งํ๋ค eager ์ด๋์ ์ค์ง๋ง revoke ๋ ์ฌ์ ํ ์ ๋ฉด์ ์ด๋ค CooperativeStickyAssignorsticky ์ ๊ฐ์ ๊ณ์ฐ cooperative ์ด๋ ๋์ ํํฐ์ ๋ง revoke ํ๋ค. KIP-429 ์ ๊ตฌํ์ฒด ๋ ์ถ์ด ์ ๋ฐ๋ก์ธ์ง๊ฐ ํต์ฌ์ด๋ค. sticky ๋ "๋ช ๊ฐ๊ฐ ์์ง์ด๋๊ฐ"๋ฅผ ์ค์ด๊ณ , cooperative ๋ "์์ง์ด์ง ์๋ ํํฐ์ ๊น์ง ๋ฉ์ถ๋๊ฐ"๋ฅผ ์์ค๋ค.
StickyAssignor๋ ์ ์๋ฆฌ์ ๋จ์ ํํฐ์ ๊น์ง ์ผ๋จ ์ ๋ถ ๋์๋ค๊ฐ ๋ค์ ๋ฐ์ง๋ง, ์ฌํ ๋น ํ ์ํ ๋ณต๊ตฌ ๋น์ฉ์ด ํฐ ์ํฌ๋ก๋์์๋ ์ฌ์ ํ ๊ฐ์ ํ๋ค.eager ์ cooperative ๋ฅผ ๊ฐ๋ฅด๋ ์ค์ฒด๋
ConsumerPartitionAssignor์supportedProtocols()๊ฐ์ด๋ค. ๋ค๋ง ์ฝ๋๋ค์ดํฐ๊ฐ ์ด ๊ฐ์ ์ง์ ๋น๊ตํ๋ ๊ฒ์ ์๋๊ณ , JoinGroup ์ ์ค๋ฆฐ assignor ์ด๋ฆ ๋ชฉ๋ก์์ ๋ชจ๋ ๋ฉค๋ฒ๊ฐ ์ง์ํ๋ ์ด๋ฆ์ ๊ณ ๋ฅธ๋ค. ์ ํ๋ assignor ์supportedProtocols()์ ๋ฐ๋ผ ๊ฐ ํด๋ผ์ด์ธํธ๊ฐ revoke ๋ฐฉ์์ ๊ฒฐ์ ํ๋ค. ๊ณต์ Javadoc์ด cooperative ํ์ฑํ์ ๋ชจ๋ ์ปจ์๋จธ๊ฐ ํด๋น assignor ๋ฅผ ์ฌ์ฉํด์ผ ํ๋ค๊ณ ๋ช ์ํ๋ ์ด์ ๋ค.Kafka 4.3 consumer ์ค์ ๋ฌธ์์ ๊ธฐ๋ณธ ๋ชฉ๋ก์
[RangeAssignor, CooperativeStickyAssignor]์ด๋ฉฐ ์์RangeAssignor๊ฐ ์ ํ๋๋ค. cooperative ๋ก ์ ํํ๋ ค๋ฉด ๋ชจ๋ ๋ฉค๋ฒ์ ๋ชฉ๋ก์์RangeAssignor๋ฅผ ์ ๊ฑฐํด์ผ ํ๋ค. [[3|consumer ํ๋กํ ์ฝ]]๋ก ๋์ด๊ฐ๋ฉด ์ด client-side ํด๋์ค ์ ํ์group.consumer.assignors์group.remote.assignor์ server-side ์ ํ์ผ๋ก ๋์ฒด๋๋ค.
๋๋ฆฐ ๋ฉค๋ฒ๊ฐ ํ๋ ์์ฌ ์์ ๋ ์ด phase ๊ฐ ์ด๋ป๊ฒ ๋์ด์ง๋์ง๋ฅผ ๊ทธ๋ฆผ์ผ๋ก ๋ณด๋ฉด ์ด๋ ๋ค.
sequenceDiagram
participant A as Consumer A
participant B as Consumer B - ๋๋ฆฐ ๋ฉค๋ฒ
participant GC as GroupCoordinator
Note over A,B: eager ๋ผ๋ฉด ์ฌ๊ธฐ์ ์ ํํฐ์
revoke โ ์๋น ์ ์ง
A->>GC: JoinGroup ๊ตฌ๋
์ ๋ณด + assignor ๋ชฉ๋ก
Note over GC: A ๋์ฐฉ. ์์ง ์ ์์ด ์๋๋ผ ๋ผ์ด๋๋ฅผ ๋ซ์ง ์์
Note over B: ๋ฌด๊ฑฐ์ด poll ์ฒ๋ฆฌ ์ค์ด๋ผ ์์ง ๋ชป ๋ณด๋
Note over A: ํํฐ์
์ ์ด๋ฏธ ๋์ ์ฑ๋ก ๋๊ธฐ
B->>GC: JoinGroup ๋ค๋ฆ๊ฒ ๋์ฐฉ
Note over GC: ์ ์ ๋์ฐฉ โ ๋ฐฐ๋ฆฌ์ด ํด์ ํต์ฌ์ ์ฝ๋๋ค์ดํฐ๊ฐ ์๊ณ ์๋ ๋ฉค๋ฒ ์ ์์ด ๋ชจ์ผ ๋๊น์ง ๋ผ์ด๋๋ฅผ ๋ซ์ง ์๋๋ค๋ ์ ์ด๋ค. ์ด๊ฒ ๋ฐฐ๋ฆฌ์ด(barrier)์ ์ค์ฒด๋ค. KIP-848 ์ ์ด ๊ตฌ์กฐ๋ฅผ group-wide synchronization barrier๋ผ๊ณ ๋ถ๋ฅธ๋ค.
์ฆ, ๊ฐ์ฅ ๋ฆ๊ฒ ์ค๋ ํ ๋ช
์ด ์ ์ฒด ์ถ๋ฐ ์๊ฐ์ ๊ฒฐ์ ํ๋ค.
๋์ ๋ฉค๋ฒ๋ง ์๋ ์ผ๋ฐ์ ์ธ ๊ทธ๋ฃน์์ ๋ฐฐ๋ฆฌ์ด๋ ๋ค์ ๋ ๊ฒฝ๋ก๋ก ์ด๋ฆฐ๋ค.
flowchart TD
P[PreparingRebalance] --> W{์๊ณ ์๋ ๋ฉค๋ฒ ์ ์์ด<br/>JoinGroup ์ ๋ณด๋๋}
W -- ์ --> C[๋ฐฐ๋ฆฌ์ด ํด์ <br/>CompletingRebalance ๋ก]
W -- ์๋์ค --> T{rebalance timeout ๊ฒฝ๊ณผ}
T -- ์๋์ค --> W
T -- ์ --> K[๋ฏธ๋์ฐฉ ๋์ ๋ฉค๋ฒ๋ฅผ ๋ช
๋จ์์ ๋นผ๊ณ ์งํ]
K --> CKIP-62์ ๋ฐ๋ผ Java client ๋ max.poll.interval.ms ๊ฐ์ JoinGroup ์ rebalance timeout ์ผ๋ก ๋ณด๋ธ๋ค. ๋ฉค๋ฒ๋ง๋ค ๊ฐ์ด ๋ค๋ฅด๋ฉด ์ฝ๋๋ค์ดํฐ๋ ๊ทธ๋ฃน์์ ๊ฐ์ฅ ํฐ rebalance timeout ์ ์ฌ์ฉํ๋ค. ์ ๊ทธ๋ฆผ์ ๋ฃจํ๊ฐ ์ค๋ ๋์๋ก eager ๋ฉค๋ฒ๊ฐ ํ ๋น์ ๋น์ด ์๊ฐ๋ ๊ธธ์ด์ง๋ค.
- ๋๋ฆฐ ๋ฉค๋ฒ ํ๋๊ฐ ์ ์์ ๋ถ์ก๋ ๊ฒฝ์ฐ: ์ด๋ค ๋ฉค๋ฒ๊ฐ ๋ฌด๊ฑฐ์ด ๋ ์ฝ๋ ์ฒ๋ฆฌ๋ฅผ ๋๋ด์ง ๋ชปํด poll ๋ก ๋์์ค์ง ์์ผ๋ฉด JoinGroup ๋ ๋ฆ์ด์ง๋ค. ์ด๋ฏธ revoke ๋ฅผ ๋ง์น eager ๋ฉค๋ฒ๋ค์ ๊ทธ ๋ฉค๋ฒ์ ํฉ๋ฅ๋ timeout ์ ๊ธฐ๋ค๋ฆฐ๋ค.
- ์๋ตํ์ง ์๋ ๋์ ๋ฉค๋ฒ๋ฅผ ์ ๊ฑฐํ๋ ๊ฒฝ์ฐ: ํ๋ก์ธ์ค๊ฐ ์ฃฝ์ด heartbeat ๊ฐ ๋๊ธฐ๋ฉด session timeout ์ผ๋ก ์ ๊ฑฐ๋๋ค. ์ด๋ฏธ ๋ฆฌ๋ฐธ๋ฐ์ค๊ฐ ์งํ ์ค์ธ๋ฐ ๋์ ๋ฉค๋ฒ๊ฐ ์ฌํฉ๋ฅํ์ง ์์ผ๋ฉด rebalance timeout ์ผ๋ก๋ ์ ๊ฑฐ๋ ์ ์๋ค. KIP-62์ ์ ์๋๋ก ๋ ์ค ๋จผ์ ์ถฉ์กฑ๋๋ ๊ฒฝ๋ก๊ฐ ๊ทธ๋ฃน์ ์งํ์ํจ๋ค.
[!WARNING]
์ ์ ๋ฉค๋ฒ์ญ(group.instance.id)์์๋ ๊ท์น์ด ๋ค๋ฅด๋ค. KIP-345์ ๋ฐ๋ผ ๋ฏธ๋์ฐฉ ์ ์ ๋ฉค๋ฒ๋ rebalance timeout ๋ง์ผ๋ก ๋ช ๋จ์์ ์ ๊ฑฐ๋์ง ์๊ณ session timeout ๊น์ง ์๋ณ์๋ฅผ ์ ์งํ๋ค. Kafka 4.3 consumer ์ค์ ๋ฌธ์๋ ์ ์ ๋ฉค๋ฒ๊ฐmax.poll.interval.ms๋ฅผ ๋๊ธฐ๋ฉด ํํฐ์ ์ ์ฆ์ ์ฌํ ๋นํ์ง ์๊ณ heartbeat ๋ฅผ ์ค๋จํ ๋ค session timeout ์ ์ฌํ ๋นํ๋ค๊ณ ์ค๋ช ํ๋ค. ์งง์ ์ฌ์์์ ํก์ํ๋ ๋์ ์ค์ ์ดํ ์ ์ฌํ ๋น์ด ๋ฆ์ด์ง ์ ์๋ ํธ๋ ์ด๋์คํ๋ค.
2-1-3. 3๋จ๊ณ: client-side assignment
์ ์์ด ๋ชจ์ด๋ฉด ์ฝ๋๋ค์ดํฐ๋ generation id ๋ฅผ ํ๋ ์ฌ๋ฆฌ๊ณ , ๋ฉค๋ฒ ์ค ํ๋๋ฅผ leader ๋ก ์ง๋ช ํ๋ค. ์ด ๋จ๊ณ์ ๊ทธ๋ฆผ์์ ์๋ต์ด ์ฐจ๋ณ๋๋ค๋ ๊ฒ์ด๋ค.
sequenceDiagram
participant A as Consumer A - leader ๋ก ์ง๋ช
๋จ
participant B as Consumer B
participant N as Consumer C - ์ ๊ท
participant GC as GroupCoordinator
Note over GC: generation++ ํ ๋ฉค๋ฒ ํ๋๋ฅผ leader ๋ก ์ง๋ช
GC-->>A: JoinGroup ์๋ต - ์ ์์ ๊ตฌ๋
์ ๋ณด ์ ๋ฌ
GC-->>B: JoinGroup ์๋ต - ๋น ๊ฐ
GC-->>N: JoinGroup ์๋ต - ๋น ๊ฐ
Note over A: partition.assignment.strategy ์ ์ง์ ํ<br/>assignor ํด๋์ค๋ฅผ A ์ JVM ์์ ์คํ
Note over GC: ๋ธ๋ก์ปค๋ ์ด ๊ณ์ฐ์ ๊ด์ฌํ์ง ์์์ฌ๊ธฐ๊ฐ classic ์ ๊ฐ์ฅ ํน์ดํ ์ง์ ์ด๋ค. ํ ๋น ๊ณ์ฐ์ ๋ธ๋ก์ปค๊ฐ ์๋๋ผ leader ์ปจ์๋จธ์ JVM ์์ ์ผ์ด๋๋ค. ์ด ์ค๊ณ๋ฅผ ๊ฐ๋ฆฌํค๋ ์ด๋ฆ์ด ๊ณง client-side assignment ์ด๊ณ , ์ด๊ฑธ ๋์
ํ ์ค๊ณ ๋ฌธ์์ ์ ๋ชฉ ์์ฒด๊ฐ "Kafka Client-side Assignment Proposal" ์ด๋ค. partition.assignment.strategy ์ ์ง์ ํ ConsumerPartitionAssignor ๊ตฌํ์ด ๋ก๋๋๊ณ ์คํ๋๋ ์๋ฆฌ๊ฐ ๋ฐ๋ก ์ฌ๊ธฐ๋ค. ์ด ์ค์ ์ด "ํด๋ผ์ด์ธํธ์ธก ์ค์ "์ธ ์ด์ ๊ฐ ์ด๊ฒ์ด๋ค.
@Configuration
class ClassicConsumerConfig {
@Bean
fun consumerFactory(): ConsumerFactory<String, String> = DefaultKafkaConsumerFactory(
mapOf(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG to "kafka:9092",
ConsumerConfig.GROUP_ID_CONFIG to "settlement-worker",
// ์ด ํด๋์ค๋ ๋ธ๋ก์ปค๊ฐ ์๋๋ผ, leader ๋ก ๋ฝํ ์ปจ์๋จธ์ JVM ์์ ์คํ๋๋ค
ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG to
listOf(CooperativeStickyAssignor::class.java.name),
// ...
),
)
}ํต์ฌ์ PARTITION_ASSIGNMENT_STRATEGY_CONFIG ์ ๊ฐ์ด ํด๋์ค ์ด๋ฆ์ด๋ผ๋ ์ ์ด๋ค. ๋ธ๋ก์ปค๋ ์ด ๋ฌธ์์ด์ด ๋ฌด์จ ์๊ณ ๋ฆฌ์ฆ์ธ์ง ์์ง ๋ชปํ๋ค. ๋ฉค๋ฒ๋ค์ด ์ ๊ณ ํ ๋ชฉ๋ก์์ ๋ชจ๋ ๋ฉค๋ฒ๊ฐ ๊ณตํต์ผ๋ก ์ง์ํ๋ assignor ์ด๋ฆ์ ๊ณจ๋ผ์ฃผ๋ ์ญํ ๊น์ง๋ง ํ๋ค.
2-1-4. 4๋จ๊ณ: sync phase
๊ณ์ฐ์ด ๋๋๋ฉด ์ ์์ด SyncGroup ์ ๋ณด๋ด๋๋ฐ, ์ด phase ์ ๊ทธ๋ฆผ์์๋ ์์ฒญ์ด ๋์นญ์ด ์๋๋ค.
sequenceDiagram
participant A as Consumer A - leader
participant B as Consumer B
participant N as Consumer C - ์ ๊ท
participant GC as GroupCoordinator
A->>GC: SyncGroup ์ ์ฒด ํ ๋นํ ํฌํจ
B->>GC: SyncGroup ๋น ์์ฒญ
N->>GC: SyncGroup ๋น ์์ฒญ
Note over GC: ํ ๋นํ๋ฅผ ์ ์ฅํ๊ณ ๋ฉค๋ฒ๋ณ๋ก ์๋ผ ์๋ต
GC-->>A: ๋ด ๋ชซ๋ง
GC-->>B: ๋ด ๋ชซ๋ง
GC-->>N: ๋ด ๋ชซ๋ง
Note over A,N: ํ ๋น ์ฝ๋ฐฑ ์คํ ํ ์๋น ์ฌ๊ฐ
Note over GC: leader ์ ํ ๋น ์์ ํ Stableleader ์ ์์ฒญ์๋ง ์ ์ฒด ํ ๋นํ๊ฐ ์ค๋ ค ์๊ณ ๋๋จธ์ง๋ ์ฌ์ค์ "๊ฒฐ๊ณผ ์ฃผ์ธ์"๋ผ๋ ๋น ์์ฒญ์ด๋ค. ์ฝ๋๋ค์ดํฐ๋ ๊ทธ ํ ๋นํ๋ฅผ ์ ์ฅํ ๋ค, ๊ฐ ๋ฉค๋ฒ์ SyncGroup ์๋ต์ ๊ทธ ๋ฉค๋ฒ ๋ชซ๋ง ์๋ผ์ ๋๋ ค์ค๋ค. generation id ๋ join phase ๊ฐ ์ฑ๊ณต์ ์ผ๋ก ๋๋ ๋ ์ด๋ฏธ ์ฆ๊ฐํ๋ฉฐ, leader ์ SyncGroup ์ ๋ด๊ธด ํ ๋น์ ์ฝ๋๋ค์ดํฐ๊ฐ ๋ฐ์ผ๋ฉด ๊ทธ๋ฃน์ Stable ๋ก ์ ํํ ์ ์๋ค. ๊ฐ ๋ฉค๋ฒ๋ ์๊ธฐ SyncGroup ์๋ต์ ๋ฐ์ ๋ค ํ ๋น ์ฝ๋ฐฑ์ ๊ฑฐ์ณ ์๋น๋ฅผ ์ฌ๊ฐํ๋ค. Spring Kafka ์์๋ revoke ์ assign ์์ ์ ๋ฆฌ์ค๋๋ก ๋ฐ์ ๋ณผ ์ ์๋ค.
class RebalanceLogger : ConsumerAwareRebalanceListener {
override fun onPartitionsRevokedBeforeCommit(
consumer: Consumer<*, *>,
partitions: Collection<TopicPartition>,
) {
// eager ๋ฉด ๋ณด์ ํํฐ์
์ ์ฒด๊ฐ, cooperative ๋ฉด ์ด๋๋ถ๋ง ๋ด๊ฒจ ์จ๋ค
log.info("revoked: {}", partitions)
}
override fun onPartitionsAssigned(
consumer: Consumer<*, *>,
partitions: Collection<TopicPartition>,
) {
log.info("assigned: {}", partitions)
}
}
@Bean
fun kafkaListenerContainerFactory(
consumerFactory: ConsumerFactory<String, String>,
) = ConcurrentKafkaListenerContainerFactory<String, String>().apply {
this.consumerFactory = consumerFactory
containerProperties.setConsumerRebalanceListener(RebalanceLogger())
}์ ์์ ์ธ ์์ ๊ถ ์ด์ ์์ onPartitionsRevokedBeforeCommit ์ ์ฐํ๋ ๋ชฉ๋ก์ eager ๋ฉด ์ ์ฒด ํ ๋น, cooperative ๋ฉด ์ด๋ ๋์ ๋ถ๋ถ์งํฉ์ด๋ค. ์์ธ์ ์ผ๋ก ๋ฉค๋ฒ๊ฐ ํ์ฑ๋๊ฑฐ๋ ์์ ๊ถ์ ์ด๋ฏธ ์์ ๊ฒฝ์ฐ์๋ revoke ๊ฐ ์๋๋ผ lost ์ฝ๋ฐฑ ๊ฒฝ๋ก๋ฅผ ํ๋ฏ๋ก, ์ด ๋ก๊ทธ ํ๋๋ง์ผ๋ก ๋ชจ๋ ๋ฆฌ๋ฐธ๋ฐ์ค ์ํฉ์ ์ค๋ช
ํ ์๋ ์๋ค.
generation id ๋ ๋จ์ํ ์นด์ดํฐ๊ฐ ์๋๋ผ ํ์ฑ(fencing) ์ฅ์น๋ค. ๊ทธ๋ฃน ๊ด๋ฆฌ๋ก ๊ตฌ๋
ํ ๋ฉค๋ฒ์ OffsetCommit ์์ฒญ์๋ generation id ์ member id ๊ฐ ์ค๋ฆฌ๊ณ , ๊ตฌ์ธ๋ ๋ฉค๋ฒ๊ฐ ๋ค๋ฆ๊ฒ ๋ณด๋ธ ์ปค๋ฐ์ ILLEGAL_GENERATION ๊ฐ์ ์ค๋ฅ๋ก ๊ฑฐ๋ถ๋๋ค. ๋คํธ์ํฌ ์ง์ฐ์ผ๋ก ๋ค์ฒ์ง ์ข๋น ๋ฉค๋ฒ๊ฐ ์ด๋ฏธ ๋จ์ ๊ฒ์ด ๋ ํํฐ์
์ ์คํ์
์ ๋ฎ์ด์ฐ๋ ์ฌ๊ณ ๋ฅผ ๋ง์ ์ฃผ๋ ์
์ด๋ค.
2-2. classic protocol ์ ์ฒด ํ๋ก์ฐ
sequenceDiagram
participant A as Consumer A
participant B as Consumer B
participant GC as GroupCoordinator
participant N as Consumer C - ์ ๊ท
Note over GC: Stable
N->>GC: JoinGroup ์ ๊ท ํฉ๋ฅ
Note over GC: rebalance trigger โ PreparingRebalance
A->>GC: Heartbeat
GC-->>A: REBALANCE_IN_PROGRESS
B->>GC: Heartbeat
GC-->>B: REBALANCE_IN_PROGRESS
Note over A,B: eager ๋ผ๋ฉด ์ฌ๊ธฐ์ ์ ํํฐ์
revoke โ ์๋น ์ ์ง
A->>GC: JoinGroup ๊ตฌ๋
์ ๋ณด + assignor ๋ชฉ๋ก
B->>GC: JoinGroup ๊ตฌ๋
์ ๋ณด + assignor ๋ชฉ๋ก
Note over GC: join phase - ์ ์ ๋์ฐฉ๊น์ง ๋๊ธฐ = ๋ฐฐ๋ฆฌ์ด
Note over GC: generation++, A ๋ฅผ leader ๋ก ์ง๋ช
GC-->>A: JoinGroup ์๋ต - ์ ์์ ๊ตฌ๋
์ ๋ณด
GC-->>B: JoinGroup ์๋ต - ๋น ๊ฐ
GC-->>N: JoinGroup ์๋ต - ๋น ๊ฐ
Note over A: client-side assignment - leader ์ JVM ์์ ์คํ
A->>GC: SyncGroup ์ ์ฒด ํ ๋นํ
B->>GC: SyncGroup ๋น ์์ฒญ
N->>GC: SyncGroup ๋น ์์ฒญ
GC-->>A: ๋ด ๋ชซ๋ง
GC-->>B: ๋ด ๋ชซ๋ง
GC-->>N: ๋ด ๋ชซ๋ง
Note over GC: sync phase ์ข
๋ฃ - Stable, ํ์ฌ generation ์ ํ ๋น ํ์ 2-3. incremental cooperative rebalancing: ๋ผ์ด๋๋ฅผ ๋ ๋ฒ ๋๋ ์ด์
eager ์ cooperative ๋ ๊ฐ์ ๊ธฐ๊ณ ์์์ ๋๋ ์๋ก ๋ค๋ฅธ ์ ๋ต์ด๋ค. ๊ณต์ ์ด๋ฆ์ผ๋ก๋ ๊ฐ๊ฐ eager rebalancing ๊ณผ incremental cooperative rebalancing(KIP-429)์ด๋ผ๊ณ ๋ถ๋ฅธ๋ค. ์ฐจ์ด๋ "์ธ์ ํํฐ์ ์ ๋๋๋"์ ์๊ณ , ๊ทธ๋์ [[2-1-2|join phase]] ์ ๋ชจ์ต์ด ๋ฌ๋ผ์ง๋ค.
2-3-1. 1๋ผ์ด๋: ์ด๋ ๋์ ์๋ณ๊ณผ revoke
cooperative ๋ JoinGroup ์ ๋ณด๋ด๊ธฐ ์ ์ ์ ์ฒด ํํฐ์ ์ ๋ฏธ๋ฆฌ revoke ํ์ง ์๋๋ค. 1๋ผ์ด๋์ leader ๋ ์๋ํ ์ต์ข ํ ๋น์ ๊ณ์ฐํ๋, ์ด์ ์์ ์๊ฐ ์์ง ๋ค๊ณ ์๋ ์ด๋ ํํฐ์ ์ ์ ์์ ์์ ๊ฒฐ๊ณผ์์ ์ ์ธํ๋ค. ๊ธฐ์กด ์์ ์๋ SyncGroup ์๋ต์์ ๋น ์ง ํํฐ์ ๋ง revoke ํ๊ณ ์ฆ์ ๋ค์ ๋ฆฌ๋ฐธ๋ฐ์ค๋ฅผ ์์ฒญํ๋ค.
sequenceDiagram
participant A as Consumer A - p3 ๋ณด์
participant B as Consumer B
participant GC as GroupCoordinator
Note over A: ํํฐ์
์ ๋์ง ์์ ์ฑ ์ฐธ์ฌ
A->>GC: JoinGroup
B->>GC: JoinGroup
Note over A: leader ์ JVM ์์ ํ ๋น ๊ณ์ฐ
A->>GC: SyncGroup ํ ๋นํ - p3 ์ด๋์ ๋ณด๋ฅ
GC-->>A: ๋ด ๋ชซ์ p3 ์ด ๋น ์ ธ ์์
Note over A: SyncGroup ์๋ต ์ฒ๋ฆฌ ์ค p3 ๋ง revoke
GC-->>B: ๋ด ๋ชซ - ์์ง p3 ์ ๋น์ด ์์
Note over A,B: 1๋ผ์ด๋ ๋์๋ p3 ๋ง ์์ ๊ถ ๊ณต๋ฐฑ1๋ผ์ด๋๊ฐ ๋๋ ์์ ์ A ๋ p3 ์ ๋์์ง๋ง B ๋ ์์ง ๋ฐ์ง ๋ชปํ๋ค. KIP-429์ ์ด์ ์์ ์์ revoke ๋ฅผ ํ์ธํ๊ธฐ ์ ์ ์ ์์ ์์๊ฒ ๊ฐ์ ํํฐ์ ์ ์ฃผ์ง ์๋๋ค. p3 ์ด ๋๋ ์ด๋ค ๋ฉค๋ฒ์ owned partitions ์๋ ๋ํ๋์ง ์๋ ๋ค์ ๋ผ์ด๋์์์ผ B ์๊ฒ ํ ๋นํ ์ ์๋ค.
2-3-2. 2๋ผ์ด๋: ์ ์์ ์์๊ฒ ํ ๋น
์ด๋ ํํฐ์ ๋ง revoke ํ ๊ธฐ์กด ์์ ์๋ ์ฆ์ JoinGroup ์ ๋ค์ ๋ณด๋ด 2๋ผ์ด๋๋ฅผ ์์ํ๋ค. ์ด ๋ผ์ด๋์์ ์ฝ๋๋ค์ดํฐ์ leader ๊ฐ p3 ์ ๊ธฐ์กด ์์ ์๊ฐ ์์์ ํ์ธํ๋ฉด ์ ์ฃผ์ธ์๊ฒ ํ ๋นํ๋ค.
sequenceDiagram
participant A as Consumer A
participant B as Consumer B
participant GC as GroupCoordinator
Note over A: p3 ๋ง revoke, p1 p2 ๋ ๊ณ์ ์๋น
A->>GC: JoinGroup ์ฌํฉ๋ฅ ํธ๋ฆฌ๊ฑฐ
B->>GC: JoinGroup
A->>GC: SyncGroup ํ ๋นํ
GC-->>A: p1 p2
GC-->>B: p3 ์ถ๊ฐ
Note over B: p3 ์๋น ์์
Note over GC: Stable๋ ๋ผ์ด๋๋ฅผ ํ๋์ ํ๋จ ํ๋ฆ์ผ๋ก ์์ถํ๋ฉด ์๋์ ๊ฐ๋ค.
flowchart TD
S[rebalance trigger] --> R1[1๋ผ์ด๋ join phase/sync phase]
R1 --> D{์ด๋ ๋์ ํํฐ์
์๋}
D -- ์์ --> ST[Stable]
D -- ์์ --> RV[ํด๋น ๋ฉค๋ฒ๊ฐ ์ด๋๋ถ๋ง revoke<br/>๋๋จธ์ง ํํฐ์
์ ๊ณ์ ์๋น]
RV --> R2[2๋ผ์ด๋ join phase/sync phase]
R2 --> AS[์ ์ฃผ์ธ์๊ฒ ํ ๋น]
AS --> ST๋ผ์ด๋ ์๋ ๋์ง๋ง ๊ฐ ๋ผ์ด๋์์ ๋ฉ์ถ๋ ๊ฑด ์ด๋๋ถ๋ฟ์ด๋ผ, ์ฒด๊ฐ ์ ์ง๊ฐ ํฌ๊ฒ ์ค์ด๋ ๋ค. ์๋ณต ํ์์ ์ ์ง ๋ฒ์๋ฅผ ๋ง๋ฐ๊พผ ์ค๊ณ๋ค. ๊ทธ๋ ๋ค๊ณ ๋ฐฐ๋ฆฌ์ด ์์ฒด๊ฐ ์ฌ๋ผ์ง ๊ฑด ์๋์ด์, [[2-1-2|์์์ ๋ณธ]] "๊ฐ์ฅ ๋๋ฆฐ ๋ฉค๋ฒ๊ฐ ์ ์ฒด๋ฅผ ๋ถ์ก๋" ์ฑ์ง์ cooperative ์์๋ ๊ทธ๋๋ก ๋จ๋๋ค. ์คํ๋ ค ๋ผ์ด๋๊ฐ ๋ ๋ฒ์ด๋ ๊ทธ ๋ฐฐ๋ฆฌ์ด๋ฅผ ๋ ๋ฒ ํต๊ณผํ๋ค.
3. consumer ํ๋กํ ์ฝ(KIP-848)
classic ํ๋กํ ์ฝ์ ํ๋ก์ฐ์ ๋น๊ตํด๋ณด๋ฉด KIP-848์ด ์ด๋ค ๊ณผ์ ์ ์์ค ๊ฑด์ง ํ์ธํ ์ ์๋ค.. [[2-2|์ด์ด ๋ถ์ธ ๊ทธ๋ฆผ]]์ ์ธ๋ก์ค ์ ์ฒด, ๊ทธ๋ฌ๋๊น join phase ์ sync phase ๋ก ์ด๋ฃจ์ด์ง ๊ทธ๋ฃน ๋จ์ ์๋ณต์ด ํต์งธ๋ก ์ฌ๋ผ์ง๋ค.
Apache Kafka 4.3 Consumer Rebalance Protocol ๋ฌธ์๋ ์ด ํ๋กํ ์ฝ์ Kafka 4.0 ๋ถํฐ GA ๋ก ์ค๋ช
ํ๋ค. ๋ธ๋ก์ปค์์๋ ํ์ฑํ๋์ด ์์ง๋ง ์ปจ์๋จธ์ ๊ธฐ๋ณธ๊ฐ์ ์ฌ์ ํ classic ์ด๋ฏ๋ก, ์ค์ ์ฌ์ฉ์๋ group.protocol=consumer ์ค์ ์ด ํ์ํ๋ค.
3-1. ์ฌ๋ผ์ง ์ธ ๊ฐ์ง
- leader ์ปจ์๋จธ์ ๊ธฐ์กด client-side assignment ์์ โ Kafka 4.3 ์ ์ผ๋ฐ ์ปจ์๋จธ์์๋ ํน์ JVM ์ด ๊ทธ๋ฃน ์ ์ฒด์ ํ ๋น์ ์ฑ ์์ง์ง ์๋๋ค
- sync phase(SyncGroup) ์์ โ ๋ณ๋์ ๋ฐฐํฌ ์๋ณต์ด ์๋ค
- join phase ์ ๋ฐฐ๋ฆฌ์ด ์์ โ ์ ์์ด ๋ชจ์ผ ๋๊น์ง ๊ธฐ๋ค๋ฆฌ๋ ์ง์ ์์ฒด๊ฐ ์๋ค
์ Consumer ๊ทธ๋ฃน์์๋ JoinGroup/SyncGroup/Heartbeat ๊ฐ ํ๋ ๋ฉค๋ฒ์ญ/ํ ๋น/์์กด ํ์ธ์ ConsumerGroupHeartbeat ๊ฐ ๋งก๋๋ค. ์ฝ๋๋ค์ดํฐ๊ฐ server-side assignment ๋ก target assignment ๋ฅผ ๊ณ์ฐํ๊ณ , ๋ฉค๋ฒ๋ณ heartbeat ์๋ต์ ํ์ฌ ์ฌ์ฉํ ์ ์๋ ํ ๋น์ ์ฆ๋ถ์ผ๋ก ์ค์ด ๋ณด๋ธ๋ค. KIP-848 ์ ๊ฐ ๋ฉค๋ฒ๊ฐ ๊ทธ ๋ชฉํ๋ก ์ด๋ํ๋ ๊ณผ์ ์ reconciliation(์๋ ด)์ด๋ผ๊ณ ๋ถ๋ฅธ๋ค. ๋ฆฌ๋ฐธ๋ฐ์ค๊ฐ "ํ ๋ฒ์ ๊ทธ๋ฃน ์ด๋ฒคํธ"์์ "๋ฉค๋ฒ๋ณ ์ํ ์๋ ด"์ผ๋ก ๋ฐ๋ ๊ฒ์ด๋ค.
3-2. reconciliation ๋ฃจํ์ ๋จ๊ณ๋ณ ์ ๊ฐ
3-2-1. 0๋จ๊ณ: server-side assignment ์ target assignment
๋ฉค๋ฒ ํฉ๋ฅ/์ดํ, ๊ตฌ๋ ๋ณ๊ฒฝ, ํ ํฝ ๋ฉํ๋ฐ์ดํฐ ๋ณ๊ฒฝ์ผ๋ก group epoch ๊ฐ ์ฌ๋ผ๊ฐ๋ฉด ์ฝ๋๋ค์ดํฐ๊ฐ ์ ํ ๋น์ ๊ณ์ฐํ๋ค. ๋ฉค๋ฒ๋ heartbeat ๋ก ๊ตฌ๋ ๊ณผ ํ์ฌ ๋ณด์ ํํฐ์ ์ ๋ณด๊ณ ํ๊ณ , ์ฝ๋๋ค์ดํฐ๋ ์ ์ฅ๋ ๊ทธ๋ฃน ์ํ๋ฅผ ๊ทผ๊ฑฐ๋ก "์ด๋ ๊ฒ ๋์ด์ผ ํ๋ค"๋ ๋ชฉํ๋ฅผ ๋ง๋ ๋ค. ์ด ๋ชฉํ ์ํ์ ๊ณต์ ์ด๋ฆ์ด target assignment ๋ค.
flowchart LR
HB[๋ฉค๋ฒ๋ค์ ConsumerGroupHeartbeat<br/>ํ์ฌ ๋ณด์ ํํฐ์
๋ณด๊ณ ] --> GC[GroupCoordinator]
GC --> AS[server-side assignor<br/>uniform ๋๋ range]
AS --> T[target assignment<br/>๋๋ฌํด์ผ ํ ๋ชฉํ ์ํ]
T --> D[๊ฐ ๋ฉค๋ฒ์ heartbeat ์๋ต์<br/>์ฆ๋ถ ์ง์๋ฅผ ์ค์ด ๋ณด๋]classic ๊ณผ ๋น๊ตํ๋ฉด ์ฐจ์ด๊ฐ ํ๋์ ๋ณด์ธ๋ค. ๋งค ๋ฆฌ๋ฐธ๋ฐ์ค๋ง๋ค ๊ณ์ฐ ์ ๋ ฅ์ ๋ค์ ๋ชจ์ผ๋ ค๊ณ ์ ์์ด ๊ฐ์ JoinGroup ๋ผ์ด๋์ ๋์ฐฉํ ๋๊น์ง ๊ธฐ๋ค๋ฆฌ๋ ์ง์ ์ด ์๋ค. ์ฝ๋๋ค์ดํฐ๋ ์ง์์ ์ผ๋ก ๋ณด๊ดํ ๊ทธ๋ฃน ์ํ๋ก target assignment ๋ฅผ ๊ณ์ฐํ๊ณ , ๊ฐ ๋ฉค๋ฒ๋ฅผ ๋ ๋ฆฝ์ ์ผ๋ก ์๋ ด์ํจ๋ค.
3-2-2. 1๋จ๊ณ: ํ ๋น ์ฐจ์ด์ ๋ฐ๋ฅธ revoke
target assignment ์ ๋ฐ๋ฅด๋ฉด p3 ์ B ์ ๋ชซ์ธ๋ฐ ์ง๊ธ์ A ๊ฐ ๋ค๊ณ ์๋ค. ๊ทธ๋ฌ๋ฉด ์ฝ๋๋ค์ดํฐ๋ A ์๊ฒ๋ง p3 ์ด ๋น ์ง ํ ๋น์ ๋ณด๋ธ๋ค. ๋ณ๋์ revoke RPC๊ฐ ์๋ ๊ฒ์ ์๋๋ฉฐ, A ๊ฐ ํ์ฌ ๋ก์ปฌ ํ ๋น๊ณผ heartbeat ์๋ต์ ํ ๋น ์ฐจ์ด๋ฅผ ๊ณ์ฐํด p3 ์ revoke ํ๋ค.
sequenceDiagram
participant A as Consumer A - p1 p2 p3 ๋ณด์
participant B as Consumer B - p4 ๋ณด์
participant GC as GroupCoordinator
A->>GC: Heartbeat ๋ณด์ p1 p2 p3 ๋ณด๊ณ
GC-->>A: ํ ๋น ์๋ต p1 p2 - p3 ์ ์ธ
Note over A: p3 ๋ง ์ ๋ฆฌ, p1 p2 ๋ ๊ณ์ ์๋น
B->>GC: Heartbeat ๋ณด์ p4 ๋ณด๊ณ
GC-->>B: ์ง์ ์์
Note over B: p4 ๋ฅผ ํ ๋ฒ๋ ๋ฉ์ถ์ง ์๊ณ ์๋น ์ค์ฌ๊ธฐ์ B ์๊ฒ ๋ด๋ ค๊ฐ ์๋ต์ด "์ง์ ์์"์ด๋ผ๋ ์ ์ด ์ค์ํ๋ค. B ๋ ์๊ธฐ ๊ธฐ์กด ํํฐ์ ์ ๊ณ์ ์๋นํ๋ค. classic eager ๋ผ๋ฉด B ๋ ์ ์ฒด ํ ๋น์ revoke ํ๊ณ ๋ฐฐ๋ฆฌ์ด ์์ ์ฐ๊ฒ ์ง๋ง, Consumer ํ๋กํ ์ฝ์์๋ ํ ๋น์ด ๋ฐ๋์ง ์๋ ๋ฉค๋ฒ๋ฅผ ์ ์ญ ๋ผ์ด๋์ ์ธ์ฐ์ง ์๋๋ค.
3-2-3. 2๋จ๊ณ: revoke ์๋ฃ ๋ณด๊ณ ์ ์ ํ ๋น
A ๊ฐ p3 ์ ๋ฆฌ๋ฅผ ๋ง์น๋ฉด ์ฆ์ ๋ณด๋ด๋ ๋ค์ heartbeat ์ ๊ทธ ์ฌ์ค์ ๋ณด๊ณ ํ๋ค. ์ฝ๋๋ค์ดํฐ๋ ๊ทธ ๋ณด๊ณ ๋ฅผ ๋ฐ๊ณ ๋์์ผ B ์๊ฒ p3 ์ ์ฌ์ฉํ ์ ์๋ ํ ๋น์ผ๋ก ๋ด๋ ค์ค๋ค.
sequenceDiagram
participant A as Consumer A
participant B as Consumer B
participant GC as GroupCoordinator
A->>GC: Heartbeat p3 revoke ์๋ฃ ๋ณด๊ณ
Note over GC: p3 ์ ์์ ๊ถ ๊ณต๋ฐฑ ํ์ธ
B->>GC: Heartbeat
GC-->>B: ํ ๋น ์๋ต p3 p4
Note over B: p4 ์ ์ด์ด p3 ์๋น ์์, reconciliation ์๋ฃ์ฃผ๋ชฉํ ์ ์ B ๊ฐ A ๋ฅผ ๊ธฐ๋ค๋ฆฐ ๊ฒ ์๋๋ผ, ์ฝ๋๋ค์ดํฐ๊ฐ B ์๊ฒ ์์ง ์๋ ค์ฃผ์ง ์์์ ๋ฟ์ด๋ผ๋ ๊ฒ์ด๋ค. B ๋ ๊ทธ๋์ ์๊ธฐ ํํฐ์ ์ ์ ์์ ์ผ๋ก ์๋นํ๊ณ ์์๋ค. ํํฐ์ ํ๋์ ๊ด์ ์์ ์ํ ์ด๋์ ๊ทธ๋ฆฌ๋ฉด ์ ์ง ๋ฒ์๊ฐ ์ด๋๊น์ง์ธ์ง๊ฐ ๋ถ๋ช ํด์ง๋ค.
stateDiagram-v2
[*] --> A์์ : ์ด๊ธฐ ์ํ
A์์ --> ๊ณต๋ฐฑ: A ๊ฐ p3 ๋ง revoke
๊ณต๋ฐฑ --> B์์ : ์ฝ๋๋ค์ดํฐ๊ฐ B ์๊ฒ assign
B์์ --> [*]: reconciliation ์๋ฃ
note right of ๊ณต๋ฐฑ
๋ฉ์ถฐ ์๋ ๊ฑด p3 ํ๋๋ฟ์ด๋ค.
p1 p2 p4 ๋ ๋ด๋ด ์๋น ์ค์ด๋ค.
end noteclassic eager ์์๋ ๋ชจ๋๊ฐ ์ ์ฒด ํ ๋น์ ๋์์ง๋ง, ์ฌ๊ธฐ์๋ ์ ์์ ์ธ ์์ ๊ถ ์ด์ ๊ธฐ์ค์ผ๋ก ์ด๋ํ๋ ํํฐ์
ํ๋๋ง ์ ๊น ๋น์ด ์๊ฒ ๋๋ค. ์ ์ง์ ๋จ์๊ฐ ๊ทธ๋ฃน์์ ํํฐ์
์ผ๋ก ๋ด๋ ค์จ ์
์ด๋ค. ํ์ฑ์ด๋ ์ธ์
๋ง๋ฃ์ฒ๋ผ ์์ ๊ถ์ ์๋ ์์ธ ๊ฒฝ๋ก์์๋ ๋ฉค๋ฒ๊ฐ ์ ์ฒด ํํฐ์
์ abandon ํ๊ณ onPartitionsLost๋ฅผ ํธ์ถํ ์ ์๋ค.
3-2-4. ์ด์ด ๋ถ์ธ ํ ๋ฒ์ reconciliation
์ธ ๋จ๊ณ๋ฅผ ํ ์ฅ์ผ๋ก ๋ถ์ด๋ฉด ์๋์ ๊ฐ๋ค.
sequenceDiagram
participant A as Consumer A
participant B as Consumer B
participant GC as GroupCoordinator
Note over GC: server-side assignor ๊ฐ target assignment ๊ณ์ฐ
A->>GC: Heartbeat ํ์ฌ ๋ณด์ ํํฐ์
๋ณด๊ณ
GC-->>A: ํ ๋น ์๋ต p1 p2 - p3 ์ ์ธ
Note over A: p3 ๋ง ์ ๋ฆฌ, p1 p2 ๋ ๊ณ์ ์๋น
B->>GC: Heartbeat
GC-->>B: ๋๊ธฐ - A ๊ฐ ์์ง ์ ๋์
A->>GC: Heartbeat p3 revoke ์๋ฃ
B->>GC: Heartbeat
GC-->>B: ํ ๋น ์๋ต p3 p4
Note over B: p3 ์๋น ์์, reconciliation ์๋ฃ[[2-2|classic ์ ์ ์ฒด ๊ทธ๋ฆผ]]๊ณผ ๋๋ํ ๋๊ณ ๋ณด๋ฉด ์ฐจ์ด๊ฐ ํํ๋ก ๋๋ฌ๋๋ค. classic eager ๋ ๋ชจ๋ ๋ฉค๋ฒ๊ฐ ๊ฐ์ ๋ผ์ด๋์์ ์ ์ฒด ํ ๋น์ ๋๊ณ , classic cooperative ๋ revoke ๋ฒ์๋ ์ค์ง๋ง ์ ์ญ ๋ฐฐ๋ฆฌ์ด๋ ํต๊ณผํ๋ค. Consumer ํ๋กํ ์ฝ์์๋ ๋ฉค๋ฒ๋ง๋ค heartbeat ๋ก ๊ฐ์ ์๋ ดํ๋ฏ๋ก ์ธ๋ก๋ก ์ ๋ ฌ๋ ๊ทธ๋ฃน ๋ฐฐ๋ฆฌ์ด ์ ์ด ์๋ค.
[!NOTE]
KIP-848 ์ group epoch ๋ ๊ตฌ๋ ์ด๋ ๋ฉค๋ฒ์ญ์ด ๋ฐ๋์ด ์ target assignment ๊ฐ ํ์ํ ์์ ์ ๋ํ๋ด๊ณ , member epoch ๋ ๊ฐ ๋ฉค๋ฒ๊ฐ ๊ทธ ๋ชฉํ์ ์ด๋๊น์ง ์๋ ดํ๋์ง๋ฅผ ๋ํ๋ธ๋ค. ์์ฒญ ํ์ฑ๊ณผOffsetCommit๊ฒ์ฆ์๋ member epoch ๊ฐ ์ฌ์ฉ๋๋ค. [[2-1-4|์์์ ๋ณธ]] ์ ์ญ generation ๊ธฐ๋ฐ ๊ฒ์ฆ์ ๋ฉค๋ฒ๋ณ ์งํ ์ํ๋ก ์ธ๋ถํํ ๊ตฌ์กฐ๋ค.
3-3. server-side assignor
ํ ๋น ๊ณ์ฐ์ด ๋ธ๋ก์ปค๋ก ์ฎ๊ฒจ์์ผ๋ ์ค์ ํค๋ ๋ฐ๋๋ค.
@Bean
fun consumerFactory(): ConsumerFactory<String, String> = DefaultKafkaConsumerFactory(
mapOf(
ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG to "kafka:9092",
ConsumerConfig.GROUP_ID_CONFIG to "settlement-worker",
"group.protocol" to "consumer",
// ์ด assignor ๋ ๋ธ๋ก์ปค์์ ์คํ๋๋ค. uniform ๋๋ range
"group.remote.assignor" to "uniform",
),
)remote ๋ผ๋ ์ด๋ฆ์ด ๊ด์ ์ ๊ทธ๋๋ก ๋๋ฌ๋ธ๋ค -- ์ปจ์๋จธ ์
์ฅ์์ ์๊ฒฉ(๋ธ๋ก์ปค)์ ์๋ assignor ๋ฅผ ์ง์ ํ๋ ๊ฒ์ด๊ธฐ ๋๋ฌธ์ด๋ค. [[2-1-3]]์ ์ค์ ์ด ํด๋์ค ์ด๋ฆ์ด์๋ ๊ฒ๊ณผ ๋ฌ๋ฆฌ, ์ฌ๊ธฐ์๋ ๋ธ๋ก์ปค๊ฐ ์์๋ฃ๋ ๋ฌธ์์ด ์๋ณ์๋ฅผ ๋๊ธด๋ค๋ ์ ๋ ๋๋น๋๋ค.
4. classic ๊ณผ consumer ๋ฅผ ๋๋ํ ๋์ผ๋ฉด
์์ ํ๋ฆ์ ๊ฐ์ ๊ธฐ์ค์ผ๋ก ๋ง์ถฐ ๋์ผ๋ฉด KIP-848 ์ด ๋ฌด์์ ๊ฐ์ ํ๋์ง๊ฐ ๋ถ๋ช ํด์ง๋ค. cooperative ๋ classic ์ ์ ์ง ๋ฒ์๋ฅผ ์ค์๊ณ , Consumer ํ๋กํ ์ฝ์ ์ ์ง๋ฅผ ๋ง๋ค์ด ๋ด๋ ๊ทธ๋ฃน ๋จ์ ๋๊ธฐํ ๊ตฌ์กฐ๋ฅผ ๋ฐ๊ฟจ๋ค.
| ๋น๊ต์ถ | classic eager | classic cooperative | consumer(KIP-848) |
|---|---|---|---|
| ํ ๋น ๊ณ์ฐ ์์น | leader ์ปจ์๋จธ | leader ์ปจ์๋จธ | GroupCoordinator |
| ํต์ฌ ์์ฒญ | JoinGroup โ SyncGroup | JoinGroup โ SyncGroup | ConsumerGroupHeartbeat |
| ๋๊ธฐํ ๋จ์ | ๊ทธ๋ฃน ์ ์ฒด ๋ฐฐ๋ฆฌ์ด | ๊ทธ๋ฃน ์ ์ฒด ๋ฐฐ๋ฆฌ์ด ์ ์ง, ์ด๋์ด ์์ผ๋ฉด ์ฐ์ ๋ผ์ด๋ | ๋ฉค๋ฒ๋ณ reconciliation |
| ์ ์ ์ฌํ ๋น์ revoke ๋ฒ์ | ๋ณด์ ํํฐ์ ์ ์ฒด | ์ด๋ ๋์ ํํฐ์ ๋ง | ์ด๋ ๋์ ํํฐ์ ๋ง |
| ์งํ ์ํ | generation ๋จ์ | generation ๋จ์ | group epoch/member epoch ๋จ์ |
| ์ํฅ์ ๋ฐ๋ ๋ฒ์ | ๊ทธ๋ฃน ์ ์ฒด | ์ด๋ ํํฐ์ ์ ์ค์ง๋ง ๋๋ฆฐ ๋ฉค๋ฒ์ ๋ฐฐ๋ฆฌ์ด ์ํฅ์ ๋จ์ | ์ด๋ ํํฐ์ ๊ณผ ๊ด๋ จ ๋ฉค๋ฒ, ๋๋จธ์ง๋ ๊ณ์ ์๋น |
ํ์์ ์ค์ํ ํ์ revoke ๋ฒ์๊ฐ ์๋๋ผ ๋๊ธฐํ ๋จ์๋ค. cooperative ์ Consumer ๋ชจ๋ ์ด๋๋ถ๋ง ๋ด๋ ค๋์ง๋ง, cooperative ๋ ์ ์์ด ๊ฐ์ ๋ผ์ด๋๋ฅผ ํต๊ณผํด์ผ ํ๋ค. Consumer ํ๋กํ ์ฝ์ ๋ฉค๋ฒ๋ง๋ค ํ์ฌ ํ ๋น์ ๋ณด๊ณ ํ๊ณ ๊ฐ์ ๋ชฉํ ํ ๋น์ผ๋ก ์๋ ดํ๋ฏ๋ก, ๋๋ฆฐ ๋ฉค๋ฒ ํ๋๊ฐ ๊ทธ๋ฃน ์ ์ฒด์ ์ถ๋ฐ ์์ ์ ๊ฒฐ์ ํ๋ ๊ตฌ์กฐ๊ฐ ์๋ค.
4-1. cooperative ๋ consumer ํ๋กํ ์ฝ์ ์ด์ ์ด๋ฆ์ด ์๋๋ค
๋์ ๊ฐ์ ๋ฌธ์ ๋ฅผ ์๋ก ๋ค๋ฅธ ์ธต์์ ํผ๋ค. KIP-429์ incremental cooperative rebalancing ์ classic ํ๋กํ ์ฝ ์์์ ํํฐ์ ์ ์ธ์ ๋์์ง๋ฅผ ๋ฐ๊พผ๋ค. 1๋ผ์ด๋์ SyncGroup ๊ฒฐ๊ณผ์์ ๊ธฐ์กด ์์ ์๊ฐ ์ด๋ ํํฐ์ ๋ง revoke ํ๊ณ , ๋ค์ ๋ผ์ด๋์์ ์ ์์ ์์๊ฒ ๋๊ธด๋ค. JoinGroup, leader ์ client-side assignment, SyncGroup, generation ์ ๊ทธ๋๋ก ๋จ๋๋ค.
KIP-848์ ๊ทธ ๋ผ์ด๋๋ฅผ ๋๊ฐ ์กฐ์จํ ์ง๋ฅผ ๋ฐ๊พผ๋ค. ์ฝ๋๋ค์ดํฐ๊ฐ target assignment ๋ฅผ ๋ง๋ค๊ณ , ๋ฉค๋ฒ๋ณ heartbeat ์๋ต์ ํ ๋น ์ฐจ์ด๋ก revoke ์ assign ์ ์ ๋ํ๋ค. cooperative ๊ฐ ์๋ณต ํ์์ ์ ์ง ๋ฒ์๋ฅผ ๋ง๋ฐ๊ฟจ๋ค๋ฉด, Consumer ํ๋กํ ์ฝ์ ์๋ณต์ ์ฑ๋ฆฝ์ํค๋ ์ ์ญ ๋ฐฐ๋ฆฌ์ด๋ฅผ ์ ๊ฑฐํ ๊ฒ์ด๋ค. KIP-848 ์์๋ cooperative ํ๋กํ ์ฝ์กฐ์ฐจ ๊ทธ๋ฃน ๋จ์ ๋ฐฐ๋ฆฌ์ด์ ์์กดํ๋ค๋ ์ ์ ๊ธฐ์กด ๊ตฌ์กฐ์ ํ๊ณ๋ก ๋ช ์ํ๋ค.
4-2. ์ค์ ์ ์์ ๊ถ๋ ๋ธ๋ก์ปค๋ก ์ด๋ํ๋ค
์คํ ์์น๊ฐ ๋ฌ๋ผ์ง๋ฉด ๊ฐ์ ์ด๋ฆ์ ์ค์ ์ ๊ทธ๋๋ก ์ฎ๊ธธ ์ ์๋ค. Kafka 4.3 ๊ณต์ ๋ฌธ์๋ Consumer ํ๋กํ ์ฝ์์ heartbeat.interval.ms, session.timeout.ms, partition.assignment.strategy๋ฅผ ์ฌ์ฉํ ์ ์๋ค๊ณ ์ ๋ฆฌํ๋ค.
| classic ์์์ ์ฑ ์ | Consumer ํ๋กํ ์ฝ์์์ ์ฑ ์ |
|---|---|
ํด๋ผ์ด์ธํธ์ heartbeat.interval.ms |
๋ธ๋ก์ปค์ group.consumer.heartbeat.interval.ms |
ํด๋ผ์ด์ธํธ์ session.timeout.ms |
๋ธ๋ก์ปค์ group.consumer.session.timeout.ms |
ํด๋ผ์ด์ธํธ์ partition.assignment.strategy |
๋ธ๋ก์ปค์ group.consumer.assignors |
| leader JVM ์ด assignor ํด๋์ค ์คํ | ์ฝ๋๋ค์ดํฐ๊ฐ server-side assignor ์คํ |
| ๊ฐ ํด๋ผ์ด์ธํธ๊ฐ assignor ์ฐ์ ์์ ์ ๊ณ | ํด๋ผ์ด์ธํธ๋ ํ์ํ ๋ group.remote.assignor๋ก ์๋ฒ assignor ์ด๋ฆ๋ง ์ ํ |
๊ธฐ๋ณธ server-side assignor ๋ uniform ๊ณผ range ๋ค. ๊ณต์ ๋ง์ด๊ทธ๋ ์ด์
ํ๋ RangeAssignor๋ฅผ range ๋ก, RoundRobinAssignor/StickyAssignor/CooperativeStickyAssignor๋ฅผ uniform ์ผ๋ก ๋์์ํจ๋ค. uniform ๊ณผ range ๋ชจ๋ ๊ธฐ์กด ํ ๋น์ ์ต๋ํ ์ ์งํ๋ sticky ์ฑ์ง์ ๊ฐ์ง๋ฏ๋ก, Consumer ํ๋กํ ์ฝ ์์ CooperativeStickyAssignor๋ฅผ ํ ๊ฒน ๋ ์น๋ ๊ตฌ์กฐ๋ ์ฑ๋ฆฝํ์ง ์๋๋ค. partition.assignment.strategy ์์ฒด๊ฐ ์ด ํ๋กํ ์ฝ์์๋ ์คํ๋์ง ์๋๋ค.
4-3. ์ ํ ์ ์ ํ์ธํ ๊ฒฝ๊ณ
Consumer ํ๋กํ ์ฝ์ Kafka 4.0 ๋ถํฐ GA ์ง๋ง, ์ค์ ํ๋๋ง ๋ฐ๊พผ๋ค๊ณ ๋ชจ๋ ๊ทธ๋ฃน์ด ๊ฐ์ ์กฐ๊ฑด์ผ๋ก ์ ํ๋๋ ๊ฒ์ ์๋๋ค.
- ๋ธ๋ก์ปค์ ํด๋ผ์ด์ธํธ ๋ฒ์ : ๋ธ๋ก์ปค๋ฟ ์๋๋ผ ์ปจ์๋จธ๋ Consumer ํ๋กํ ์ฝ์ ์ง์ํด์ผ ํ๋ค. Java client ๋ Kafka 4.0 ๋ถํฐ ์ ์ ์ง์ํ๋ค.
- ์จ๋ผ์ธ ์ ํ ์กฐ๊ฑด: rolling ๋ฐฉ์์ ๋ฌด์ค๋จ ์ ํ์ ๋ธ๋ก์ปค์
group.consumer.migration.policy๊ฐ upgrade ๋ฅผ ํ์ฉํ๊ณ , classic ๊ทธ๋ฃน์ assignor ๊ฐ ์ปค์คํ ๋ฉํ๋ฐ์ดํฐ๋ฅผ ํ ๋น ์ ๋ณด์ ์ฌ์ง ์๋ ๊ฒฝ์ฐ์ ์ง์๋๋ค. ์กฐ๊ฑด์ ๋ง์กฑํ์ง ์์ผ๋ฉด ๊ทธ๋ฃน์ ๋น์ด ๋ค ํ๋กํ ์ฝ์ ๋ฐ๊พธ๋ offline ์ ํ์ด ํ์ํ๋ค. - ์ปค์คํ
ํ ๋น ๋ก์ง: Kafka 4.3 ๊ธฐ์ค client-side assignor ๋ Consumer ํ๋กํ ์ฝ์์ ์ง์๋์ง ์๋๋ค. ๋ธ๋ก์ปค์
ConsumerGroupPartitionAssignor๊ตฌํ์ ๋ฐฐํฌํ๋ server-side ํ์ฅ์ ๊ฐ๋ฅํ์ง๋ง, ํด๋ผ์ด์ธํธ๋ง๋ค ์ฝ๋๋ฅผ ์ฃ๋ ๊ธฐ์กด ๋ฐฉ์๊ณผ ์ด์ ๊ฒฝ๊ณ๊ฐ ๋ค๋ฅด๋ค. - ์ค์ ์์ ๊ถ: heartbeat ์ session timeout ์ด ๋ธ๋ก์ปค ์ค์ ์ผ๋ก ์ด๋ํ๋ค. ํน์ ์ ํ๋ฆฌ์ผ์ด์ ๋ง์ ๊ฐ์ ์กฐ์ ํ๋ ์ด์ ๋ฐฉ์์ด๋ผ๋ฉด ์ ํ ์ ์ ํด๋ฌ์คํฐ/๊ทธ๋ฃน ์ค์ ์ ์ฑ ์ ํจ๊ป ์ ๋ฆฌํด์ผ ํ๋ค.
๋ฐ๋ผ์ ์ ํ ํ๋จ์ "CooperativeStickyAssignor ๋ณด๋ค ๋น ๋ฅธ๊ฐ"๋ก ๋๋์ง ์๋๋ค. ์ปค์คํ ๋ฉํ๋ฐ์ดํฐ์ assignor ์ ์กด์ฌ ์ฌ๋ถ, timeout ์ค์ ์ ์์ ์ฃผ์ฒด, rolling migration ๊ฐ๋ฅ์ฑ๊น์ง ํ์ธํด์ผ ํ๋ค. ์ด ์กฐ๊ฑด์ ํต๊ณผํ๋ฉด Consumer ํ๋กํ ์ฝ์ cooperative ์ ์ ์ง ๋ฒ์ ์ถ์๋ฅผ ์ ์งํ๋ฉด์ classic ์ ์ ์ญ ๋ฐฐ๋ฆฌ์ด๊น์ง ์ ๊ฑฐํ๋ค.
๋ง๋ฌด๋ฆฌ
- classic ์ ๋ฆฌ๋ฐธ๋ฐ์ค๋ ๊ทธ๋ฃน ๋จ์ ์๋ณต์ด๋ค. rebalance trigger(heartbeat ์๋ต์ผ๋ก ๊ฐ์ ํต๋ณด) โ join phase(์ ์ JoinGroup ๋ฐฐ๋ฆฌ์ด) โ client-side assignment(leader ์ JVM) โ sync phase(SyncGroup) ์์ผ๋ก ํ ๋ฐํด๋ฅผ ๋๋ค. ์ด ์ค ํ๋กํ ์ฝ์ด phase ๋ก ์ด๋ฆ ๋ถ์ธ ๊ฒ์ ๊ฐ์ด๋ฐ ๋๋ฟ์ด๋ค.
- eager ๋ฆฌ๋ฐธ๋ฐ์ค๋ฅผ ๊ธธ๊ฒ ๋ง๋๋ ํต์ฌ ๋๊ธฐ ์ง์ ์ join phase ์ ๋ฐฐ๋ฆฌ์ด๋ค. ๊ฐ์ฅ ๋ฆ๊ฒ ์ฌํฉ๋ฅํ๋ ๋ฉค๋ฒ๊ฐ ๋ผ์ด๋ ์๋ฃ ์์ ์ ๊ฒฐ์ ํ ์ ์๊ธฐ ๋๋ฌธ์ด๋ค. ๋ค๋ง trigger ์ธ์ง ์ ์๋ ๊ธฐ์กด ํ ๋น์ ๊ณ์ ์๋นํ ์ ์์ผ๋ฏ๋ก ์ฝ๋๋ค์ดํฐ์
PreparingRebalance์ฒด๋ฅ ์๊ฐ ์ ์ฒด๊ฐ ๊ณง ์๋น ์ ์ง ์๊ฐ์ ์๋๋ค. - incremental cooperative rebalancing ์ ์ด ๊ธฐ๊ณ๋ฅผ ๊ทธ๋๋ก ์ฐ๋ "1๋ผ์ด๋ ๋์ ์ด๋๋ถ๋ง revoke ํ๊ณ , ๋ค์ ๋ผ์ด๋์ ์ ์์ ์์๊ฒ ๋๊ธด๋ค"๋ก ์ ์ง ๋ฒ์๋ฅผ ์ค์ธ ์ ๋ต์ด๋ค. ์๋ณต ํ์๋ฅผ ๋๋ ค ์ ์ง ๋ฒ์๋ฅผ ์ฐ ์ ์ด๋ค.
- KIP-848 ์ Consumer ํ๋กํ ์ฝ์ ๊ทธ๋ฃน ์ ์ฒด๊ฐ ํจ๊ป ๋๋ ์๋ณต์ ์์ ๊ณ , target assignment ๊ณ์ฐ โ ํ ๋น ์ฐจ์ด์ ๋ฐ๋ฅธ revoke โ ์๋ฃ ๋ณด๊ณ โ ์ ํ ๋น์ผ๋ก ์ด์ด์ง๋ reconciliation์ผ๋ก ๋์ฒดํ๋ค. ์ ์ ์ฌํ ๋น์์ ๋ฉ์ถ๋ ๋จ์๊ฐ ๊ทธ๋ฃน ์ ์ฒด์์ ์ด๋ ํํฐ์ ์ผ๋ก ๋ด๋ ค์จ๋ค.
group.protocol=consumer์์๋partition.assignment.strategy๋ฅผ ์ฌ์ฉํ ์ ์์ผ๋ฏ๋กCooperativeStickyAssignor๋ฅผ ๊ฒน์ณ ์ ์ฉํ ์ ์๋ค. ์ ์ ์ฌํ ๋น์ incremental ๋์์ Consumer ํ๋กํ ์ฝ์ด ๋ด๋นํ๊ณ , ๊ธฐ๋ณธ server-side uniform/range assignor ๋ ๊ธฐ์กด ํ ๋น์ ์ต๋ํ ์ ์งํ๋ sticky ์ฑ์ง์ ์ ๊ณตํ๋ค.
์ถ์ฒ
- Kafka Client-side Assignment Proposal โ join phase ์ sync phase, leader ์ ํ ๋น ๊ณ์ฐ
- KIP-62: Allow consumer to send heartbeats from a background thread โ
max.poll.interval.ms, rebalance timeout, heartbeat thread - KIP-345: Introduce static membership protocol โ ์ ์ ๋ฉค๋ฒ์ session timeout ๋์
- KIP-429: Kafka Consumer Incremental Rebalance Protocol โ incremental cooperative rebalancing
- KIP-848: The Next Generation of the Consumer Rebalance Protocol โ server-side assignment, target assignment, reconciliation, epoch
- Apache Kafka 4.3 Consumer Configs โ classic ๊ธฐ๋ณธ๊ฐ, assignor ๋ชฉ๋ก, timeout ์ค์
- Apache Kafka 4.3 ConsumerRebalanceListener โ ๋ฆฌ๋ฐธ๋ฐ์ค trigger ์ revoke ๋ฒ์
- Apache Kafka 4.3 CooperativeStickyAssignor โ cooperative ํ์ฑํ ์กฐ๊ฑด
- Apache Kafka 4.3 Consumer Rebalance Protocol โ GA ๋ฒ์ , ์ค์ ๋ณํ, assignor ๋์, ์ ํ ์กฐ๊ฑด๊ณผ ์ ํ์ฌํญ
- Apache Kafka 4.3 Broker Configs โ Consumer ํ๋กํ ์ฝ migration policy
'Infra > Kafka' ์นดํ ๊ณ ๋ฆฌ์ ๋ค๋ฅธ ๊ธ
| Kafka ํํฐ์ ์ฆ๊ฐ ์์ด ์ฒ๋ฆฌ๋ ๋๋ฆฌ๊ธฐ (feat. Parallel Consumer, Share Group) (0) | 2026.08.29 |
|---|