Java 클라이언트를 V5로 마이그레이션하기

Java 클라이언트를 V5로 마이그레이션하기 (Migrate to V5)

이 가이드는 기존 Java 애플리케이션을 현재 클라이언트 SDK(org.apache.pulsar.client.api)에서 스케일러블 토픽에 사용되는 V5 클라이언트 SDK(org.apache.pulsar.client.api.v5)로 마이그레이션하는 방법을 설명해요. 꼭 해야 하는 것은 아니고, 스케일러블 토픽 또는 V5 API가 필요할 때 진행하면 돼요. 두 SDK가 같은 JVM 안에서 나란히 돌 수 있으므로 한 번에 하나씩 점진적으로 옮겨갈 수 있어요.

출처: 문서

본문

이 가이드는 기존 Java 애플리케이션을 현재 클라이언트 SDK(org.apache.pulsar.client.api)에서 스케일러블 토픽에 사용되는 V5 클라이언트 SDK(org.apache.pulsar.client.api.v5)로 마이그레이션하는 방법을 설명해요.

마이그레이션할 필요는 없어요. 현재 SDK는 완전히 지원되며, 자바가 아닌 애플리케이션이나 스케일러블 토픽이 필요 없는 애플리케이션, 그리고 비영속 토픽(V5 클라이언트가 지원하지 않음)에는 여전히 올바른 선택이에요. 스케일러블 토픽이나 V5 API를 원할 때 Java 애플리케이션을 마이그레이션하세요.

마이그레이션 동작 방식 (How migration works)

두 SDK는 독립적이며 같은 JVM 안에서 나란히 실행될 수 있어요. 그래서 한 번에 모두 바꾸지 않고, 프로듀서나 컨슈머 하나씩 점진적으로 마이그레이션할 수 있어요. 대표적인 경로는 다음과 같아요.

  • pulsar-client-v5 의존성을 추가해요(아직 마이그레이션하지 않은 v4 코드를 위한 pulsar-client-original도 함께 — Dependencies 참고).
  • 프로듀서와 컨슈머를 V5 API로 옮겨요. V5 클라이언트는 기존 persistent:// 토픽에서 동작하므로, 어떤 토픽도 변경하지 않고 할 수 있어요.
  • 준비가 되면 토픽 자체를 스케일러블 토픽으로 마이그레이션해요 — 이는 V5 애플리케이션에는 투명한 별도의 서버 쪽 단계예요.

사전 요구사항 (Prerequisites)

  • Java 17. V5 클라이언트는 Java 17이 필요해요(현재 SDK는 Java 8+를 지원).
  • 의존성. pulsar-client-v5를 추가하고, v4 코드가 남아 있는 동안에는 pulsar-client-original도 추가해요 — 아래 Dependencies 참고.

의존성 (Dependencies)

pulsar-client-v5는 이미 unshaded v4 클라이언트(pulsar-client-original)를 번들하고 있어요. 점진적으로 마이그레이션하는 동안에는 v4 API를 쓰는 코드를 위해 shade된 pulsar-client가 아니라 pulsar-client-original에 의존하세요 — shade된 pulsar-client는 클라이언트 클래스의 두 번째 충돌 복사본을 추가할 거예요. 전부 V5 API로 옮기면 pulsar-client-v5 하나만으로 충분해요.

Pulsar BOM을 사용해 모든 Pulsar 아티팩트를 한 버전으로 유지하세요.

Maven

<!-- in your <properties> block -->
<pulsar.version>5.0.0-M2</pulsar.version>

<!-- in your <dependencyManagement> block -->
<dependency>
  <groupId>org.apache.pulsar</groupId>
  <artifactId>pulsar-bom</artifactId>
  <version>${pulsar.version}</version>
  <type>pom</type>
  <scope>import</scope>
</dependency>

<!-- in your <dependencies> block; version comes from the BOM -->
<dependency>
  <groupId>org.apache.pulsar</groupId>
  <artifactId>pulsar-client-v5</artifactId>
</dependency>
<!-- only while v4 code remains; remove once fully migrated -->
<dependency>
  <groupId>org.apache.pulsar</groupId>
  <artifactId>pulsar-client-original</artifactId>
</dependency>

Gradle

def pulsarVersion = '5.0.0-M2'

