배경
PR #511 Copilot 피드백 — 이슈 #510 의 후속.
PR #511 의 BufferDrainer 는 전역 PUBLISHER_REDIS_BUFFER_ENABLED flag 만 보고 시작. 같은 바이너리를 여러 stage/인스턴스로 띄우면 각 instance 가 동일 Kafka backlog 값을 읽고 각자 DrainBatch 만큼 drain → 인스턴스 수만큼 Kafka 로 재주입되어 TargetBacklog 초과 가능.
현 시점 영향: 단일 인스턴스 운영에서는 무영향. 다중 인스턴스 / 수평 확장 시 (이슈 #443 stage toggle 로 fetcher-only 노드 여러 개 띄울 때) 발생.
현상 (다중 인스턴스 시나리오)
Instance A: backlog=2000 < target=3000 → drain 100 → publish 100 → backlog +100
Instance B: backlog=2000 < target=3000 → drain 100 → publish 100 → backlog +100 (B 는 A 의 publish 를 못 봄)
Instance C: backlog=2000 < target=3000 → drain 100 → publish 100 → backlog +100
...
실제 backlog: 2000 + 100*N — target_backlog 의 N배 초과
Redis RPOP 자체는 atomic — 같은 항목이 두 instance 에 동시 pop 되지는 않음. 그러나 각 instance 가 자기 drainBatch 만큼은 publish 하므로 total Kafka rate 가 N배 증가.
변경 방향
옵션 A — Redis 기반 leader election (권장)
drainer goroutine 의 매 tick 시작 시 Redis 분산 락 (SET NX EX) 으로 leader 인지 확인. leader 만 drain 실행.
// drainer tick:
acquired := redis.SET("publisher:buffer:drainer:leader", instanceID, "NX", "EX", "31s")
if !acquired { return } // 다른 instance 가 leader
// drain logic
- TTL: drain interval + 1 (예: 30s interval → TTL 31s)
- instance ID 는 hostname + pid + random (재기동 시 새 ID)
- leader 가 죽으면 TTL 만료 후 다음 instance 가 자동 leader 승격
- 기존
pkg/redis/lock.go 의 AcquireLock 패턴 재사용 가능
옵션 B — Token bucket per-target_backlog
drain 시점에 Redis 의 atomic counter (DECR) 로 token 차감. token 0 이면 drain skip. period 마다 counter 보충 (SET token target_backlog).
복잡도 높음, leader election 보다 일관성 약함.
옵션 C — Stage role gating
scheduler stage 와 drainer 를 결합 — STAGES_SCHEDULER_ENABLED=true 인 instance 만 drainer 실행. scheduler 가 보통 1대만 띄우므로 자연 leader.
단점: scheduler stage 와 drainer 의 책임 결합 — 향후 scheduler 비활성 환경에서 drainer 도 함께 비활성 (별도 fetcher-only 클러스터)
완료 조건
Why
우선순위
Low — 단일 인스턴스 운영에서는 무영향. 다중 인스턴스 운영 시작 직전에 처리 권장.
관련
배경
PR #511 Copilot 피드백 — 이슈 #510 의 후속.
PR #511 의
BufferDrainer는 전역PUBLISHER_REDIS_BUFFER_ENABLEDflag 만 보고 시작. 같은 바이너리를 여러 stage/인스턴스로 띄우면 각 instance 가 동일 Kafka backlog 값을 읽고 각자DrainBatch만큼 drain → 인스턴스 수만큼 Kafka 로 재주입되어TargetBacklog초과 가능.현 시점 영향: 단일 인스턴스 운영에서는 무영향. 다중 인스턴스 / 수평 확장 시 (이슈 #443 stage toggle 로 fetcher-only 노드 여러 개 띄울 때) 발생.
현상 (다중 인스턴스 시나리오)
Redis RPOP자체는 atomic — 같은 항목이 두 instance 에 동시 pop 되지는 않음. 그러나 각 instance 가 자기 drainBatch 만큼은 publish 하므로 total Kafka rate 가 N배 증가.변경 방향
옵션 A — Redis 기반 leader election (권장)
drainer goroutine 의 매 tick 시작 시 Redis 분산 락 (
SET NX EX) 으로 leader 인지 확인. leader 만 drain 실행.pkg/redis/lock.go의AcquireLock패턴 재사용 가능옵션 B — Token bucket per-target_backlog
drain 시점에 Redis 의 atomic counter (
DECR) 로 token 차감. token 0 이면 drain skip. period 마다 counter 보충 (SET token target_backlog).복잡도 높음, leader election 보다 일관성 약함.
옵션 C — Stage role gating
scheduler stage 와 drainer 를 결합 —
STAGES_SCHEDULER_ENABLED=true인 instance 만 drainer 실행. scheduler 가 보통 1대만 띄우므로 자연 leader.단점: scheduler stage 와 drainer 의 책임 결합 — 향후 scheduler 비활성 환경에서 drainer 도 함께 비활성 (별도 fetcher-only 클러스터)
완료 조건
BufferDrainer.acquireLeader(ctx)메소드 추가 — Redis SET NX EX 기반Why
kafka backlog exceeds thresholddrop 해소 충분우선순위
Low — 단일 인스턴스 운영에서는 무영향. 다중 인스턴스 운영 시작 직전에 처리 권장.
관련
pkg/redis/lock.go— 기존 분산 락 패턴 (재사용 후보)