OurNight; engineering — POS
0x10. requirement
QR 주문 → 입금 확인 → 조리 → 서빙 → 정산 팝업 주점의 주문·입금 확인·조리·서빙 정보를 실시간으로 연결해 주문 누락과 현장 소통 비용을 줄인다.
0x11. 홀 POS
1. 주문·입금 관리
- 신규 주문과 입금 확인 요청을 대기열에 표시한다.
- 테이블 번호, 주문 내역, 주문 차수, 총금액을 보여준다.
- 스태프가 입금을 확인하고 주문을 확정한다.
- 확정된 주문을 주방으로 전달한다.
- 홀에서 주문을 직접 추가하거나 취소·수량 변경할 수 있다.
2. 서빙 관리
- 조리 완료된 메뉴를 서빙 대기열에 표시한다.
- 테이블 번호, 메뉴명, 수량을 보여준다.
- 서빙 후 메뉴 단위로 완료 처리한다.
- 모든 메뉴가 서빙되면 주문을 완료 처리한다.
3. 호출 관리
- 물, 티슈, 수저, 컵 등의 고객 요청을 표시한다.
- 요청을 처리 중 또는 완료 상태로 변경한다.
- 여러 스태프가 사용할 경우 담당자를 지정해 중복 응대를 막는다.
4. 테이블 관리
- 테이블별 사용 여부와 진행 중 주문을 표시한다.
- 손님 퇴장 후 테이블을 정리하고 새 세션으로 초기화한다.
- 손님이 자리를 옮기면 주문을 새 테이블로 이전한다.
- 합석 시 두 테이블의 주문과 세션을 하나로 병합한다.
- 테이블 이용시간과 정리 필요 여부를 확인한다.
0x12. 주방 POS
1. 주문서 관리
- 입금이 확인된 주문만 주방에 표시한다.
- 테이블 번호, 메뉴명, 수량, 옵션, 주문 경과시간을 제공한다.
- 주문을 들어온 순서대로 표시한다.
- 메뉴 단위로 다음 상태를 관리한다.
조리 대기 → 조리 중 → 조리 완료 → 서빙 완료- 조리 완료 시 홀 POS에 알림을 보낸다.
- 진행 중 주문과 완료 주문을 구분한다.
2. 메뉴·재고 관리
- 메뉴를 즉시 품절 처리할 수 있다.
- 품절 상태를 고객 메뉴판에 실시간 반영한다.
- 재고 추가 또는 품절 해제가 가능해야 한다.
0x13. 도어 POS
- 남은 테이블 수와 현재 혼잡도를 표시한다.
- 대기 발생 여부를 설정한다.
- 입장권을 QR 또는 인증 코드로 검증한다.
- 유효한 입장권만 승인한다.
0x14. Admin
- 실시간 매출과 순매출을 조회한다.
- 주문 수와 메뉴별 판매량을 조회한다.
- 진행 중·완료된 주문 내역을 조회한다.
- 취소·환불 요청과 처리 이력을 관리한다.
5. 공통 요구사항
- 홀·주방·도어 화면의 상태를 실시간 동기화한다.
- 여러 스태프의 중복 처리를 방지한다.
- 모든 상태 변경에 처리자와 처리 시각을 기록한다.
- 네트워크 재연결 후 누락된 변경 사항을 복구한다.
- 중복 클릭이나 재요청으로 주문이 여러 번 처리되지 않게 한다.
0x20. DESIGN
0x21. happy path
0x22. first connection or server restart
상황
- POS 처음 엶
- server restart
문제
- Redis reader는 현재 stream tail 근처에서 시작한다.
- SSE connection 이전에 발생한 event는 reader가 놓칠 수 있다.
- SSE event만으로 현재 POS state를 만들 수 없다.
design
SSE 연결
→ full baseline snapshot load
→ snapshot watermark 모음
→ 가장 작은 watermark까지 covered 처리
→ 이후 SSE event patch0x23. temporary network loss
상황
- Wi-Fi가 잠깐 끊겼다.
- browser가 sleep 상태에 들어갔다.
- SSE connection만 잠시 사라졌다.
- Redis Stream에는 누락 event가 아직 남아 있다.
design
감지
EventSource onerror
또는
45초 watchdog timeout
또는
online / focus / visibilitychangerecovery
기존 EventSource cleanup
→ backoff + jitter
→ 새 ticket 발급
→ lastAppliedEventId 포함해 reconnect
→ Redis Stream에서 (lastEventId, +) replay
→ dedupe
→ Redux patch0x24. incremental recovery invalidated
0x23를 더 이상 믿을 수 없는 상태
상황
| 상황 | incremental recovery가 깨지는 이유 |
|---|---|
| 장기 유실 | 빠진 event 범위가 너무 크거나 알 수 없다. |
| stream trim | 필요한 event가 이미 없어졌다. |
| replay 과다 | 기술적으로 가능해도 replay 비용과 queue risk가 너무 크다. |
| unknown event | event를 어떻게 state에 반영해야 할지 해석할 수 없다. |
| payload parse 실패 | event meaning 자체를 만들 수 없다. |
| queue overflow | 중간 event가 모두 유지됐다는 보장이 없다. |
문제
Level 1: event replay Level 2: delta calculation + partial reload Level 3: full baseline reload
1단계: event replay
다음 조건이면 replay를 믿을 수 없다.
- stream trim으로 gap expired
- replay가 300개 초과
- invalid stream ID
- unknown event type
- invalid payload
- replay exception
2단계: delta calculation + partial reload
sinceEventId 전달
→ Redis Stream scan
→ changed order/table/menu/call ID 계산3단계: full baseline reload
다음 상황이면 delta recovery도 포기
/changesevent가 500개 초과- changed entity가 종류별 100개 초과
- replay gap expired
- payload/type를 믿을 수 없음
sinceEventId가 없음- stream 내부에
RELOAD_REQUIRED존재
design flow
0x25. high traffic
design flow
0x26. slow consumer (single device failure)
design flow
0x27. redis publish failure
상황
- 주문 DB transaction은 성공했다.
- outbox row도 저장됐다.
- Redis가 잠시 unavailable이다.
design
REST command transaction commit
→ outbox row remains unpublished
→ relay publish 실패
→ attempts 증가
→ nextAttemptAt 설정
→ 최대 60초 backoff
→ 이후 retry0x2F. WHY SSE?
0x20. IMPL — core flow (happy)
대표 사례로 “결제 확인 및 주문 확정” 케이스를 설명하겠다.
0x21. server state 확정
POST /api/pos/{storeId}/hall/orders/{orderId}/confirm
→ HallOrderController.confirmOrder()
→ OrderService.confirmOrder()
→ OrderCommandService.confirmOrder()0x22. outbox pattern
// PosNotificationPublisher
@Component
class PosNotificationPublisher(... 생략 ...) {
@TransactionalEventListener(phase = TransactionPhase.BEFORE_COMMIT)
fun handleCookingRequest(event: CookingRequestedEvent) {
if (!isCurrentSlotOrderEvent(event.storeId, event.orderId)) return
enqueueOutbox(event.storeId, "KITCHEN_COOK", event)
}
private fun enqueueOutbox(storeId: UUID, type: String, payload: Any) {
outboxRepository.save(
PosEventOutbox(
storeId = storeId,
eventType = type,
payload = objectMapper.writeValueAsString(payload),
)
)
}
}- TransactionPhase.BEFORE_COMMIT
- domain state change + outbox insert -> 같은 Transaction에 참여
0x23. outbox relay가 Redis Stream에 publish
@Component
class PosOutboxRelay(...) {
@Scheduled(fixedDelay = 500)
@Transactional
fun publishPendingEvents() {
val events = outboxRepository.lockPublishable(batchSize)
if (events.isEmpty()) return
events.forEach { event ->
try {
val streamId = publishToRedis(event)
event.markPublished(streamId)
} catch (e: Exception) {
event.markFailed(e)
// log ...
}
}
}
}
interface PosEventOutboxRepository : JpaRepository<PosEventOutbox, UUID> {
@Query(
value = """
select *
from pos_event_outbox
where published_at is null
and next_attempt_at <= now()
order by created_at
limit :limit
for update skip locked
""",
nativeQuery = true,
)
fun lockPublishable(@Param("limit") limit: Int): List<PosEventOutbox>
}0x24. 매장별 SSE reader가 stream event를 읽음
outline
GET /api/pos/{storeId}/notifications/subscribe
→ SseNotificationService.subscribe()
→ registerEmitter()
→ startRedisReader()
→ readStreamLoop()code
fun subscribe(
expectedStoreId: UUID,
ticket: String,
connectionId: String?,
lastEventId: String?
): SseEmitter {
val storeIdStr = redisTemplate.opsForValue().getAndDelete("ticket:$ticket")
?: throw BusinessException(ErrorCode.INVALID_TICKET)
if (storeIdStr != expectedStoreId.toString()) {
throw BusinessException(ErrorCode.INVALID_TICKET)
}
val emitter = SseEmitter(60 * 60 * 1000L)
val chanKey = channelKey(storeIdStr)
val sKey = streamKey(storeIdStr)
val safeConnectionId = connectionId?.takeIf { it.isNotBlank() }?.take(128) ?: UUID.randomUUID().toString()
registerEmitter(chanKey, emitter, sKey, safeConnectionId)
markActive(chanKey)
logger.info("[SSE-CONNECT] storeId={} connectionId={} lastEventId={}", storeIdStr, safeConnectionId, lastEventId)
enqueueOutbound(chanKey, emitter, SseMsg("connect", "connected", "0-0"))
if (!lastEventId.isNullOrBlank()) {
replayMissingEvents(sKey, lastEventId, emitter)
markActive(chanKey)
}
return emitter
}
private fun registerEmitter(chanKey: String, emitter: SseEmitter, streamKey: String, connectionId: String) {
connectionIds[emitter] = connectionId
outboundQueueOf(emitter)
startEmitterSender(chanKey, emitter, connectionId)
emitter.onCompletion { unregisterEmitter(chanKey, emitter, "completion") }
emitter.onTimeout { emitter.complete(); unregisterEmitter(chanKey, emitter, "timeout") }
emitter.onError {
logger.warn("[SSE-ERROR] storeId={} connectionId={}", chanKey, connectionId)
unregisterEmitter(chanKey, emitter, "error")
}
emitters.compute(chanKey) { _, list ->
val newList = list ?: CopyOnWriteArrayList()
newList.add(emitter)
newList
}
startRedisReader(chanKey, streamKey)
}
private fun startRedisReader(chanKey: String, streamKey: String) {
streamReaders.compute(chanKey) { _, existing ->
if (existing != null) return@compute existing
try {
val startId = currentStreamTailId(streamKey)
val task = streamReaderExec.submit { readStreamLoop(chanKey, streamKey, startId) }
logger.info("Redis Stream Reader Started: chanKey={} streamKey={} startId={}", chanKey, streamKey, startId)
task
} catch (e: Exception) {
logger.error("Failed to start redis stream reader", e)
null
}
}
}
private fun readStreamLoop(chanKey: String, streamKey: String, initialLastSeenId: String) {
var lastSeenId = initialLastSeenId
val options = StreamReadOptions.empty()
.count(100)
.block(Duration.ofSeconds(2))
while (!Thread.currentThread().isInterrupted && emitters.containsKey(chanKey)) {
try {
val records = redisTemplate.opsForStream<String, String>().read(
options,
StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))
) ?: emptyList()
for (record in records) {
lastSeenId = record.id.toString()
enqueueRecord(chanKey, streamKey, record)
}
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
} catch (e: Exception) {
logger.warn("Redis stream read failed: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
Thread.sleep(1000)
}
}
logger.info("Redis Stream Reader Stopped: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
}0x25. SseEmitter가 event id와 payload를 브라우저에 전달
outline
매장 queue
→ broadcastNow()
→ emitter별 outbound queue
→ startEmitterSender()
→ sendToEmitter()
→ emitter.send()code
private fun broadcastNow(chanKey: String, msg: SseMsg) {
emitters[chanKey]?.forEach { emitter ->
enqueueOutbound(chanKey, emitter, msg)
}
}
private fun startEmitterSender(chanKey: String, emitter: SseEmitter, connectionId: String) {
val task = emitterSenderExec.submit {
while (!Thread.currentThread().isInterrupted && outboundQueues.containsKey(emitter)) {
val msg = try {
outboundQueues[emitter]?.poll(30, TimeUnit.SECONDS)
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
null
} ?: continue
if (!sendToEmitter(emitter, msg.name, msg.data, msg.id)) {
unregisterEmitter(chanKey, emitter, "send-failed")
break
}
}
logger.info("[SSE-SENDER-STOP] storeId={} connectionId={}", chanKey, connectionId)
}
emitterSenderTasks[emitter] = task
}
private fun sendToEmitter(emitter: SseEmitter, name: String, data: String, id: String): Boolean {
try {
emitter.send(
SseEmitter.event()
.id(id)
.name(name)
.data(data)
.reconnectTime(2000)
)
return true
} catch (e: Exception) {
logger.debug("Failed to send SSE event: ${e.message}")
emitter.completeWithError(e)
return false
}
}0x26. FE: Redux에 patch하고 lastAppliedEventId 저장
0. outline
PosRoot
→ PosRealtimeBridge
→ ticket 발급
→ EventSource 생성
→ MessageEvent 수신
→ eventAdapter
→ Redux realtime action dispatch
→ listener middleware에서 entity patch
→ markStreamEventApplied
→ lastAppliedEventId 갱신1. bridge mount
// PosRoot.tsx
export default function PosRoot() {
const { storeId = "" } = useParams();
return (
<StoreProvider>
<PosRealtimeBridge storeId={storeId} />
<div className="contents font-sans">
<Outlet />
</div>
</StoreProvider>
);
}2. EventSource 생성과 event id 수신
// sseClient.ts
export function subscribeSse(params: {
url: string;
withCredentials?: boolean;
handlers: Record<string, (raw: string, meta: SseEventMeta) => void>; // eventName -> handler(rawJson)
onError?: (ev?: Event) => void;
onOpen?: () => void;
}): SseSubscription {
const es = new EventSource(params.url, { withCredentials: params.withCredentials ?? true });
for (const [eventName, handler] of Object.entries(params.handlers)) {
es.addEventListener(eventName, (ev) => {
const message = ev as MessageEvent;
handler(message.data, { lastEventId: message.lastEventId });
});
}
es.onopen = () => params.onOpen?.();
es.onerror = (ev) => params.onError?.(ev);
return {
close: () => es.close(),
};
}3. lastAppliedEventId를 reconnect URL에 포함
// PosRealtimeBridge.tsx
export function PosRealtimeBridge({ storeId }: Props) {
const dispatch = useAppDispatch();
const lastAppliedEventId = useAppSelector((state) => state.pos.stream.lastAppliedEventId);
// ....// PosRealtimeBridge.tsx
export function PosRealtimeBridge({ storeId }: Props) {
const dispatch = useAppDispatch();
const lastAppliedEventId = useAppSelector((state) => state.pos.stream.lastAppliedEventId);
// ....
// ...
useEffect(() => {
//. ...
const reconcileAfterResume = () => {
//. ...
const connectionId = createRequestId();
const url = getSseUrl(storeId, ticket, connectionId, lastAppliedEventIdRef.current);4. SSE 이벤트를 Redux action으로 변환
// listener.ts
export const dispatchPosRealtimeEventWithId = (
event: PosRealtimeEvent,
streamEventId: string,
dedupeEventId?: string,
): PosRealtimeAction & { meta: { streamEventId: string; dedupeEventId?: string } } =>
({
type: `pos/realtime/${event.type}`,
payload: event.payload,
meta: { streamEventId, dedupeEventId },
}) as PosRealtimeAction & { meta: { streamEventId: string; dedupeEventId?: string } };5. Redux entity patch
// listener.ts
posListenerMiddleware.startListening({
effect: async (action, api) => {
// ...
switch (action.type) {
// ...
case "pos/realtime/payment/resolved": {
const payload = action.payload;
const pos = getPosState(api);
const orderExists = !!pos.entities.ordersById[payload.orderId];
if (orderExists) {
api.dispatch(
posActions.patchOrderStatus({
orderId: payload.orderId,
status: toPosOrderStatus(payload.currentOrderStatus),
}),
);
} else {
invalidatePayments(api);
}
api.dispatch(posActions.removePaymentQueue(payload.orderId));
invalidateHall(api);
invalidateKitchen(api);
invalidateTables(api);
break;
}6. lastAppliedEventId 저장
//listener.ts
posListenerMiddleware.startListening({
// .....
effect: async (action, api) => {
//......
if (isReplayableStreamEvent) {
api.dispatch(
posActions.markStreamEventApplied({
eventId: streamEventId,
dedupeId,
snapshotScopes: getSnapshotScopesForRealtimeAction(action.type),
}),
);
}// slice.ts
export const posSlice = createSlice({
name: "pos",
initialState: initialPosDomainState,
reducers: {
markStreamEventApplied(
state,
action: PayloadAction<{ eventId: string; dedupeId?: string; snapshotScopes: StreamSnapshotScope[] }>
) {
const eventId = normalizeRedisStreamId(action.payload.eventId);
rememberAppliedEventId(state, eventId);
rememberAppliedDedupeId(state, action.payload.dedupeId);
if (isRedisStreamIdAfter(eventId, state.stream.lastAppliedEventId)) {
state.stream.lastAppliedEventId = eventId;
}
for (const scope of action.payload.snapshotScopes) {
const prev = state.stream.lastAppliedEventIdBySnapshot[scope];
if (isRedisStreamIdAfter(eventId, prev)) {
state.stream.lastAppliedEventIdBySnapshot[scope] = eventId;
}
}
},0x27. 다음 이벤트 또는 Redis Stream에서 대기
XREAD 결과 있음
→ 각 record 처리
→ lastSeenId 갱신
→ 다시 XREAD
XREAD 결과 없음
→ 최대 2초 blocking
→ 결과가 없으면 빈 목록
→ while loop를 돌며 다시 XREAD
새 event publish
→ blocking read 반환
→ record 처리// SseNotificationService.kt
private fun readStreamLoop(chanKey: String, streamKey: String, initialLastSeenId: String) {
var lastSeenId = initialLastSeenId
val options = StreamReadOptions.empty()
.count(100)
.block(Duration.ofSeconds(2))
while (!Thread.currentThread().isInterrupted && emitters.containsKey(chanKey)) {
try {
val records = redisTemplate.opsForStream<String, String>().read(
options,
StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))
) ?: emptyList()
for (record in records) {
lastSeenId = record.id.toString()
enqueueRecord(chanKey, streamKey, record)
}
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
} catch (e: Exception) {
logger.warn("Redis stream read failed: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
Thread.sleep(1000)
}
}
logger.info("Redis Stream Reader Stopped: chanKey={} streamKey={} lastSeenId={}", chanKey, streamKey, lastSeenId)
}
0x2F. 정리
0x30. IMPL — first connection or server restart
outline
1. SSE 연결
2. 연결 성공 시 baseline 복구 실행
3. 최초 연결(`lastAppliedEventId = 0-0`)이면 full baseline 실행
4. 각 snapshot 응답의 watermark를 scope별로 저장
5. 모든 snapshot 중 가장 작은 watermark까지 stream cursor를 전진
6. SSE 이벤트를 Redux 상태에 patch0x31. Redis reader가 현재 stream tail에서 시작
//SseNotificationService.kt
val startId = currentStreamTailId(streamKey)
val task = streamReaderExec.submit {
readStreamLoop(chanKey, streamKey, startId)
}tail 조회
//SseNotificationService.kt
private fun currentStreamTailId(streamKey: String): String =
try {
val info = redisTemplate.opsForStream<String, String>().info(streamKey)
info.lastGeneratedId() ?: info.lastEntryId() ?: "0-0"
} catch (_: Exception) {
"0-0"
}이 startId 이후 이벤트만 읽는다.
// SseNotificationService.kt -> readStreamLoop
StreamOffset.create(streamKey, ReadOffset.from(lastSeenId))0x32. SSE 연결 성공 후 baseline 복구
// PosRealtimeBridge.tsx
const sub = subscribeSse({
url,
handlers,
onOpen: () => {
touchSseActivity();
reconnectDelayMsRef.current = INITIAL_RECONNECT_DELAY_MS;
startWatchdog();
reloadBaseline();
},
});onOpen에서 reloadBaseline()을 호출하므로 순서는
EventSource 연결 성공
→ baseline 복구 시작0x33. 최초 POS 연결이면 full baseline으로 전환
초기 stream cursor는 [initialState.ts (line 36)]에서 0-0
initialState.ts
stream: {
lastAppliedEventId: "0-0",
}연결 후 실제 호출되는 것은 reloadChangedPosBaseline이지만, cursor가 0-0이면 full baseline으로 전환
//reloadChangedPosBaseline.ts
if (!sinceEventId || sinceEventId === "0-0") {
await dispatch(reloadPosBaseline(storeId)).unwrap();
return { storeId, mode: "full" as const };
}즉 최초 POS 오픈 흐름은
SSE open
→ reloadChangedPosBaseline(sinceEventId = "0-0")
→ reloadPosBaseline()
→ full snapshot load재연결이고 유효한 cursor가 있으면 전체 baseline 대신 변경분 복구를 먼저 시도
0x34. full baseline snapshot 병렬 조회
//reloadPosBaseline.ts
await Promise.all([
dispatch(loadPaymentQueue({ storeId, force: true })).unwrap(),
dispatch(loadServingQueue({ storeId, force: true })).unwrap(),
dispatch(loadCallQueue({ storeId, force: true })).unwrap(),
dispatch(loadHallTableZones(storeId)).unwrap(),
dispatch(loadHallTables({ storeId, force: true })).unwrap(),
dispatch(loadMenuCatalog(storeId)).unwrap(),
dispatch(loadHallOrders(...)).unwrap(),
dispatch(loadKitchenOrders(...)).unwrap(),
]);이 부분이 full baseline snapshot load에 해당
0x35. 서버가 snapshot watermark를 응답 헤더에 추가
// PosStreamWatermarkService.kt
fun currentWatermark(storeId: UUID): String {
val key = PosStreamKey.store(storeId)
val info = redisTemplate.opsForStream<String, String>().info(key)
return info.lastGeneratedId() ?: info.lastEntryId() ?: "0-0"
}응답에는 X-Pos-Stream-Watermark 헤더로 포함
ResponseEntity.ok()
.header(POS_STREAM_WATERMARK_HEADER, watermark)
.body(body)프론트는 [streamSnapshot.ts (line 11)]에서 이 헤더를 읽음
0x36. snapshot별 watermark 저장
예를 들어 결제 snapshot은 loadPaymentQueue.ts (line 45)에서 저장
dispatch(posActions.markStreamSnapshotLoaded({
eventId: snapshot.streamWatermark,
snapshotScopes: ["payments"],
}));Reducer는
// slice.ts
for (const scope of action.payload.snapshotScopes) {
const prev = state.stream.lastAppliedEventIdBySnapshot[scope];
if (isRedisStreamIdAfter(eventId, prev)) {
state.stream.lastAppliedEventIdBySnapshot[scope] = eventId;
}
}payments, servings, calls, hallOrders, kitchenOrders, tables별 watermark가 각각 관리됩니다.
0x37. 가장 작은 watermark까지 covered 처리
//reloadPosBaseline.ts
const coveredWatermark = BASELINE_SNAPSHOT_SCOPES
.map((scope) => stream.lastAppliedEventIdBySnapshot[scope])
.reduce((min, eventId) =>
compareRedisStreamIds(eventId, min) < 0 ? eventId : min
);
dispatch(posActions.markStreamWatermarkCovered({
eventId: coveredWatermark,
}));Reducer는
// reducer.ts
if (isRedisStreamIdAfter(eventId, state.stream.lastAppliedEventId)) {
state.stream.lastAppliedEventId = eventId;
}가장 작은 값을 사용하는 이유는 일부 snapshot만 최신일 때 stream cursor를 너무 앞으로 보내서 이벤트를 건너뛰지 않도록 하기 위해서
0x38. 이후 SSE 이벤트 patch
SSE 이벤트를 Redux action으로 변환하고 dispatch하는 부분
PosRealtimeBridge.tsx
const event = adaptRealtimeEvent(eventName, parsed);
const reduxAction = dispatchPosRealtimeEventWithId(
event,
meta.lastEventId,
dedupeEventId,
);
dispatch(reduxAction);-
실제 entity patch와 stream cursor 갱신은 listener.ts (
posListenerMiddleware.startListening)에서 처리 -
이벤트 적용이 끝나면 해당 Redis Stream ID까지 적용됐다고 기록
0x3F. 정리
0x40. IMPL — temporary network loss
- 학교 와이파이는 터질 가능성이 늘 존재한다
- 이에 다음 전략을 사용
| 단계 | 전략 | 역할 |
|---|---|---|
| 1 | Redis Stream cursor | Redis Stream에서 읽은 도메인 이벤트의 SSE event id에 Redis Stream ID를 넣는다. |
| 2 | Dedupe/idempotent apply | 프론트가 이벤트를 적용하면서 lastAppliedEventId를 갱신한다. |
| 3 | Client watchdog/reconnect | 연결 오류, idle timeout, online/focus 복귀 시 재연결한다. |
| 4 | Last-Event-ID replay | 재연결 URL에 lastEventId를 실어 누락 구간 replay를 요청한다. |
| 5 | Snapshot reconciliation | replay가 불가능하면 변경분 reload 또는 full baseline으로 전환한다. |
0x41. outline
0x42. 설명
이 구조에서 중요한 점은 네트워크 유실을 “예외”로 보지 않는다는 것이다. POS 단말은 와이파이, 태블릿 sleep, 브라우저 background 상태의 영향을 자주 받기 때문에, 유실 후 재동기화가 기본 경로로 들어가 있다.
서버 전송부는 SSE event id를 명시적으로 넣는다. Redis Stream에서 읽은 도메인 이벤트는 id가 Redis Stream ID이고, PING, connect, 일부 즉시 reload 신호처럼 stream record가 아닌 제어 이벤트는 0-0을 쓴다.
private fun sendToEmitter(emitter: SseEmitter, name: String, data: String, id: String): Boolean {
emitter.send(
SseEmitter.event()
.id(id)
.name(name)
.data(data)
.reconnectTime(2000)
)
return true
}프론트 EventSource 래퍼는 브라우저의 MessageEvent.lastEventId를 handler로 넘긴다.
es.addEventListener(eventName, (ev) => {
const message = ev as MessageEvent;
handler(message.data, { lastEventId: message.lastEventId });
});프론트가 연결을 다시 만들 때는 마지막 적용 지점을 URL에 싣는다.
const url = getSseUrl(storeId, ticket, connectionId, lastAppliedEventIdRef.current);그래도 복구가 불가능한 상황은 full baseline으로 떨어진다.
| full baseline으로 전환되는 조건 | 이유 |
|---|---|
| Redis Stream이 trim되어 gap이 사라진 경우 | lastEventId 이후 이벤트를 완전히 재생할 수 없다. |
| replay 이벤트가 300개를 넘는 경우 | 단일 연결 replay가 너무 커져 SSE 송신 큐를 압박할 수 있다. |
/notifications/changes 결과가 500개를 넘는 경우 | 변경분 계산 자체가 너무 커져 partial reload의 장점이 줄어든다. |
| payload parse 실패 | 이벤트 의미를 신뢰할 수 없다. |
| unknown event type | 프론트/백엔드 이벤트 스키마가 어긋났을 가능성이 있다. |
| SSE queue overflow | 중간 이벤트를 모두 유지하지 못했으므로 cursor만 믿을 수 없다. |
| slow consumer | 해당 단말이 이벤트 소비 속도를 따라오지 못한다. |
| baseline partial reload 대상 entity가 너무 많은 경우 | 개별 재조회보다 전체 snapshot이 더 안정적이다. |
이 점에서 현재 구현은 “유실을 절대 막는다”가 아니라 “유실을 감지하면 빠르게 기준 상태를 다시 잡는다”는 설계다.
0x50. IMPL — incremental recovery invalidated
SSE replay 시도
→ replay 불신 조건이면 RELOAD_REQUIRED
→ 클라이언트 reloadChangedPosBaseline 실행
→ 유효한 lastAppliedEventId가 있으면 /changes 호출
→ delta 계산 가능 + entity 100개 이하이면 partial reload
→ 아니면 full baseline reload
→ watermark까지 cursor 전진
→ SSE event 적용 재개1단계: Event replay
0x31. replay 시작
클라이언트가 lastEventId를 전달하면 서버가 누락 이벤트를 replay
// SseNotificationService.kt
if (!lastEventId.isNullOrBlank()) {
replayMissingEvents(sKey, lastEventId, emitter)
markActive(chanKey)
}실제 replay 구현은 다음 코드
//SseNotificationService.kt
private fun replayMissingEvents(
streamKey: String,
lastEventId: String,
emitter: SseEmitter,
) {
try {
if (isReplayGapExpired(streamKey, lastEventId)) {
sendReloadSignal(emitter)
return
}
val range = Range.rightOpen(lastEventId, "+")
val events = redisTemplate.opsForStream<String, String>()
.range(streamKey, range)
?: return
var count = 0
for (record in events) {
if (count++ >= 300) {
enqueueReplay(
emitter,
SseMsg("RELOAD_REQUIRED", "{}", "0-0"),
)
return
}
val data = record.value
val type = data["type"] ?: "UNKNOWN"
if (type !in knownEventNames) {
enqueueReplay(
emitter,
SseMsg(
"RELOAD_REQUIRED",
"{}",
record.id.toString(),
),
)
return
}
enqueueReplay(/* event */)
}
} catch (e: Exception) {
sendReloadSignal(emitter)
}
}replay 불신 조건 대응
| 조건 | 구현 위치 |
|---|---|
| stream trim / gap expired | isReplayGapExpired() |
| invalid stream ID | parseStreamId() 실패 시 gap expired로 처리 |
| replay 300개 초과 | count++ >= 300 |
| unknown event type | type !in knownEventNames |
| replay exception | catch에서 sendReloadSignal() |
| invalid payload | 서버 replay가 아닌 클라이언트 safeJson()에서 감지 |
Stream trim과 invalid ID 검사
// SseNotificationService.kt
private fun isReplayGapExpired(
streamKey: String,
lastEventId: String,
): Boolean {
val lastId = parseStreamId(lastEventId) ?: return true
val info = try {
redisTemplate.opsForStream<String, String>().info(streamKey)
} catch (_: Exception) {
return false
}
val firstId = info.firstEntryId() ?: return false
return compareStreamIds(
parseStreamId(firstId) ?: return false,
lastId,
) > 0
}lastEventId형식이 잘못되면parseStreamId()가null을 반환하므로 replay를 신뢰하지 ㄴㄴ- Redis Stream 첫 ID가
lastEventId보다 뒤라면 필요한 이벤트가 trim된 것으로 판단
Unknown event 감지
서버가 신뢰하는 event 목록은 다음 코드에 정의되어 있음.
// SseNotificationService.kt
private val knownEventNames = setOf(
"RELOAD_REQUIRED",
"HALL_PAYMENT",
"HALL_PAYMENT_RESOLVED",
// ...
"KITCHEN_MENU_CHANGED",
)stream reader와 replay 양쪽에서 이 목록에 없는 이벤트를 발견하면 RELOAD_REQUIRED로 전환
0x32. Payload parse 실패 감지
서버 replay는 payload JSON을 검증하지 않고 클라이언트로 전달. 클라이언트가 JSON 파싱 실패를 감지
//sseClient.ts
export function safeJson<T>(raw: string): T | null {
try {
return JSON.parse(raw) as T;
} catch {
return null;
}
}// PosRealtimeBridge.tsx
const parsed = safeJson<unknown>(raw);
if (!parsed) {
console.error("[SSE PARSE ERROR]", {
storeId,
eventName,
raw,
});
reloadBaseline();
return;
}event name을 프론트 도메인 event로 매핑하지 못하는 경우도 recovery를 실행
const event = adaptRealtimeEvent(eventName, parsed);
if (!event) {
reloadBaseline();
return;
}매핑 구현은 eventAdapter.ts (line 13)에 있습니다.
0x33. Queue overflow
서버에는 세 종류의 queue가 있다.
// 매장 broadcast queue
ArrayBlockingQueue<SseMsg>(2000)
// emitter별 replay queue
ArrayBlockingQueue<SseMsg>(500)
// emitter별 outbound queue
ArrayBlockingQueue<SseMsg>(500)Broadcast queue overflow
if (q.offer(msg)) {
markActive(chanKey)
return
}
logger.warn("queue overflow: chanKey=$chanKey")
q.clear()
q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
markActive(chanKey)중간 이벤트를 모두 버리고 RELOAD_REQUIRED만 남김
Replay queue overflow
if (!q.offer(msg)) {
q.clear()
q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
}Outbound queue overflow
logger.warn(
"[SSE-SLOW-CONSUMER] outbound queue overflow"
)
q.clear()
q.offer(SseMsg("RELOAD_REQUIRED", "{}", "0-0"))
emitter.complete()
unregisterEmitter(chanKey, emitter, "slow-consumer")이 경우에는 RELOAD_REQUIRED를 넣은 직후 연결을 finalize. 따라서 control event 전달보다 연결 종료 → 재연결 → baseline recovery가 실질적인 복구 경로
2단계: Delta calculation + partial reload
0x34. /notifications/changes 호출
프론트 API 코드
// posStreamChangesApi.ts
const sp = new URLSearchParams();
sp.set("since", sinceEventId);
sp.set("limit", "500");
const res = await apiRequest(
`/api/pos/${storeId}/notifications/changes?${sp.toString()}`,
);백엔드 endpoint는
// PosNotificationController.kt
@GetMapping("/changes")
fun getChangesSince(
@RequestParam since: String?,
@RequestParam(
required = false,
defaultValue = "500",
) limit: Int,
) = ResponseEntity.ok(
posStreamChangeService.collectChanges(
storeId = storeId,
sinceEventId = since,
maxEvents = limit.coerceIn(1, 1000),
),
)현재 프론트가 전달하는 기준은 500개
0x35. Redis Stream scan
// PosStreamChangeService.kt
val records = redisTemplate.opsForStream<String, String>()
.range(
streamKey,
Range.rightOpen(normalizedSince, "+"),
)
?: emptyList()이후 event payload에서 변경된 entity ID를 collect
val orderIds = linkedSetOf<UUID>()
val tableIds = linkedSetOf<UUID>()
val menuIds = linkedSetOf<UUID>()
val callRequestIds = linkedSetOf<UUID>()대표적인 계산
when (type) {
"HALL_PAYMENT",
"HALL_SERVE",
"KITCHEN_COOK" -> {
payload.uuid("orderId")?.let(orderIds::add)
payload.uuid("tableId")?.let(tableIds::add)
payload.uuid("menuId")?.let(menuIds::add)
}
"CUSTOMER_CALL" -> {
payload.uuid("requestId")?.let(callRequestIds::add)
payload.uuid("tableId")?.let(tableIds::add)
}
"KITCHEN_MENU_CHANGED" -> {
payload.uuid("menuId")?.let(menuIds::add)
}
}0x36. Delta 계산을 포기하는 조건
| 조건 | 반환 reason |
|---|---|
invalid sinceEventId | INVALID_SINCE_EVENT_ID |
sinceEventId == 0-0 | NO_BASELINE_EVENT_ID |
| stream trim | REPLAY_GAP_EXPIRED |
| event 500개 초과 | TOO_MANY_EVENTS |
| type field 없음 | MISSING_EVENT_TYPE |
stream 안에 RELOAD_REQUIRED | STREAM_RELOAD_REQUIRED |
| payload parse 실패 | INVALID_EVENT_PAYLOAD |
| unknown event | UNKNOWN_EVENT_TYPE:{type} |
예를 들어 event 개수 제한은 다음
if (records.size > maxEvents) {
return reloadRequired(
normalizedSince,
currentWatermark,
"TOO_MANY_EVENTS",
)
}payload 검증은 다음 코드
val payload = parsePayload(record.value["payload"])
?: return reloadRequired(
normalizedSince,
currentWatermark,
"INVALID_EVENT_PAYLOAD",
)0x37. Partial reload 조건
프론트 복구 진입점은 다음
/changes가 delta 계산을 포기했다면 full baseline으로 전환
// reloadChangedPosBaseline.ts
if (changes.requiresReload) {
await dispatch(reloadPosBaseline(storeId)).unwrap();
return {
storeId,
mode: "full",
reason: changes.reason,
};
}변경 entity가 종류별 100개를 초과해도 full baseline으로 전환
const PARTIAL_RELOAD_ENTITY_LIMIT = 100;
if (
changes.changedOrderIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
changes.changedTableIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
changes.changedMenuIds.length > PARTIAL_RELOAD_ENTITY_LIMIT ||
changes.changedCallRequestIds.length > PARTIAL_RELOAD_ENTITY_LIMIT
) {
await dispatch(reloadPosBaseline(storeId)).unwrap();
return {
storeId,
mode: "full",
reason: "TOO_MANY_CHANGED_ENTITIES",
};
}0x38. 변경 entity만 다시 조회
// reloadChangedPosBaseline.ts
await Promise.all([
orderIds.length > 0
? fetchHallOrdersByIds(storeId, orderIds)
: Promise.resolve(),
orderIds.length > 0
? fetchKitchenOrdersByIds(storeId, orderIds)
: Promise.resolve(),
orderIds.length > 0
? dispatch(loadPaymentQueue({ storeId, force: true })).unwrap()
: Promise.resolve(),
orderIds.length > 0
? dispatch(loadServingQueue({ storeId, force: true })).unwrap()
: Promise.resolve(),
changes.changedCallRequestIds.length > 0
? dispatch(loadCallQueue({ storeId, force: true })).unwrap()
: Promise.resolve(),
changes.changedMenuIds.length > 0
? dispatch(loadMenuCatalog(storeId)).unwrap()
: Promise.resolve(),
tableIds.length > 0
? reloadChangedTables(storeId, tableIds, dispatch, getState)
: Promise.resolve(),
]);Partial snapshot 반영 후 /changes가 반환한 현재 watermark까지 cursor를 전진
dispatch(
posActions.markStreamWatermarkCovered({
eventId: changes.currentWatermark,
}),
);3단계: Full baseline reload
0x39. Full baseline 전환 조건
//reloadChangedPosBaseline.ts
if (!sinceEventId || sinceEventId === "0-0") {
await dispatch(reloadPosBaseline(storeId)).unwrap();
}그 외에도 다음 조건에서 full baseline을 execute
/changes요청 실패changes.requiresReload- 변경 entity 종류별 100개 초과
0x3A. 주요 snapshot 병렬 조회
// reloadPosBaseline.ts
await Promise.all([
dispatch(loadPaymentQueue({ storeId, force: true })).unwrap(),
dispatch(loadServingQueue({ storeId, force: true })).unwrap(),
dispatch(loadCallQueue({ storeId, force: true })).unwrap(),
dispatch(loadHallTableZones(storeId)).unwrap(),
dispatch(loadHallTables({ storeId, force: true })).unwrap(),
dispatch(loadMenuCatalog(storeId)).unwrap(),
dispatch(loadHallOrders(/* ... */)).unwrap(),
dispatch(loadKitchenOrders(/* ... */)).unwrap(),
]);0x3B. 최소 watermark 계산
const coveredWatermark = BASELINE_SNAPSHOT_SCOPES
.map(
(scope) =>
stream.lastAppliedEventIdBySnapshot[scope],
)
.reduce(
(min, eventId) =>
compareRedisStreamIds(eventId, min) < 0
? eventId
: min,
);
dispatch(
posActions.markStreamWatermarkCovered({
eventId: coveredWatermark,
}),
);필수 scope는 다음 여섯 개
payments
servings
calls
hallOrders
kitchenOrders
tables0x3C. Recovery가 시작되는 클라이언트 감지 지점
PosRealtimeBridge.tsx (line 323)에서 다음 상황을 감지
- SSE payload JSON parse 실패
- event mapping 실패
RELOAD_REQUIRED수신- SSE open
- resume/focus/online 주기 reconcile
if (!parsed) {
reloadBaseline();
return;
}
if (!event) {
reloadBaseline();
return;
}
if (eventName === "RELOAD_REQUIRED") {
reloadBaseline();
}watchdog timeout은 즉시 baseline을 호출하지 않고 먼저 연결을 끊고 재연결
// PosRealtimeBridge.tsx
cleanupSubscription();
scheduleReconnect();재연결 성공 후 onOpen에서 recovery가 실행
onOpen: () => {
startWatchdog();
reloadBaseline();
}0x60. IMPL — high traffic
Event traffic 증가
→ Redis Stream batch read
→ bounded queue에 임시 저장
→ 16개 shard worker가 매장별 queue 처리
→ emitter별 outbound queue로 전달
→ 한계 초과 시 RELOAD_REQUIRED 또는 연결 종료
→ 클라이언트 partial/full baseline 복구0x61. Event traffic 증가 → batch 처리
Redis Stream reader는 이벤트를 하나씩 요청하지 않고 한 번에 최대 100개씩 읽.
val options = StreamReadOptions.empty()
.count(100)
.block(Duration.ofSeconds(2))val records = redisTemplate.opsForStream<String, String>().read(
options,
StreamOffset.create(streamKey, ReadOffset.from(lastSeenId)),
) ?: emptyList()
for (record in records) {
lastSeenId = record.id.toString()
enqueueRecord(chanKey, streamKey, record)
}짧은 burst가 들어오면 Redis에서 batch로 읽어 서버 queue에 넣음
Outbox → Redis Stream 발행도 batch
//PosOutboxRelay.kt
private val streamMaxLen = 10_000L
private val batchSize = 100
@Scheduled(fixedDelay = 500)
@Transactional
fun publishPendingEvents() {
val events = outboxRepository.lockPublishable(batchSize)
}즉 500ms마다 최대 100개의 outbox event를 Redis Stream으로 발행
0x62. Bounded queue로 짧은 burst buffer
// 매장별 broadcast queue
private fun queueOf(chanKey: String) =
queues.computeIfAbsent(chanKey) {
ArrayBlockingQueue(2000)
}
// SSE 재연결 replay용 emitter별 queue
private fun replayQueueOf(emitter: SseEmitter) =
replayQueues.computeIfAbsent(emitter) {
ArrayBlockingQueue(500)
}
// 실제 브라우저 송신용 emitter별 queue
private fun outboundQueueOf(emitter: SseEmitter) =
outboundQueues.computeIfAbsent(emitter) {
ArrayBlockingQueue(500)
}| Queue | 범위 | 크기 |
|---|---|---|
queues | 매장별 broadcast | 2,000 |
replayQueues | SSE emitter별 replay | 500 |
outboundQueues | SSE emitter별 실제 송신 | 500 |
모두 ArrayBlockingQueue를 써, 메모리가 무한정 증가하지 않게 설계함.
0x63. Worker shard로 전송 작업 분산
private val senderWorkers = 16
private val senderExec = Executors.newFixedThreadPool(senderWorkers)
private fun shard(chanKey: String) =
(chanKey.hashCode() and Int.MAX_VALUE) % senderWorkers
private val activeKeys =
Array(senderWorkers) {
ConcurrentLinkedQueue<String>()
}매장 ID인 chanKey의 hash를 이용해 16개 worker 중 하나에 할당
private fun markActive(chanKey: String) {
if (activeFlag.putIfAbsent(chanKey, true) == null) {
activeKeys[shard(chanKey)].offer(chanKey)
}
}activeFlag는 같은 매장 key가 worker queue에 반복적으로 쌓이는 것을 막음
0x64. Worker가 queue drain
private fun startSenders() {
repeat(senderWorkers) { idx ->
senderExec.submit {
val myKeys = mutableListOf<String>()
while (!Thread.currentThread().isInterrupted) {
while (true) {
val key = activeKeys[idx].poll() ?: break
myKeys.add(key)
}
// 매장별 queue 처리
}
}
}
}한 worker가 특정 매장의 이벤트를 무제한으로 처리하지 않도록 한 번에 최대 50개만 꺼냅니다.
var processed = 0
while (processed < 50) {
val msg = q.poll() ?: break
broadcastNow(chanKey, msg)
processed++
}replay도 emitter별로 한 번에 최대 20개만 처리
drainReplay(
chanKey,
limitPerEmitter = 20,
)이 제한은 하나의 매장이나 emitter가 worker를 장시간 독점하는 것을 완화
0x65. 처리 한계 안쪽 → SSE 전달 계속
Broadcast 단계는 매장에 연결된 모든 emitter의 outbound queue에 이벤트를 넣음
private fun broadcastNow(
chanKey: String,
msg: SseMsg,
) {
emitters[chanKey]?.forEach { emitter ->
enqueueOutbound(chanKey, emitter, msg)
}
}각 emitter는 별도 sender task가 outbound queue를 소비해 브라우저로 전송
while (
!Thread.currentThread().isInterrupted &&
outboundQueues.containsKey(emitter)
) {
val msg = outboundQueues[emitter]
?.poll(30, TimeUnit.SECONDS)
?: continue
if (!sendToEmitter(
emitter,
msg.name,
msg.data,
msg.id,
)
) {
unregisterEmitter(
chanKey,
emitter,
"send-failed",
)
break
}
}따라서 느린 브라우저 하나의 실제 SSE 송신이 매장 broadcast worker를 직접 막지 않음
0x66. 처리 한계 초과 → incremental apply 포기
매장 broadcast queue overflow
private fun enqueue(
chanKey: String,
msg: SseMsg,
) {
val q = queueOf(chanKey)
if (q.offer(msg)) {
markActive(chanKey)
return
}
logger.warn("queue overflow: chanKey=$chanKey")
q.clear()
q.offer(
SseMsg(
"RELOAD_REQUIRED",
"{}",
"0-0",
),
)
markActive(chanKey)
}overflow가 발생하면 밀려 있던 이벤트를 모두 버리고 RELOAD_REQUIRED로 교체
중간 event를 모두 전달할 수 없음
→ event cursor 기반 incremental apply를 더 이상 신뢰하지 않음
→ snapshot 기반 복구 요청Replay queue overflow
private fun enqueueReplay(
emitter: SseEmitter,
msg: SseMsg,
) {
val q = replayQueueOf(emitter)
if (!q.offer(msg)) {
q.clear()
q.offer(
SseMsg(
"RELOAD_REQUIRED",
"{}",
"0-0",
),
)
}
}누락 이벤트 replay가 queue 처리량을 초과해도 replay를 포기
Outbound queue overflow
if (q.offer(msg)) return
logger.warn(
"[SSE-SLOW-CONSUMER] " +
"storeId={} connectionId={} " +
"outbound queue overflow",
chanKey,
connectionIds[emitter],
)
q.clear()
q.offer(
SseMsg(
"RELOAD_REQUIRED",
"{}",
"0-0",
),
)
emitter.complete()
unregisterEmitter(
chanKey,
emitter,
"slow-consumer",
)특정 브라우저가 SSE 소비 속도를 따라가지 못하면 slow consumer로 판단하고 그 연결만 종료. 다른 POS 연결은 계속 처리
0x67. RELOAD_REQUIRED → 복구 시작
클라이언트는 RELOAD_REQUIRED를 받으면 recovery를 실행
dispatch(reduxAction);
if (eventName === "RELOAD_REQUIRED") {
reloadBaseline();
}reloadBaseline()은 실제로 reloadChangedPosBaseline()을 실행
await dispatch(
reloadChangedPosBaseline({
storeId,
sinceEventId:
lastAppliedEventIdRef.current,
}),
).unwrap();0x68. Partial reload 또는 full baseline
유효한 cursor가 있으면 /changes를 이용해 partial reload를 시도
lastAppliedEventId
→ /notifications/changes
→ 변경된 order/table/menu/call ID 계산
→ 해당 entity만 재조회변경량이 많거나 delta 계산을 신뢰할 수 없으면 full baseline으로 전환
if (changes.requiresReload) {
await dispatch(
reloadPosBaseline(storeId),
).unwrap();
}if (
changes.changedOrderIds.length > 100 ||
changes.changedTableIds.length > 100 ||
changes.changedMenuIds.length > 100 ||
changes.changedCallRequestIds.length > 100
) {
await dispatch(
reloadPosBaseline(storeId),
).unwrap();
}0x70. IMPL — slow consumer (single device failure)
특정 client 송신 지연
→ emitter별 outbound queue에 이벤트 적재
→ queue limit 500 안쪽이면 계속 송신
→ queue가 가득 차면 해당 emitter만 종료
→ FE EventSource error 감지
→ backoff + jitter 재연결
→ 연결 성공 후 partial/full baseline 복구0x71. Client별 bounded queue로 지연 격리
각 SSE emitter마다 별도의 outbound queue와 sender task를 생성했음.
// SseNotificationService.kt
private val replayQueues =
ConcurrentHashMap<SseEmitter, ArrayBlockingQueue<SseMsg>>()
private fun replayQueueOf(emitter: SseEmitter) =
replayQueues.computeIfAbsent(emitter) {
ArrayBlockingQueue(500)
}
private val outboundQueues =
ConcurrentHashMap<SseEmitter, ArrayBlockingQueue<SseMsg>>()
private fun outboundQueueOf(emitter: SseEmitter) =
outboundQueues.computeIfAbsent(emitter) {
ArrayBlockingQueue(500)
}client별 실제 송신 queue 크기를 500으로 제한했음.
매장 broadcast 이벤트는 연결된 emitter 각각의 outbound queue에 복사했음.
private fun broadcastNow(
chanKey: String,
msg: SseMsg,
) {
emitters[chanKey]?.forEach { emitter ->
enqueueOutbound(chanKey, emitter, msg)
}
}따라서 특정 client의 queue가 밀려도 다른 client의 outbound queue와 직접 섞이지 않았음.
0x72. Client별 sender task로 송신 분리
emitter가 등록될 때 해당 emitter 전용 queue와 sender task를 생성했음.
private fun registerEmitter(
chanKey: String,
emitter: SseEmitter,
streamKey: String,
connectionId: String,
) {
connectionIds[emitter] = connectionId
outboundQueueOf(emitter)
startEmitterSender(chanKey, emitter, connectionId)
// ...
}전용 sender task가 해당 emitter의 queue만 소비했음.
private fun startEmitterSender(
chanKey: String,
emitter: SseEmitter,
connectionId: String,
) {
val task = emitterSenderExec.submit {
while (
!Thread.currentThread().isInterrupted &&
outboundQueues.containsKey(emitter)
) {
val msg = try {
outboundQueues[emitter]
?.poll(30, TimeUnit.SECONDS)
} catch (e: InterruptedException) {
Thread.currentThread().interrupt()
null
} ?: continue
if (!sendToEmitter(
emitter,
msg.name,
msg.data,
msg.id,
)
) {
unregisterEmitter(
chanKey,
emitter,
"send-failed",
)
break
}
}
}
emitterSenderTasks[emitter] = task
}특정 client의 네트워크 송신이 늦어도 매장 broadcast worker가 그 client의 emitter.send()를 직접 기다리지 않도록 구성했음.
0x73. Queue limit 안쪽이면 이벤트 전달 계속
outbound queue에 공간이 있으면 offer()가 성공하고 별도 조치 없이 반환했음.
private fun enqueueOutbound(
chanKey: String,
emitter: SseEmitter,
msg: SseMsg,
) {
val q = outboundQueues[emitter] ?: return
if (q.offer(msg)) return
// overflow 처리
}queue가 500개 미만인 동안 sender task가 순차적으로 이벤트를 계속 전송했음.
0x74. Queue limit 초과 시 slow consumer 연결 종료
500개 queue가 가득 차서 offer()가 실패하면 해당 emitter를 slow consumer로 판단했음.
private fun enqueueOutbound(
chanKey: String,
emitter: SseEmitter,
msg: SseMsg,
) {
val q = outboundQueues[emitter] ?: return
if (q.offer(msg)) return
logger.warn(
"[SSE-SLOW-CONSUMER] " +
"storeId={} connectionId={} " +
"outbound queue overflow",
chanKey,
connectionIds[emitter],
)
q.clear()
q.offer(
SseMsg(
"RELOAD_REQUIRED",
"{}",
"0-0",
),
)
emitter.complete()
unregisterEmitter(
chanKey,
emitter,
"slow-consumer",
)
}처리 순서는 다음과 같았음.
outbound queue overflow
→ 기존 적체 이벤트 제거
→ RELOAD_REQUIRED 삽입
→ emitter 종료
→ 해당 client 연결 제거0x75. 해당 client만 제거
unregisterEmitter()는 문제가 발생한 emitter의 queue와 sender task만 먼저 제거했음.
private fun unregisterEmitter(
chanKey: String,
emitter: SseEmitter,
reason: String,
) {
val connectionId = connectionIds.remove(emitter)
replayQueues.remove(emitter)?.clear()
outboundQueues.remove(emitter)?.clear()
emitterSenderTasks.remove(emitter)?.cancel(true)
emitters.computeIfPresent(chanKey) { _, list ->
list.remove(emitter)
if (list.isEmpty()) {
cancelStreamReader(chanKey)
queues.remove(chanKey)?.clear()
activeFlag.remove(chanKey)
null
} else {
list
}
}
}같은 매장에 다른 emitter가 남아 있으면 매장 Redis reader와 broadcast queue를 유지했음. 따라서 단일 slow consumer 장애가 다른 POS 단말의 연결 종료로 전파되지 않았음.
0x76. FE가 연결 종료 감지
서버가 emitter를 종료하면 브라우저의 EventSource error handler가 실행됐음.
const sub = subscribeSse({
url,
withCredentials: true,
handlers,
onError: (error) => {
console.error("[SSE ERROR]", {
storeId,
connectionId,
error,
mySeq,
});
if (
closedRef.current ||
connectSeqRef.current !== mySeq
) {
return;
}
if (subRef.current === sub) {
cleanupSubscription();
}
scheduleReconnect();
},
});기존 EventSource를 정리한 다음 재연결을 예약했음.
0x77. Backoff와 jitter로 재연결
재연결 지연은 2초부터 시작해 최대 30초까지 증가하도록 설정했음.
const INITIAL_RECONNECT_DELAY_MS = 2_000;
const MAX_RECONNECT_DELAY_MS = 30_000;const scheduleReconnect = () => {
if (closedRef.current) return;
if (connectingRef.current) return;
if (reconnectTimerRef.current != null) return;
const baseDelayMs =
reconnectDelayMsRef.current;
const jitterMs = Math.floor(
Math.random() *
Math.min(2_000, baseDelayMs * 0.3),
);
const delayMs = baseDelayMs + jitterMs;
reconnectDelayMsRef.current = Math.min(
baseDelayMs * 2,
MAX_RECONNECT_DELAY_MS,
);
reconnectTimerRef.current =
window.setTimeout(() => {
reconnectTimerRef.current = null;
void connect();
}, delayMs);
};동시에 여러 POS가 재연결되더라도 정확히 같은 시점에 요청하지 않도록 jitter를 추가했음.
연결에 성공하면 backoff 값을 다시 2초로 초기화했음.
onOpen: () => {
touchSseActivity();
reconnectDelayMsRef.current =
INITIAL_RECONNECT_DELAY_MS;
startWatchdog();
reloadBaseline();
},0x78. 재연결 성공 후 recovery 실행
SSE 연결이 다시 열리면 reloadBaseline()을 호출했음.
실제 함수명과 달리 항상 full baseline만 실행하는 것은 아니며, lastAppliedEventId를 기준으로 partial recovery를 먼저 시도했음.
const reloadBaseline = () => {
if (baselineReloadRef.current) return;
const p = (async () => {
await dispatch(
reloadChangedPosBaseline({
storeId,
sinceEventId:
lastAppliedEventIdRef.current,
}),
).unwrap();
})();
baselineReloadRef.current = p;
};0x79. Partial reload 또는 full baseline
유효한 cursor가 없으면 바로 full baseline을 실행했음.
if (!sinceEventId || sinceEventId === "0-0") {
await dispatch(
reloadPosBaseline(storeId),
).unwrap();
return {
storeId,
mode: "full" as const,
};
}유효한 cursor가 있으면 /changes를 통해 변경 entity만 조회했음.
await Promise.all([
orderIds.length > 0
? fetchHallOrdersByIds(storeId, orderIds)
: Promise.resolve(),
orderIds.length > 0
? fetchKitchenOrdersByIds(storeId, orderIds)
: Promise.resolve(),
changes.changedCallRequestIds.length > 0
? dispatch(
loadCallQueue({
storeId,
force: true,
}),
).unwrap()
: Promise.resolve(),
changes.changedMenuIds.length > 0
? dispatch(loadMenuCatalog(storeId)).unwrap()
: Promise.resolve(),
]);partial recovery를 신뢰할 수 없으면 full baseline으로 전환했음.
if (changes.requiresReload) {
await dispatch(
reloadPosBaseline(storeId),
).unwrap();
return {
storeId,
mode: "full" as const,
reason: changes.reason,
};
}0x80. IMPL — redis publish failure
[Redis 장애에도 이벤트가 유실되지 않고 재시도된다] Redis Stream publish는 DB transaction과 분리되어 있으며, Redis 장애 시 Outbox row가 unpublished 상태로 남아 backoff 후 재시도된다.
정상 Flow
Outbox pending
→ Relay가 조회
→ Redis Stream publish
→ stream ID 저장
→ published 처리실패 flow
Redis publish 실패
→ attempts 증가
→ nextAttemptAt 설정
→ 이후 relay에서 재조회model
@Entity
@Table(name = "pos_event_outbox")
class PosEventOutbox(
// 생략 ...
@Column(name = "published_at")
var publishedAt: Instant? = null,
@Column(name = "redis_stream_id", length = 80)
var redisStreamId: String? = null,
@Column(name = "attempts", nullable = false)
var attempts: Int = 0,
@Column(name = "next_attempt_at", nullable = false)
var nextAttemptAt: Instant = Instant.now(),
@Column(name = "last_error", columnDefinition = "TEXT")
var lastError: String? = null,
) {
fun markPublished(streamId: String, now: Instant = Instant.now()) {
redisStreamId = streamId
publishedAt = now
lastError = null
}
fun markFailed(error: Throwable) {
attempts += 1
lastError = (error.message ?: error::class.java.simpleName).take(1000)
val delaySeconds = minOf(60L, 1L shl attempts.coerceAtMost(6))
nextAttemptAt = Instant.now().plus(Duration.ofSeconds(delaySeconds))
}
}outbox relay
@Component
class PosOutboxRelay(...) {
@Scheduled(fixedDelay = 500)
@Transactional
fun publishPendingEvents() {
val events = outboxRepository.lockPublishable(batchSize)
if (events.isEmpty()) return
events.forEach { event ->
try {
val streamId = publishToRedis(event)
event.markPublished(streamId)
} catch (e: Exception) {
event.markFailed(e)
// log ...
}
}
}
}
interface PosEventOutboxRepository : JpaRepository<PosEventOutbox, UUID> {
@Query(
value = """
select *
from pos_event_outbox
where published_at is null
and next_attempt_at <= now()
order by created_at
limit :limit
for update skip locked
""",
nativeQuery = true,
)
fun lockPublishable(@Param("limit") limit: Int): List<PosEventOutbox>
}
consideration
lockPublishable()이 어떤 row를 선택하는지- 여러 relay instance가 같은 row를 동시에 처리할 수 있는지
- Redis publish 성공 후
markPublished()가 같은 DB transaction에서 commit되는지 - Redis publish는 성공했는데 DB commit이 실패할 수 있는지
- 실패 row가 영구적으로 retry되는지, 최대 횟수가 있는지
- poison event가 계속 retry되면 어떻게 되는지
한계
- 현재 구조는 유실을 줄이지만 exactly-once 발행은 아니다
- 아래 경우에 중복 발행 가능
Redis XADD 성공
→ 애플리케이션 종료 또는 DB commit 실패
→ Outbox에는 published_at이 여전히 null
→ 재시작 후 같은 Outbox 이벤트를 다시 XADD대처법: FE deduplication
0. BE: SSE에 Outbox id inject
private fun enqueueRecord(chanKey: String, streamKey: String, message: MapRecord<String, String, String>) {
try {
val type = message.value["type"] ?: "UNKNOWN"
val payload = message.value["payload"] ?: "{}"
val id = message.id.toString()
if (type !in knownEventNames) {
logger.warn("Unknown SSE event type: type={} id={} streamKey={}", type, id, streamKey)
enqueue(chanKey, SseMsg("RELOAD_REQUIRED", "{}", id))
return
}
enqueue(chanKey, SseMsg(type, payloadWithOutboxEventId(payload, message.value["eventId"]), id))
} catch (e: Exception) {
logger.error("Error processing redis message", e)
}
}
private fun payloadWithOutboxEventId(payload: String, outboxEventId: String?): String {
if (outboxEventId.isNullOrBlank()) return payload
return try {
val node = objectMapper.readTree(payload)
if (!node.isObject) return payload
(node as com.fasterxml.jackson.databind.node.ObjectNode)
.put("__outboxEventId", outboxEventId)
objectMapper.writeValueAsString(node)
} catch (_: Exception) {
payload
}
}1. FE: 추출 및 dedupe
function readOutboxDedupeEventId(value: unknown): string | undefined {
if (!value || typeof value !== "object") return undefined;
const outboxEventId = (value as { __outboxEventId?: unknown }).__outboxEventId;
return typeof outboxEventId === "string" && outboxEventId.trim()
? `outbox:${outboxEventId}`
: undefined;
}posListenerMiddleware.startListening({
predicate: (action): action is PosRealtimeAction =>
typeof action.type === "string" && action.type.startsWith("pos/realtime/"),
effect: async (action, api) => {
const meta = (action as { meta?: { streamEventId?: string; dedupeEventId?: string } }).meta;
const streamEventId = normalizeRedisStreamId(meta?.streamEventId);
const dedupeEventId = meta?.dedupeEventId;
const dedupeId = dedupeEventId ?? (streamEventId !== "0-0" ? `stream:${streamEventId}` : undefined);
const isReplayableStreamEvent = streamEventId !== "0-0";
if (dedupeId) {
const alreadyApplied = getPosState(api).stream.recentAppliedDedupeIds.includes(dedupeId);
if (alreadyApplied) {
if (isReplayableStreamEvent) {
api.dispatch(
posActions.markStreamEventApplied({
eventId: streamEventId,
dedupeId,
snapshotScopes: getSnapshotScopesForRealtimeAction(action.type),
}),
);
}
return;
}
}
switch (action.type) {
case "pos/realtime/payment/open": {
// ... 이하 생략 ...0x90. PS
0x91. SSE authentication and connection
SSE는 요청 헤더에 커스텀 Authorization을 싣기 어렵다. [(draft yet): SSE; Server-Sent Event]
1. 현재 [(draft yet): Authentication & Authorization]
- kakao Login + Standard Token-based Authorization Scheme을 사용 중
2. ticket 기반 연결
- FE가
/api/pos/{storeId}/notifications/ticket을 먼저 호출한다. - BE는 로그인 사용자와 매장 권한을
storeService.assertHost(storeId, user.id)로 확인한다. - 권한이 맞으면 Redis에 30초 TTL의 1회용 티켓을 저장한다.
- FE는
/api/pos/{storeId}/notifications/subscribe?ticket=...&connectionId=...&lastEventId=...로 EventSource를 연다. subscribe엔드포인트는 Spring Security에서permitAll이지만, 실제 접근 제어는 ticket 검증으로 처리한다.- BE는 티켓을
getAndDelete로 꺼내므로 같은 티켓은 재사용할 수 없다.
3. 의도
- 이 구조의 의도는 EventSource의 헤더 제약을 우회하면서도, SSE URL만 알아서는 접속할 수 없게 만드는 것이다.
- 주의할 점은 티켓 TTL이 30초라서, 티켓 발급 후 구독까지의 시퀀스가 빠르게 이어져야 한다는 것이다. 프론트는 연결 시도마다 새 티켓을 발급받는다.
- 실제 서버 코드는 티켓을 조회하는 동시에 삭제한다. 이 때문에 같은 티켓으로 두 번 연결할 수 없다.
fun createTicket(storeId: UUID): String {
val ticket = UUID.randomUUID().toString().replace("-", "")
redisTemplate.opsForValue().set("ticket:$ticket", storeId.toString(), Duration.ofSeconds(30))
return ticket
}
fun subscribe(
expectedStoreId: UUID,
ticket: String,
connectionId: String?,
lastEventId: String?
): SseEmitter {
val storeIdStr = redisTemplate.opsForValue().getAndDelete("ticket:$ticket")
?: throw BusinessException(ErrorCode.INVALID_TICKET)
if (storeIdStr != expectedStoreId.toString()) {
throw BusinessException(ErrorCode.INVALID_TICKET)
}
val emitter = SseEmitter(60 * 60 * 1000L)
...
}- 프론트는 EventSource 연결 전에 항상 ticket API를 호출하고,
lastEventId가0-0이 아닐 때만 구독 URL에 붙인다.
async function getPosTicket(storeId: string): Promise<string> {
const res = await apiRequest(`/api/pos/${storeId}/notifications/ticket`, {
method: "GET",
credentials: "include",
});
...
}
function getSseUrl(
storeId: string,
ticket: string,
connectionId: string,
lastEventId: string,
): string {
const params = new URLSearchParams({ ticket, connectionId });
if (lastEventId !== "0-0") params.set("lastEventId", lastEventId);
return `/api/pos/${storeId}/notifications/subscribe?${params.toString()}`;
}