dependencies {
    implementation enforcedPlatform("org.apache.pulsar:pulsar-bom:${pulsarVersion}")

    implementation 'org.apache.pulsar:pulsar-client-v5'
    // only while v4 code remains; remove once fully migrated
    implementation 'org.apache.pulsar:pulsar-client-original'
}

API 매핑 (API mapping)

가장 큰 변화는 컨슈머 모델이에요. 네 가지 서브스크립션 타입이 세 가지 목적별 컨슈머 타입으로 합쳐져요.

현재 SDK V5 SDK
org.apache.pulsar.client.api.PulsarClient org.apache.pulsar.client.api.v5.PulsarClient
Exclusive / Failover 서브스크립션 Stream consumer — 순서대로 소비, 누적 ack
Shared / Key_Shared 서브스크립션 Queue consumer — 개별 ack, 부정 ack, 데드 레터
Reader Checkpoint consumer — Checkpoint를 통한 외부 위치
Schema.STRING, Schema.JSON(T.class), Schema.AVRO(T.class) Schema.string(), Schema.json(T.class), Schema.avro(T.class)
consumer.acknowledge(msg) consumer.acknowledge(msg.id())
reader.seek(messageId) 빌드 시 startPosition(...)에 전달되는 Checkpoint
타임아웃과 지연이 long 밀리초 Duration / Instant
빌더 옵션 세터 불변 구성 레코드 (DeadLetterPolicy, BackoffPolicy, …)

V5 클라이언트가 아직 다루지 않는 것은 계속 현재 SDK를 사용하세요: Reader 스타일의 임의 탐색, TableView, 비영속 토픽. 다른 언어 SDK의 스케일러블 토픽 지원은 계획 중이에요. 오늘날 V5 클라이언트는 Java 전용이에요.

클라이언트 (Client)

클라이언트 빌더는 거의 동일해요. 패키지만 ...client.api에서 ...client.api.v5로 바뀌어요.

// Current
import org.apache.pulsar.client.api.PulsarClient;
// V5
import org.apache.pulsar.client.api.v5.PulsarClient;

PulsarClient client = PulsarClient.builder()
        .serviceUrl("pulsar://localhost:6650")
        .build();

프로듀서 (Producers)

프로듀서와 메시지 빌더는 거의 변경 없이 이어져요. 소문자 스키마 팩토리에 주목하고, 같은 코드가 기존 persistent:// 토픽이나 topic:// 스케일러블 토픽에서 동작한다는 점을 확인하세요.

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

// V5
Producer<String> producer = client.newProducer(Schema.string())
        .topic("persistent://public/default/orders")   // or topic://... for a scalable topic
        .create();

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

컨슈머 (Consumers)

현재 서브스크립션 타입과 일치하는 V5 컨슈머를 고르세요.

Exclusive 또는 Failover → Stream consumer

순서대로 소비하며 누적 ack를 사용해요.

// Current
Consumer<String> consumer = client.newConsumer(Schema.STRING)
        .topic("persistent://public/default/orders")
        .subscriptionName("my-sub")
        .subscriptionType(SubscriptionType.Failover)
        .subscribe();
Message<String> msg = consumer.receive();
consumer.acknowledgeCumulative(msg);

// V5
StreamConsumer<String> consumer = client.newStreamConsumer(Schema.string())
        .topic("persistent://public/default/orders")
        .subscriptionName("my-sub")
        .subscribe();
Message<String> msg = consumer.receive();
consumer.acknowledgeCumulative(msg.id());

Shared 또는 Key_Shared → Queue consumer

개별 ack, 부정 ack, 데드 레터 지원과 함께 병렬 소비를 제공해요.

// Current
Consumer<String> consumer = client.newConsumer(Schema.STRING)
        .topic("persistent://public/default/orders")
        .subscriptionName("workers")
        .subscriptionType(SubscriptionType.Shared)
        .subscribe();
Message<String> msg = consumer.receive();
consumer.acknowledge(msg);            // or consumer.negativeAcknowledge(msg);

// V5
QueueConsumer<String> consumer = client.newQueueConsumer(Schema.string())
        .topic("persistent://public/default/orders")
        .subscriptionName("workers")
        .subscribe();
Message<String> msg = consumer.receive();
consumer.acknowledge(msg.id());       // or consumer.negativeAcknowledge(msg.id());

Reader → Checkpoint consumer

자신의 위치를 추적하는 코드(MessageId에서 시작한 Reader)에는 직렬화 가능한 Checkpoint를 가진 체크포인트 컨슈머를 사용해요.

