트랜잭션 사용하기
트랜잭션 사용하기 (Use transactions)
Pulsar 트랜잭션은 기본적으로 서버 쪽이자 프로토콜 수준의 기능이에요. 이 튜토리얼에서는 Java 클라이언트에서 Pulsar 트랜잭션 API를 사용해 메시지를 보내고 받는 전체 과정을 단계별로 안내할게요. 트랜잭션이 첫 도입이라면 이 문서를 처음부터 따라가면 됩니다.
출처: 문서
본문
Pulsar 트랜잭션은 주로 서버 쪽이자 프로토콜 수준의 기능이에요. 이 튜토리얼은 Java 클라이언트에서 Pulsar 트랜잭션 API를 사용해 메시지를 보내고 받는 모든 단계를 안내해요.
현재 Pulsar 트랜잭션 API는 Pulsar 2.8.0 이상 버전에서 사용할 수 있어요. Java, Go, .NET 클라이언트에서만 사용할 수 있으니 참고하세요.
사전 요구사항 (Prerequisites)
- Pulsar 2.8.0 이상 버전을 시작하세요.
단계 (Steps)
Pulsar 트랜잭션 API를 사용하려면 다음 단계를 완료하세요.
- 트랜잭션 활성화.
broker.conf또는standalone.conf파일에서 다음 설정을 지정할 수 있어요.
//mandatory configuration, used to enable transaction coordinator
transactionCoordinatorEnabled=true
//mandatory configuration, used to create systemTopic used for transaction buffer snapshot
systemTopicEnabled=true
참고: 기본적으로 Pulsar 트랜잭션은 비활성화되어 있어요.
- 트랜잭션 코디네이터 메타데이터 초기화. 트랜잭션 코디네이터는 파티션 토픽의 장점(예: 로드 밸런싱)을 활용할 수 있어요.
bin/pulsar initialize-transaction-coordinator-metadata -cs 127.0.0.1:2181 -c standalone
출력:
Transaction coordinator metadata setup success
-
Pulsar 클라이언트를 만들고 트랜잭션을 활성화. 클라이언트는 시스템 토픽에서 트랜잭션 코디네이터를 알아야 하므로, 클라이언트 역할이 시스템 네임스페이스
pulsar/system에 대해 produce/consume 권한을 가지고 있는지 확인하세요. -
프로듀서와 컨슈머 생성.
-
메시지 생산 및 수신.
-
트랜잭션 생성.
-
트랜잭션으로 메시지 생산 및 ack. 참고로 현재 메시지는 누적(cumulatively)이 아니라 개별적으로(individually) ack할 수 있어요.
-
트랜잭션 종료.
팁: 아래 코드 스니펫은 3단계~8단계의 예시예요.
PulsarClient client = PulsarClient.builder()
// Step 3: create a Pulsar client and enable transactions.
.enableTransaction(true)
.serviceUrl(jct.serviceUrl)
.build();
// Step 4: create three producers to produce messages to input and output topics.
ProducerBuilder<String> producerBuilder = client.newProducer(Schema.STRING);
Producer<String> inputProducer = producerBuilder.topic(inputTopic)
.sendTimeout(0, TimeUnit.SECONDS).create();
Producer<String> outputProducerOne = producerBuilder.topic(outputTopicOne)
.sendTimeout(0, TimeUnit.SECONDS).create();
Producer<String> outputProducerTwo = producerBuilder.topic(outputTopicTwo)
.sendTimeout(0, TimeUnit.SECONDS).create();
// Step 4: create three consumers to consume messages from input and output topics.
Consumer<String> inputConsumer = client.newConsumer(Schema.STRING)
.subscriptionName("your-subscription-name").topic(inputTopic).subscribe();
Consumer<String> outputConsumerOne = client.newConsumer(Schema.STRING)
.subscriptionName("your-subscription-name").topic(outputTopicOne).subscribe();
Consumer<String> outputConsumerTwo = client.newConsumer(Schema.STRING)
.subscriptionName("your-subscription-name").topic(outputTopicTwo).subscribe();
int count = 2;
// Step 5: produce messages to input topics.
for (int i = 0; i < count; i++) {
inputProducer.send("Hello Pulsar! count : " + i);
}
// Step 5: consume messages and produce them to output topics with transactions.
for (int i = 0; i < count; i++) {
// Step 5: the consumer successfully receives messages.
Message<String> message = inputConsumer.receive();
// Step 6: create transactions.
// The transaction timeout is specified as 10 seconds.
// If the transaction is not committed within 10 seconds, the transaction is automatically aborted.
Transaction txn = null;
try {
txn = client.newTransaction()
.withTransactionTimeout(10, TimeUnit.SECONDS).build().get();
// Step 6: you can process the received message with your use case and business logic.
// Step 7: the producers produce messages to output topics with transactions
outputProducerOne.newMessage(txn).value("Hello Pulsar! outputTopicOne count : " + i).send();
outputProducerTwo.newMessage(txn).value("Hello Pulsar! outputTopicTwo count : " + i).send();
// Step 7: the consumers acknowledge the input message with the transactions *individually*.
inputConsumer.acknowledgeAsync(message.getMessageId(), txn).get();
// Step 8: commit transactions.
txn.commit().get();
} catch (ExecutionException e) {
if (!(e.getCause() instanceof PulsarClientException.TransactionConflictException)) {
// If TransactionConflictException is not thrown,
// you need to redeliver or negativeAcknowledge this message,
// or else this message will not be received again.
inputConsumer.negativeAcknowledge(message);
}
// If a new transaction is created,
// then the old transaction should be aborted.
if (txn != null) {
txn.abort();
}
}
}
// Final result: consume messages from output topics and print them.
for (int i = 0; i < count; i++) {
Message<String> message = outputConsumerOne.receive();
System.out.println("Receive transaction message: " + message.getValue());
}
for (int i = 0; i < count; i++) {
Message<String> message = outputConsumerTwo.receive();
System.out.println("Receive transaction message: " + message.getValue());
}
// Step 3: create a Pulsar client and enable transactions.
client, err := pulsar.NewClient(pulsar.ClientOptions{
URL: "<serviceUrl>",
EnableTransaction: true,
})
if err != nil {
log.Fatalf("create client fail, err = %v", err)
}
defer client.Close()
// Step 4: create three producers to produce messages to input and output topics.
inputTopic := "inputTopic"
outputTopicOne := "outputTopicOne"
outputTopicTwo := "outputTopicTwo"
subscriptionName := "your-subscription-name"
inputProducer, _ := client.CreateProducer(pulsar.ProducerOptions{
Topic: inputTopic,
SendTimeout: 0,
})
defer inputProducer.Close()
outputProducerOne, _ := client.CreateProducer(pulsar.ProducerOptions{
Topic: outputTopicOne,
SendTimeout: 0,
})
defer outputProducerOne.Close()
outputProducerTwo, _ := client.CreateProducer(pulsar.ProducerOptions{
Topic: outputTopicTwo,
SendTimeout: 0,
})
defer outputProducerTwo.Close()
// Step 4: create three consumers to consume messages from input and output topics.
inputConsumer, _ := client.Subscribe(pulsar.ConsumerOptions{
Topic: inputTopic,
SubscriptionName: subscriptionName,
})
defer inputConsumer.Close()
outputConsumerOne, _ := client.Subscribe(pulsar.ConsumerOptions{
Topic: outputTopicOne,
SubscriptionName: subscriptionName,
})
defer outputConsumerOne.Close()
outputConsumerTwo, _ := client.Subscribe(pulsar.ConsumerOptions{
Topic: outputTopicTwo,
SubscriptionName: subscriptionName,
})
defer outputConsumerTwo.Close()
// Step 5: produce messages to input topics.
ctx := context.Background()
count := 2
for i := 0; i < count; i++ {
inputProducer.Send(ctx, &pulsar.ProducerMessage{
Payload: []byte(fmt.Sprintf("Hello Pulsar! count : %d", i)),
})
}
// Step 5: consume messages and produce them to output topics with transactions.
for i := 0; i < count; i++ {
// Step 5: the consumer successfully receives messages.
message, err := inputConsumer.Receive(ctx)
if err != nil {
log.Printf("receive message from %s fail, err = %v", inputTopic, err)
continue
}
// Step 6: create transactions.
// The transaction timeout is specified as 10 seconds.
// If the transaction is not committed within 10 seconds, the transaction is automatically aborted.
txn, err := client.NewTransaction(10 * time.Second)
if err != nil {
log.Printf("create txn fail, err = %v", err)
continue
}
// Step 6: you can process the received message with your use case and business logic.
// processMessage(message)
// Step 7: the producers produce messages to output topics with transactions
_, err = outputProducerOne.Send(context.Background(), &pulsar.ProducerMessage{
Transaction: txn,
Payload: []byte(fmt.Sprintf("Hello Pulsar! outputTopicOne count : %d", i)),
})
if err != nil {
log.Printf("send to producerOne fail %v", err)
txn.Abort(ctx)
}
_, err = outputProducerTwo.Send(context.Background(), &pulsar.ProducerMessage{
Transaction: txn,
Payload: []byte(fmt.Sprintf("Hello Pulsar! outputTopicTwo count : %d", i)),
})
if err != nil {
log.Printf("send to producerTwo fail %v", err)
txn.Abort(ctx)
}
// Step 7: the consumers acknowledge the input message with the transactions *individually*.
err = inputConsumer.AckWithTxn(message, txn)
if err != nil {
log.Printf("ack message fail %v", err)
txn.Abort(ctx)
}
// Step 8: commit transactions.
err = txn.Commit(ctx)
if err != nil {
log.Printf("commit txn fail %v", err)
}
}
// Final result: consume messages from output topics and print them.
for i := 0; i < count; i++ {
message, _ := outputConsumerOne.Receive(ctx)
log.Printf("Receive transaction message: %s", string(message.Payload()))
}
for i := 0; i < count; i++ {
message, _ := outputConsumerTwo.Receive(ctx)
log.Printf("Receive transaction message: %s", string(message.Payload()))
}
출력:
Receive transaction message: Hello Pulsar! count : 1
Receive transaction message: Hello Pulsar! count : 2
Receive transaction message: Hello Pulsar! count : 1
Receive transaction message: Hello Pulsar! count : 2
관련 주제 (Related topics)
- 트랜잭션과 함께 쓸 수 있는 더 많은 기능을 보려면 Pulsar 트랜잭션 - 고급 기능 문서를 확인하세요.
더 알아보기 (Learn more)
- Pulsar 트랜잭션이란? — 트랜잭션의 기본 개념을 살펴봐요.
- Pulsar 트랜잭션, 왜 필요한가요? — 트랜잭션 도입 배경을 알아봐요.
- 트랜잭션 동작 원리 — 트랜잭션의 내부 동작을 살펴봐요.
- 트랜잭션 고급 기능 — 배치 메시지 ack, 인증, 정확히 한 번 의미를 확인해요.
- 트랜잭션 모니터링 — 트랜잭션 메트릭을 살펴봐요.