Java 클라이언트(V5)

Java 클라이언트(V5) (Java client V5)

V5 Java 클라이언트는 스케일러블 토픽을 위해 만들어진 클라이언트 SDK예요. 세 가지 목적별 컨슈머 타입과 현대적이고 동기(sync) 우선의 API를 제공하고, 기존의 파티셔닝·비파티셔닝 토픽에서도 동작하므로 토픽을 마이그레이션하기 전에 먼저 도입할 수 있어요. Java 17이 필요하다는 점만 확인하고 시작하면 돼요.

출처: 문서

본문

V5 Java 클라이언트는 스케일러블 토픽을 위해 만들어진 클라이언트 SDK예요. 세 가지 목적별 컨슈머 타입과 현대적이고 동기 우선의 API를 제공하며, 기존의 파티셔닝·비파티셔닝 토픽에서도 동작해요 — 그래서 어떤 토픽도 마이그레이션하기 전에 먼저 도입할 수 있죠.

현재 Java 클라이언트와 어떻게 비교되는지는 Java 클라이언트 SDK를, 메시징 모델과 더 깊은 API 살펴보기는 Messaging을 참고하세요. 아래 언급된 모든 타입은 API 참조로 연결돼요.

V5 클라이언트는 Java 17을 요구해요. 다른 언어 SDK의 스케일러블 토픽 지원은 계획 중이에요. 오늘날 V5 클라이언트는 Java에서 사용할 수 있어요.

설치 (Install)

V5 클라이언트는 Maven Central에 pulsar-client-v5로 게시돼요(5.0.0-M2부터 사용 가능). API는 org.apache.pulsar.client.api.v5 패키지 아래에 있어요.

Maven

<dependency>
  <groupId>org.apache.pulsar</groupId>
  <artifactId>pulsar-client-v5</artifactId>
  <version>5.0.0-M2</version>
</dependency>

Gradle

dependencies {
    implementation "org.apache.pulsar:pulsar-client-v5:5.0.0-M2"
}

클라이언트 만들기 (Create a client)

PulsarClient는 API의 진입점이에요. 하나를 만들고 모든 프로듀서·컨슈머에서 공유하며, 종료 시 닫으면 돼요. PulsarClient.builder()는 연결을 구성하는 PulsarClientBuilder를 반환해요.

import org.apache.pulsar.client.api.v5.PulsarClient;

PulsarClient client = PulsarClient.builder()
        .serviceUrl("pulsar://localhost:6650")
        .build();
// ... create producers and consumers ...
client.close();

서비스 URL은 pulsar://(또는 pulsar+ssl://) 체계를 사용해요. 인증, TLS, 연산 타임아웃, 메모리 제한은 모두 같은 빌더에서 설정돼요.

메시지 생산 (Produce messages)

클라이언트에서 스키마와 토픽으로 프로듀서를 만들어요. client.newProducer(schema)ProducerBuilder를 반환하고, create()를 호출하면 프로듀서를 얻어요. 각 producer.newMessage()는 전송 전에 키, 값, 기타 메시지별 속성을 설정하는 MessageBuilder를 반환해요.

Producer<String> producer = client.newProducer(Schema.string())
        .topic("topic://public/default/orders")
        .create();

producer.newMessage()
        .key("user-123")
        .value("order placed")
        .send();

비차단 전송을 하려면 producer.async()CompletableFuture를 반환하는 연산들을 가진 AsyncProducer를 돌려줘요.

var async = producer.async();
async.newMessage().value("order placed").send()
        .thenAccept(id -> System.out.println("sent " + id));

async.flush().join();

메시지 소비 (Consume messages)

V5 클라이언트는 기존의 네 가지 서브스크립션 타입을 세 가지 컨슈머 타입으로 대체해요. 어떻게 소비할 것인지에 따라 하나를 고르면 돼요. 전체 의미는 Consumers를 참고하세요. 각 컨슈머는 MessageId로 식별되는 메시지를 건네줘요.

스트림 컨슈머 (Stream consumer)

StreamConsumer는 누적 ack와 함께 순서대로 소비를 제공해요 — 메시지를 ack하면 그 앞의 모든 메시지도 함께 ack돼요.

StreamConsumer<String> consumer = client.newStreamConsumer(Schema.string())
        .topic("topic://public/default/orders")
        .subscriptionName("my-sub")
        .subscribe();

Message<String> msg = consumer.receive(Duration.ofSeconds(5)); // null on timeout
if (msg != null) {
    process(msg.value());
    consumer.acknowledgeCumulative(msg.id());
}

큐 컨슈머 (Queue consumer)

QueueConsumer는 개별 ack, 부정 ack, 데드 레터(dead-letter) 지원과 함께 병렬 소비를 제공해요.

QueueConsumer<String> consumer = client.newQueueConsumer(Schema.string())
        .topic("topic://public/default/orders")
        .subscriptionName("workers")
        .subscribe();

Message<String> msg = consumer.receive(Duration.ofSeconds(5));
if (msg != null) {
    try {
        process(msg.value());
        consumer.acknowledge(msg.id());
    } catch (Exception e) {
        consumer.negativeAcknowledge(msg.id());
    }
}

체크포인트 컨슈머 (Checkpoint consumer)

CheckpointConsumer는 자신의 위치를 추적하는 스트림 처리 엔진에 적합해요 — 서브스크립션이나 ack가 없어요. 직렬화 가능한 Checkpoint를 캡처하고 나중에 그것으로 복원할 수 있어요.

CheckpointConsumer<String> consumer = client.newCheckpointConsumer(Schema.string())
        .topic("topic://public/default/orders")
        .startPosition(Checkpoint.earliest())   // earliest(), latest(), or a saved checkpoint
        .create();

Message<String> msg = consumer.receive(Duration.ofSeconds(5));
byte[] state = consumer.checkpoint().toByteArray();   // persist externally
// resume later with: .startPosition(Checkpoint.fromByteArray(state))

스키마 (Schemas)

프로듀서나 컨슈머를 만들 때 스키마를 전달해요. V5 스키마 팩토리는 소문자 메서드예요.

Schema.string()            // StringSchema
.json(Order.class)         // JSON-encoded POJO
.avro(Order.class)         // Avro-encoded POJO

원시 타입 팩토리(Schema.int32(), Schema.bool(), Schema.bytes(), …)와 Schema.protobuf(...)도 사용할 수 있어요.

트랜잭션 (Transactions)

트랜잭션은 메시지를 생산하고 소비된 메시지를 ack하는 것을 원자적으로 할 수 있게 해줘요. client.newTransaction()으로 시작하고, 메시지 빌더의 .transaction(txn)으로 생산을, 인자 두 개짜리 acknowledge로 ack를 바인딩한 다음 커밋하거나 중단해요.

Transaction txn = client.newTransaction();
try {
    producer.newMessage().transaction(txn).value(result).send();
    consumer.acknowledge(msg.id(), txn);
    txn.commit();
} catch (Exception e) {
    txn.abort();
}

다음으로? (What's next)

더 알아보기 (Learn more)