Redis Stream Coordinator
Design docs

Member Data-Plane Boundary

Redis Stream Coordinator는 message processor가 아니다. 이 문서는 coordinator와 member data-plane의 경계만 정의한다.

Coordinator Owns

Member Data-Plane Owns

At-Least-Once Processing Model

이 프로젝트의 기본 처리 모델은 at-least-once이다. Redis Stream consumer group delivery는 XREADGROUP/XACK 기준으로 중복 전달될 수 있고, consumer handler도 장애, pending recovery, rebalance 이후 다시 실행될 수 있다.

단일 처리 보장은 제공하지 않는다. 실제 application handler는 하나의 message를 처리하면서 DB update, Redis write, HTTP call, external API call, local cache update 등 여러 비즈니스 side effect를 함께 수행할 수 있다. 이 프로젝트는 그런 side effect들과 Redis Stream ACK를 하나의 원자적 transaction으로 묶는 protocol을 제공하지 않는다.

Application 책임:

권장 처리 흐름:

XREADGROUP
  -> memberEpoch / assignmentEpoch / shard ownership 확인
  -> business handler 실행
  -> application-level duplicate guard 또는 retry policy 적용
  -> XACKDEL 또는 XACK

보장 경계:

Contract

Coordinator heartbeat response의 assignment.assignedShards에서 기존 owned shard가 빠지면 member는 해당 shard의 신규 read를 중단하고 local in-flight가 0이 된 뒤 heartbeat의 revokingShards.state=REVOKED로 보고한다.

OK response의 assignment.assignedShards에 새 shard가 포함되면 member는 shard를 ownedShards에 반영하고, 실제 Redis Stream read/recovery는 member data-plane 정책에 따라 수행한다.

Metadata correction 중 SYNC_METADATAREVOKE_PENDING은 drain-only response이다. Member는 이미 소유한 shard 중 assignment.assignedShards에 남은 shard만 유지할 수 있고, 처음 등장한 shard는 이후 OK response를 받을 때까지 read를 시작하면 안 된다.

Coordinator가 UNKNOWN_MEMBER_ID, FENCED_MEMBER_EPOCH을 반환하거나 member epoch mismatch가 발생하면 member는 read/ack를 중단하고 local ownership/revoke state를 비운 뒤, 같은 memberIdmemberEpoch=0 full heartbeat로 rejoin한다.

Consumer module은 UNKNOWN_MEMBER_ID 이후 stale ownedShards를 계속 보고하면 안 된다. Rejoin은 empty ownership report에서 시작하며, 이후 OK heartbeat로 내려온 shard만 다시 read할 수 있다.

Coordinator는 member가 보고한 ownedShards/revokingShards를 server-side target assignment와 이전에 수락한 current assignment 기준으로 검증한다. 아직 pending인 shard, 다른 live member가 소유 중인 shard, 또는 더 이상 허용되지 않는 shard를 owned로 보고하면 stale ownership으로 보고 fencing한다. 이미 처리된 terminal REVOKED duplicate report는 fencing하지 않고 무시한다.