Redis Stream Coordinator
Design docs

Group Metadata and Assignment Model

Group Metadata

{
  "streamPrefix": "summary",
  "consumerGroup": "summary-group",
  "groupEpoch": 12,
  "metadataVersion": 8,
  "assignmentEpoch": 11,
  "state": "STABLE",
  "shardCount": 2,
  "readableShardSet": [1, 2],
  "updatedAt": "2026-05-15T10:00:00Z"
}

groupEpoch은 group metadata가 바뀔 때 증가한다. Stream version 값은 정수이며 prefix 문자열을 붙이지 않는다.

Trigger:

Group State

state는 group이 target assignment에 어느 정도 수렴했는지를 나타낸다.

State Meaning
EMPTY 살아 있는 member가 없고 readable shard owner도 없다.
ASSIGNING groupEpoch > assignmentEpoch라서 coordinator가 새 target assignment를 계산해야 한다.
RECONCILING target assignment는 계산됐지만 일부 member가 revoke/assign을 아직 완료하지 않았다.
STABLE 모든 active member의 current assignment가 target assignment와 같은 epoch으로 수렴했다.

State transition:

EMPTY
  -> ASSIGNING      first live member joins

STABLE
  -> ASSIGNING      member join, member leave, idle removal, or readable metadata change
  -> EMPTY          last member leaves and no readable shard can be owned

ASSIGNING
  -> RECONCILING    target assignment persisted for current groupEpoch

RECONCILING
  -> STABLE         all active members converge
  -> ASSIGNING      another metadata change happens before convergence

Epoch Model

epoch은 시간이 아니라 stale actor를 막기 위한 세대 번호이다. coordinator는 epoch을 단조 증가시키고, member는 자신이 들고 있는 epoch이 최신인지 검증한 뒤 assignment를 적용한다.

Epoch Owner Increments When Used For
groupEpoch coordinator member join/leave, idle removal, shard migration, readable metadata 변경 새 target assignment가 필요한지 판단
assignmentEpoch coordinator target assignment를 새 group metadata 기준으로 저장할 때 target assignment가 어떤 group epoch 기준인지 표시
memberEpoch coordinator member가 새 assignment를 적용해야 하거나 fencing될 때 stale member heartbeat/read/ack 차단

Invariant:

Coordinator가 UNKNOWN_MEMBER_ID를 반환하면 consumer는 local ownership이 더 이상 유효하지 않다고 판단해야 한다. Consumer는 read/ack를 중단하고 local owned shard와 revoking shard 상태를 비운 뒤, 같은 memberIdmemberEpoch=0으로 full heartbeat를 다시 보낸다. Group metadata가 여전히 존재하면 이 heartbeat는 rejoin 요청이며, coordinator는 새 member epoch과 assignment를 내려줄 수 있다.

Member Metadata

{
  "memberId": "018f8d27-8d71-7a7c-a7e7-b3f77c613d7f",
  "memberName": "summary-consumer-a",
  "state": "ACTIVE",
  "memberEpoch": 11,
  "metadataVersion": 8,
  "runtimeMaxConcurrency": 12,
  "activeConsumerWorkers": 8,
  "currentAssignment": [
    {"shardIndex": 1},
    {"shardIndex": 3}
  ],
  "revoking": [],
  "lastHeartbeatAt": "2026-05-15T10:00:00Z",
  "memberLeaseExpiresAt": "2026-05-15T10:00:15Z"
}

Member metadata는 member lifecycle의 source of truth이다.

Logical Member와 Worker Capacity

설계에서는 논리 coordinator member와 member 내부 worker capacity를 분리한다.

Annotation 기반 consumer는 Kafka와 유사한 listener concurrency 모델을 사용한다. @StreamListener(concurrency = "4")는 같은 애플리케이션 프로세스 안에 4개의 논리 coordinator member를 만든다. Starter는 생성된 base member ID에 -m0, -m1, -m2, -m3 같은 suffix를 붙여 member ID를 만든다. 각 논리 member는 별도 heartbeat loop, Redis consumer name, assignment state, member epoch, metadata version, shard ownership을 가진다. Annotation listener에는 하나의 member로 합치는 toggle이 없다.

