트랜잭션 API
트랜잭션 API (Transactions API)
트랜잭션의 모든 메시지는 트랜잭션 커밋 후에만 컨슈머가 사용할 수 있어요. 트랜잭션이 중단(abort)되면 이 트랜잭션의 모든 쓰기와 확인(acknowledgment)이 롤백돼요. 이번에는 Pulsar 트랜잭션 API를 실제로 설정하고 사용하는 방법을 단계별로 살펴볼게요.
출처: 문서
본문
사전 준비 (Prerequisites)
- Pulsar에서 트랜잭션을 활성화하려면
broker.conf파일이나standalone.conf파일의 파라미터를 설정해야 해요.
transactionCoordinatorEnabled=true
- 트랜잭션 코디네이터가 파티셔닝된 토픽의 이점(예: 로드밸런싱)을 활용할 수 있도록 트랜잭션 코디네이터 메타데이터를 초기화해요.
bin/pulsar initialize-transaction-coordinator-metadata -cs 127.0.0.1:2181 -c standalone
트랜잭션 코디네이터 메타데이터를 초기화한 뒤에 트랜잭션 API를 사용할 수 있어요. 사용 가능한 API는 다음과 같아요.
Pulsar 클라이언트 초기화 (Initialize Pulsar client)
트랜잭션 클라이언트에 대해 트랜잭션을 활성화하고 트랜잭션 코디네이터 클라이언트를 초기화할 수 있어요.
PulsarClient pulsarClient = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.enableTransaction(true)
.build();
트랜잭션 시작 (Start transactions)
다음과 같은 방법으로 트랜잭션을 시작할 수 있어요.
Transaction txn = pulsarClient
.newTransaction()
.withTransactionTimeout(5, TimeUnit.MINUTES)
.build()
.get();
트랜잭션 메시지 생성 (Produce transaction messages)
새 트랜잭션 메시지를 생성할 때는 트랜잭션 파라미터가 필요해요. Pulsar의 트랜잭션 메시지 의미론은 read-committed이므로, 컨슈머는 트랜잭션이 커밋되기 전에는 진행 중인 트랜잭션 메시지를 받을 수 없어요.
producer.newMessage(txn).value("Hello Pulsar Transaction".getBytes()).sendAsync();
트랜잭션으로 메시지 확인 (Acknowledge the messages with the transaction)
트랜잭션 확인은 트랜잭션 파라미터가 필요해요. 트랜잭션 확인은 메시지 상태를 보류-ack(pending-ack) 상태로 표시해요. 트랜잭션이 커밋되면 보류-ack 상태는 ack 상태가 돼요. 트랜잭션이 중단되면 보류-ack 상태는 미확인(unacknowledged) 상태가 돼요.
Message<byte[]> message = consumer.receive();
consumer.acknowledgeAsync(message.getMessageId(), txn);
트랜잭션 커밋 (Commit transactions)
트랜잭션이 커밋되면 컨슈머가 트랜잭션 메시지를 받고 보류-ack 상태가 ack 상태가 돼요.
txn.commit().get();
트랜잭션 중단 (Abort transaction)
트랜잭션이 중단되면 트랜잭션 확인이 취소되고 보류-ack 메시지가 재전달돼요.
txn.abort().get();
예제 (Example)
다음 예제는 트랜잭션에서 메시지가 어떻게 처리되는지 보여줘요.
PulsarClient pulsarClient = PulsarClient.builder()
.serviceUrl(getPulsarServiceList().get(0).getBrokerServiceUrl())
.statsInterval(0, TimeUnit.SECONDS)
.enableTransaction(true)
.build();
String sourceTopic = "public/default/source-topic";
String sinkTopic = "public/default/sink-topic";
Producer<String> sourceProducer = pulsarClient
.newProducer(Schema.STRING)
.topic(sourceTopic)
.create();
sourceProducer.newMessage().value("hello pulsar transaction").sendAsync();
Consumer<String> sourceConsumer = pulsarClient
.newConsumer(Schema.STRING)
.topic(sourceTopic)
.subscriptionName("test")
.subscriptionType(SubscriptionType.Shared)
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscribe();
Producer<String> sinkProducer = pulsarClient
.newProducer(Schema.STRING)
.topic(sinkTopic)
.sendTimeout(0, TimeUnit.MILLISECONDS)
.create();
Transaction txn = pulsarClient
.newTransaction()
.withTransactionTimeout(5, TimeUnit.MINUTES)
.build()
.get();
// source message acknowledgment and sink message produce belong to one transaction,
// they are combined into an atomic operation.
Message<String> message = sourceConsumer.receive();
sourceConsumer.acknowledgeAsync(message.getMessageId(), txn);
sinkProducer.newMessage(txn).value("sink data").sendAsync();
txn.commit().get();
트랜잭션에서 배치 메시지 활성화 (Enable batch messages in transactions)
트랜잭션에서 배치 메시지를 활성화하려면 배치 인덱스 확인 기능을 활성화해야 해요. 트랜잭션 ack는 배치 인덱스 확인이 충돌하는지 검사해요.
배치 인덱스 확인을 활성화하려면 broker.conf 또는 standalone.conf 파일에서 acknowledgmentAtBatchIndexLevelEnabled를 true로 설정해야 해요.
acknowledgmentAtBatchIndexLevelEnabled=true
그리고 컨슈머 빌더에서 enableBatchIndexAcknowledgment(true) 메서드를 호출해야 해요.
Consumer<byte[]> sinkConsumer = pulsarClient
.newConsumer()
.topic(transferTopic)
.subscriptionName("sink-topic")
.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)
.subscriptionType(SubscriptionType.Shared)
.enableBatchIndexAcknowledgment(true) // enable batch index acknowledgment
.subscribe();
더 알아보기 (Learn more)
- 트랜잭션의 기본 개념은 Transactions 문서에서 확인해요.
- 트랜잭션이 보장하는 의미론은 Transactions Guarantee 문서를 참고해요.
- 트랜잭션을 활성화하는 전체 가이드는 Pulsar transactions 문서를 봐요.
- 구성 파라미터의 기본값은 Transactions 문서에서 살펴봐요.