// Current
Reader<String> reader = client.newReader(Schema.STRING)
        .topic("persistent://public/default/orders")
        .startMessageId(MessageId.earliest)
        .create();
Message<String> msg = reader.readNext();

// V5
CheckpointConsumer<String> consumer = client.newCheckpointConsumer(Schema.string())
        .topic("persistent://public/default/orders")
        .startPosition(Checkpoint.earliest())
        .create();
Message<String> msg = consumer.receive();
byte[] state = consumer.checkpoint().toByteArray();   // persist; restore via Checkpoint.fromByteArray(...)

Checkpoint는 리더와 함께 저장하던 MessageId를 대체해요. 둘은 호환되지 않아요 — 저장된 MessageIdCheckpoint로 사용할 수 없죠. 그래서 리더를 마이그레이션할 때는 깔끔한 전환(cutover)을 계획해야 해요.

여러 토픽: 패턴 서브스크립션 → 네임스페이스 컨슈머

현재 SDK에서 단일 컨슈머는 토픽 목록이나 정규식 패턴으로 많은 토픽에 연결할 수 있어요.

// Current -- pattern subscription
Consumer<String> consumer = client.newConsumer(Schema.STRING)
        .topicsPattern(Pattern.compile("persistent://tenant/ns/orders-.*"))
        .subscriptionName("workers")
        .subscriptionType(SubscriptionType.Shared)
        .subscribe();

V5 SDK에서는 스트림 또는 큐 컨슈머가 이름 패턴이 아니라 토픽 속성으로 좁힐 수 있는 전체 네임스페이스에 연결돼요. topic(...) 대신 namespace(...)를 설정해요.

// V5 -- every scalable topic in the namespace
QueueConsumer<String> consumer = client.newQueueConsumer(Schema.string())
        .namespace("tenant/ns")
        .subscriptionName("workers")
        .subscribe();

// V5 -- only topics whose properties match every filter (AND semantics)
Map<String, String> filters = Map.ofEntries(
        Map.entry("team", "orders"),
        Map.entry("tier", "gold"));

QueueConsumer<String> consumer = client.newQueueConsumer(Schema.string())
        .namespace("tenant/ns", filters)
        .subscriptionName("workers")
        .subscribe();

일치하는 집합은 실시간(live)이에요. 네임스페이스에서 토픽이 만들어지거나 삭제되거나, 속성이 변경되면 컨슈머가 자동으로 연결·해제돼요. v4 패턴 서브스크립션이 새로 생성된 토픽을 추적하는 방식과 같죠. topic(...) 또는 namespace(...) 중 하나만 설정하고 둘 다는 하지 마세요.

필터링은 이름 정규식이 아니라 토픽 속성에 의한 것이므로, 토픽을 만들 때 속성으로 태그(pulsar-admin scalable-topics create ... --property team=orders)하고 일치하는 필터로 선택하세요. 네임스페이스 소비는 스트림과 큐 컨슈머에서 사용 가능해요. 체크포인트 컨슈머는 단일 토픽이에요.

스키마, 트랜잭션, 구성 (Schemas, transactions, and configuration)

  • 스키마Schema.STRING / Schema.JSON(...) 상수를 소문자 팩토리 메서드 Schema.string() / Schema.json(...)로 바꿔요. Schemas 참고.
  • 트랜잭션 — 모델은 변경 없음. .transaction(txn)으로 생산을, 인자 두 개짜리 acknowledge로 ack를 바인딩해요. Transactions 참고.
  • 구성 — 옵션 세터가 불변 레코드(DeadLetterPolicy, BackoffPolicy, BatchingPolicy, …)가 되고, 시간 값은 long 밀리초 대신 Duration / Instant를 사용해요.

토픽 마이그레이션 (Migrating the topics)

V5 API를 도입해도 토픽을 변경할 필요는 없어요 — V5 클라이언트는 기존 persistent:// 토픽에서 동작해요. 스케일러블 토픽의 혜택(자동 분할/병합, 고정 파티션 수 없음)을 얻으려면 서버에서 토픽을 마이그레이션해요.

bin/pulsar-admin scalable-topics migrate persistent://public/default/orders

이것은 단방향 연산이며 연결된 V5 클라이언트에는 투명해요. Migrate a regular topic 참고.

다음으로? (What's next)

더 알아보기 (Learn more)