Bean 기반 consumer는 CoordinatorConsumerProperties.consumer(...)로 하나의 coordinator member를 명시적으로 등록할 수 있다. 이 경로에서 runtimeMaxConcurrency는 그 member 하나의 local worker capacity이다. Handler execution 또는 poll worker를 몇 개까지 동시에 실행할 수 있는지를 의미하며, 추가 coordinator member를 만들지는 않는다.

이 분리는 ownership 규칙을 단순하게 유지한다. Coordinator는 shard를 live member에 assign하고, 각 member는 자신이 소유한 shard들을 local worker capacity 안에서 어떻게 순회할지 결정한다.

Target Assignment

Target assignment는 coordinator가 계산한 desired state이다.

{
  "streamPrefix": "summary",
  "consumerGroup": "summary-group",
  "assignmentEpoch": 12,
  "groupEpoch": 12,
  "members": {
    "member-a": [
      {"shardIndex": 0}
    ],
    "member-b": [
      {"shardIndex": 1}
    ]
  }
}

Invariant:

Current Assignment

Current assignment는 member가 실제로 적용 중이라고 보고한 state이다.

{
  "memberId": "member-a",
  "memberEpoch": 11,
  "owned": [
    {"shardIndex": 0}
  ],
  "revoking": [
    {"shardIndex": 4}
  ],
  "ackedRevocations": [
    {"shardIndex": 4}
  ]
}

Coordinator는 current assignment를 보고 revoke ack가 들어온 shard만 새 owner에게 assign한다. 단, heartbeat report는 그대로 신뢰하지 않는다. Coordinator가 수락하는 ownership report는 다음 범위로 제한된다.

member가 아직 pendingShards 상태인 shard나 다른 active member가 소유한 shard를 ownedShards로 보고하면 coordinator는 해당 member를 fenced 처리하고 FENCED_MEMBER_EPOCH을 반환한다.

Member Removal State

Member 제거는 두 경로로 진행된다.

EXPIRED member metadata는 즉시 삭제하지 않는다. 운영 디버깅과 group 이력 추적을 위해 장기간 유지한다.

Sticky Partition Assignment

이 설계는 sticky partition assignment를 고정 전제로 둔다. member가 선택하거나 heartbeat로 보고할 값이 아니다.

목표:

Sticky partition assignment는 현재 target assignment를 입력으로 사용해 movement cost를 최소화한다. scale-out에서는 필요한 일부 shard만 신규 member로 이동하고, scale-in에서는 사라지는 member의 shard만 살아 있는 member에게 재분배한다.

Sticky Assignment Algorithm

Input:

Steps:

  1. previous owner가 live member이면 기존 shard owner를 우선 유지한다.
  2. excluded member가 가진 shard를 unassigned set으로 옮긴다.
  3. unassigned shard를 {shardIndex} 기준으로 deterministic sort한다.
  4. live member를 현재 shard load, optional member weight, memberId 순으로 정렬한다.
  5. 각 unassigned shard를 가장 덜 찬 member에게 배정한다.
  6. previous owner와 target owner가 다르면 coordinator는 previous owner의 assignment.assignedShards에서 해당 shard를 제외한다.
  7. coordinator는 revoke ack 또는 expired member fencing을 확인한 뒤 target owner의 assignment.assignedShards에 해당 shard를 포함한다.

maxConcurrency 처리:

Annotation listener에서 concurrency = 4는 독립적으로 assign되는 coordinator member 네 개를 뜻한다. Bean 기반 integration에서 runtimeMaxConcurrency = 4는 하나의 member가 최대 네 개의 handler execution을 병렬로 실행할 수 있다는 뜻이다. 어떤 경우에도 member가 shard 네 개까지만 소유할 수 있다는 뜻은 아니다.

하나의 member가 local worker 수보다 많은 shard를 소유하면 built-in polling adapter는 owned shard를 round-robin으로 순회한다. 뒤쪽 shard index도 계속 poll/processing되어야 하며, assignment 불균형이 영구적인 shard starvation으로 이어지면 안 된다.