트랜잭션 사용하기

트랜잭션 사용하기 (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를 사용하려면 다음 단계를 완료하세요.

  1. 트랜잭션 활성화. 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 트랜잭션은 비활성화되어 있어요.

  1. 트랜잭션 코디네이터 메타데이터 초기화. 트랜잭션 코디네이터는 파티션 토픽의 장점(예: 로드 밸런싱)을 활용할 수 있어요.
bin/pulsar initialize-transaction-coordinator-metadata -cs 127.0.0.1:2181 -c standalone

출력:

Transaction coordinator metadata setup success
  1. Pulsar 클라이언트를 만들고 트랜잭션을 활성화. 클라이언트는 시스템 토픽에서 트랜잭션 코디네이터를 알아야 하므로, 클라이언트 역할이 시스템 네임스페이스 pulsar/system에 대해 produce/consume 권한을 가지고 있는지 확인하세요.

  2. 프로듀서와 컨슈머 생성.

  3. 메시지 생산 및 수신.

  4. 트랜잭션 생성.

  5. 트랜잭션으로 메시지 생산 및 ack. 참고로 현재 메시지는 누적(cumulatively)이 아니라 개별적으로(individually) ack할 수 있어요.

  6. 트랜잭션 종료.

: 아래 코드 스니펫은 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

더 알아보기 (Learn more)