Pulsar 트랜잭션, 왜 필요한가요?

Pulsar 트랜잭션, 왜 필요한가요?

Pulsar 트랜잭션(txn)은 이벤트 스트리밍 애플리케이션이 메시지를 소비하고, 처리하고, 생산하는 작업을 하나의 원자적(atomic) 연산으로 묶을 수 있게 해줘요. 이 기능이 왜 개발됐는지 그 이유를 정리해서 살펴볼게요. 핵심은 "정확히 한 번 처리"라는 강한 보장을 여러 파티션에 걸쳐 얻고 싶다는 요구에서 출발했어요.

출처: 문서

본문

Pulsar 트랜잭션(txn)은 이벤트 스트리밍 애플리케이션이 메시지를 하나의 원자적 연산으로 소비하고, 처리하고, 생산할 수 있게 해줘요. 이 기능을 개발한 이유는 다음과 같이 정리할 수 있어요.

스트림 처리에 대한 수요 (Demand of stream processing)

스트림 처리의 성장과 함께, 더 강한 처리 보장을 원하는 스트림 처리 애플리케이션에 대한 수요도 함께 커졌어요. 예를 들어 금융 업계에서는 스트림 처리 엔진으로 사용자의 대변(debit)과 차변(credit)을 처리해요. 이런 유형의 사용 사례는 예외 없이 모든 메시지가 정확히 한 번 처리되는 것을 요구해요.

다시 말해서, 스트림 처리 애플리케이션이 메시지 A를 소비하고 그 결과를 메시지 B(B = f(A))로 생산한다면, 정확히 한 번 처리 보장은 "B가 성공적으로 생산된 경우에만 A가 소비된 것으로 표시될 수 있고, 그 반대도 마찬가지"임을 의미해요.

Pulsar 트랜잭션 API는 스트림 처리의 메시지 전달 의미와 처리 보장을 강화해요. 스트림 처리 애플리케이션이 메시지를 하나의 원자적 연산으로 소비하고, 처리하고, 생산할 수 있게 해주죠. 즉, 트랜잭션에 속한 메시지 배치가 많은 토픽 파티션에서 수신되고, 생산되고, 확인될 수 있어요. 트랜잭션에 포함된 모든 작업은 하나의 단위로 성공하거나 실패해요.

멱등 프로듀서의 한계 (Limitation of idempotent producer)

데이터 손실이나 중복을 피하는 것은 Pulsar 멱등 프로듀서(idempotent producer)로 해결할 수 있지만, 이 방식은 여러 파티션에 걸친 쓰기에 대해서는 보장을 제공하지 못해요.

Pulsar에서 가장 높은 수준의 메시지 전달 보장은 단일 파티션에서 정확히 한 번 의미와 함께 멱등 프로듀서를 사용하는 것이에요. 즉, 각 메시지가 데이터 손실이나 중복 없이 정확히 한 번 영속된다는 뜻이죠. 하지만 이 솔루션에는 몇 가지 한계가 있어요.

  • 단조 증가하는 시퀀스 ID 때문에, 이 솔루션은 단일 파티션과 단일 프로듀서 세션 내에서만 동작해요(즉, 메시지 하나를 생산할 때만 해당). 그래서 하나의 파티션이나 여러 파티션에 여러 메시지를 생산할 때는 원자성이 없어요. 이런 경우 메시지 생산·수신 과정에서 장애가 발생하면(예: 클라이언트/브로커/bookie 크래시, 네트워크 장애 등) 메시지가 다시 처리되고 재전달되면서 데이터 손실이나 중복이 발생할 수 있어요. 프로듀서 입장에서 보면, 프로듀서가 메시지 전송을 재시도하면 일부 메시지가 여러 번 영속되고, 재시도하지 않으면 일부 메시지는 한 번 영속되고 다른 메시지는 손실돼요. 컨슈머 입장에서 보면, 컨슈머는 브로커가 메시지를 받았는지 알 수 없으므로 ack 재전송을 시도하지 않을 수 있고, 그 결과 중복 메시지를 받게 돼요.
  • 마찬가지로 Pulsar Function의 경우에도, 멱등 함수가 단일 이벤트에 대해 정확히 한 번 의미를 보장할 뿐, 여러 이벤트를 처리하거나 여러 결과를 생산하는 것을 정확히 보장하지는 못해요. 예를 들어 함수가 여러 이벤트를 받아 하나의 결과를 생산한다면(예: 윈도우 함수), 함수는 결과를 생산한 뒤 들어오는 메시지를 확인(ack)하기 전에, 심지어 개별 이벤트를 확인하는 사이에 실패할 수 있어요. 그러면 들어오는 메시지의 전부(또는 일부)가 재전달·재처리되고 새 결과가 생산돼요. 하지만 많은 시나리오는 여러 파티션과 세션에 걸친 원자적 보장을 필요로 해요.
  • 컨슈머는 메시지를 한 번만 확인(ack)하기 위해 더 많은 메커니즘에 의존해야 해요. 예를 들어 컨슈머는 MessageID를 그 ack 상태와 함께 저장해야 해요. 토픽이 언로드된 후, 토픽이 다시 로드될 때 서브스크립션이 이 MessageID의 ack 상태를 메모리에서 복구할 수 있어야 하죠.

더 알아보기 (Learn more)