토픽 관리

토픽 관리 (Manage topics)

Pulsar에는 영속(persistent) 토픽과 비-영속(non-persistent) 토픽이 있어요. 영속 토픽은 메시지를 게시하고 소비하기 위한 논리적 엔드포인트예요. 이 페이지에서는 토픽 리소스 관리부터 비-파티션 토픽, 파티션 토픽, 구독 관리까지 pulsar-admin CLI, REST API, Java admin API로 정리했어요.

출처: 문서

본문

tip

이 페이지는 자주 사용하는 일부 작업만 보여줘요.

  • Pulsar admin 명령, 플래그, 설명 등 최신·전체 정보는 Pulsar admin docs를 참고해요.
  • REST API 파라미터, 응답, 샘플 등 최신·전체 정보는 REST API doc을 참고해요.
  • Java admin API 클래스, 메서드, 설명 등 최신·전체 정보는 Java admin API doc을 참고해요.

Pulsar에는 영속과 비-영속 토픽이 있어요. 영속 토픽은 메시지를 게시하고 소비하기 위한 논리적 엔드포인트예요. 영속 토픽의 이름 구조는 다음과 같아요.

persistent://tenant/namespace/topic

비-영속 토픽은 실시간으로 게시되는 메시지만 소비하고 영속 보장이 필요 없는 애플리케이션에서 사용돼요. 이렇게 하면 메시지 영속화의 오버헤드를 제거해 메시지 게시 지연을 줄여줘요. 비-영속 토픽의 이름 구조는 다음과 같아요.

non-persistent://tenant/namespace/topic

note

토픽 이름: 하위 호환성 때문에 "/," 같은 일부 특수 문자는 토픽 이름의 일부로 허용돼요. 하지만 토픽 이름의 일부로 특수 문자를 사용하지 않는 것을 권장해요.

토픽 리소스 관리 (Manage topic resources)

영속 토픽이든 비-영속 토픽이든 pulsar-admin 도구, REST API, Java로 토픽 리소스를 얻을 수 있어요.

note

REST API에서 :schema는 영속(persistent) 또는 비-영속(non-persistent)을 나타내요. :tenant, :namespace, :x는 변수이며, 사용할 때 실제 tenant, namespace, x 이름으로 바꿔주세요.

예를 들어 GET /admin/v2/persistent/{tenant}/{namespace}을 보면, REST API에서 영속 토픽 목록을 얻으려면 https://pulsar.apache.org/admin/v2/persistent/my-tenant/my-namespace를 사용해요. 비-영속 토픽 목록을 얻으려면 https://pulsar.apache.org/admin/v2/non-persistent/my-tenant/my-namespace를 사용해요.

토픽 목록 (List of topics)

주어진 네임스페이스 아래의 토픽 목록을 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics list my-tenant/my-namespace

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}

Java

String namespace = "my-tenant/my-namespace";
admin.topics().getList(namespace);

권한 부여 (Grant permission)

클라이언트 역할에 주어진 토픽에서 특정 동작을 수행할 권한을 다음 방법으로 부여할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics grant-permission \
    --actions produce,consume \
    --role application1 \
    persistent://test-tenant/ns1/tp1

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/permissions/{role}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String role = "test-role";
Set<AuthAction> actions  = Sets.newHashSet(AuthAction.produce, AuthAction.consume);
admin.topics().grantPermission(topic, role, actions);

권한 가져오기 (Get permission)

권한을 다음 방법으로 가져올 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics permissions persistent://test-tenant/ns1/tp1

예시 출력:

application1    [consume, produce]

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/permissions

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getPermissions(topic);

권한 철회 (Revoke permission)

클라이언트 역할에 부여된 권한을 다음 방법으로 철회할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics revoke-permission \
    --role application1 \
    persistent://test-tenant/ns1/tp1

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}/permissions/{role}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String role = "test-role";
admin.topics().revokePermissions(topic, role);

토픽 삭제 (Delete topic)

토픽을 다음 방법으로 삭제할 수 있어요. 활성 구독이나 프로듀서가 토픽에 연결되어 있으면 삭제할 수 없어요.

pulsar-admin · REST API · Java

pulsar-admin topics delete persistent://test-tenant/ns1/tp1

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().delete(topic);

토픽 언로드 (Unload topic)

토픽을 다음 방법으로 언로드할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics unload persistent://test-tenant/ns1/tp1

REST API: PUT /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/unload

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().unload(topic);

토픽 절단 (Truncate topic)

토픽을 다음 방법으로 절단(truncate)할 수 있어요. truncate 연산은 모든 커서를 토픽의 끝으로 이동하고 비활성 레저(ledger)를 모두 삭제해서, 이미 소비된 메시지가 사용하는 저장 공간을 해제해요.

pulsar-admin · REST API · Java

pulsar-admin topics truncate persistent://test-tenant/ns1/tp1

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}/truncate

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().truncate(topic);

통계 가져오기 (Get stats)

토픽의 상세 통계는 Pulsar statistics를 참고해요.

다음은 토픽 상태의 예시예요.

{
  "msgRateIn" : 0.0,
  "msgThroughputIn" : 0.0,
  "msgRateOut" : 0.0,
  "msgThroughputOut" : 0.0,
  "bytesInCounter" : 504,
  "msgInCounter" : 9,
  "bytesOutCounter" : 2296,
  "msgOutCounter" : 41,
  "averageMsgSize" : 0.0,
  "msgChunkPublished" : false,
  "storageSize" : 504,
  "backlogSize" : 0,
  "filteredEntriesCount" : 100,
  "earliestMsgPublishTimeInBacklogs": 0,
  "offloadedStorageSize" : 0,
  "publishers" : [ {
    "accessMode" : "Shared",
    "msgRateIn" : 0.0,
    "msgThroughputIn" : 0.0,
    "averageMsgSize" : 0.0,
    "chunkedMessageRate" : 0.0,
    "producerId" : 0,
    "metadata" : { },
    "address" : "/127.0.0.1:65402",
    "connectedSince" : "2021-06-09T17:22:55.913+08:00",
    "clientVersion" : "2.9.0-SNAPSHOT",
    "producerName" : "standalone-1-0"
  } ],
  "waitingPublishers" : 0,
  "subscriptions" : {
    "sub-demo" : {
      "msgRateOut" : 0.0,
      "msgThroughputOut" : 0.0,
      "bytesOutCounter" : 2296,
      "msgOutCounter" : 41,
      "msgRateRedeliver" : 0.0,
      "chunkedMessageRate" : 0,
      "msgBacklog" : 0,
      "backlogSize" : 0,
      "earliestMsgPublishTimeInBacklog": 0,
      "msgBacklogNoDelayed" : 0,
      "blockedSubscriptionOnUnackedMsgs" : false,
      "msgDelayed" : 0,
      "unackedMessages" : 0,
      "type" : "Exclusive",
      "activeConsumerName" : "20b81",
      "msgRateExpired" : 0.0,
      "totalMsgExpired" : 0,
      "lastExpireTimestamp" : 0,
      "lastConsumedFlowTimestamp" : 1623230565356,
      "lastConsumedTimestamp" : 1623230583946,
      "lastAckedTimestamp" : 1623230584033,
      "lastMarkDeleteAdvancedTimestamp" : 1623230584033,
      "filterProcessedMsgCount": 100,
      "filterAcceptedMsgCount": 100,
      "filterRejectedMsgCount": 0,
      "filterRescheduledMsgCount": 0,
      "consumers" : [ {
        "msgRateOut" : 0.0,
        "msgThroughputOut" : 0.0,
        "bytesOutCounter" : 2296,
        "msgOutCounter" : 41,
        "msgRateRedeliver" : 0.0,
        "chunkedMessageRate" : 0.0,
        "consumerName" : "20b81",
        "availablePermits" : 959,
        "unackedMessages" : 0,
        "avgMessagesPerEntry" : 314,
        "blockedConsumerOnUnackedMsgs" : false,
        "lastAckedTimestamp" : 1623230584033,
        "lastConsumedTimestamp" : 1623230583946,
        "metadata" : { },
        "address" : "/127.0.0.1:65172",
        "connectedSince" : "2021-06-09T17:22:45.353+08:00",
        "clientVersion" : "2.9.0-SNAPSHOT"
      } ],
      "allowOutOfOrderDelivery": false,
      "consumersAfterMarkDeletePosition" : { },
      "nonContiguousDeletedMessagesRanges" : 0,
      "nonContiguousDeletedMessagesRangesSerializedSize" : 0,
      "durable" : true,
      "replicated" : false
    }
  },
  "replication" : { },
  "deduplicationStatus" : "Disabled",
  "nonContiguousDeletedMessagesRanges" : 0,
  "nonContiguousDeletedMessagesRangesSerializedSize" : 0,
  "ownerBroker" : "localhost:8080"
}

토픽의 상태를 얻으려면 다음 방법을 사용할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics stats persistent://test-tenant/ns1/tp1

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/stats

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getStats(topic);

내부 통계 가져오기 (Get internal stats)

토픽 내부의 상세 통계는 Pulsar statistics를 참고해요.

다음은 토픽의 내부 통계 예시예요.

{
    "entriesAddedCounter":0,
    "numberOfEntries":0,
    "totalSize":0,
    "currentLedgerEntries":0,
    "currentLedgerSize":0,
    "lastLedgerCreatedTimestamp":"2021-01-22T21:12:14.868+08:00",
    "lastLedgerCreationFailureTimestamp":null,
    "waitingCursorsCount":0,
    "pendingAddEntriesCount":0,
    "lastConfirmedEntry":"3:-1",
    "state":"LedgerOpened",
    "ledgers":[
        {
            "ledgerId":3,
            "entries":0,
            "size":0,
            "offloaded":false,
            "metadata":null
        }
    ],
    "cursors":{
        "test":{
            "markDeletePosition":"3:-1",
            "readPosition":"3:-1",
            "waitingReadOp":false,
            "pendingReadOps":0,
            "messagesConsumedCounter":0,
            "cursorLedger":4,
            "cursorLedgerLastEntry":1,
            "individuallyDeletedMessages":"[]",
            "lastLedgerSwitchTimestamp":"2021-01-22T21:12:14.966+08:00",
            "state":"Open",
            "numberOfEntriesSinceFirstNotAckedMessage":0,
            "totalNonContiguousDeletedMessagesRange":0,
            "properties":{

            }
        }
    },
    "schemaLedgers":[
        {
            "ledgerId":1,
            "entries":11,
            "size":10,
            "offloaded":false,
            "metadata":null
        }
    ],
    "compactedLedger":{
        "ledgerId":-1,
        "entries":-1,
        "size":-1,
        "offloaded":false,
        "metadata":null
    }
}

토픽의 내부 상태를 얻으려면 다음 방법을 사용할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics stats-internal persistent://test-tenant/ns1/tp1

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/internalStats

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getInternalStats(topic);

메시지 픽 (Peek messages)

주어진 토픽의 특정 구독에 대해 일정 수의 메시지를 픽(peek)할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics peek-messages \
    --count 10 --subscription my-subscription \
    persistent://test-tenant/ns1/tp1

예시 출력:

Message ID: 77:2
Publish time: 1668674963028
Event time: 0
         +-------------------------------------------------+
         |  0  1  2  3  4  5  6  7  8  9  a  b  c  d  e  f |
+--------+-------------------------------------------------+----------------+
|00000000| 68 65 6c 6c 6f 2d 31                            |hello-1         |
+--------+-------------------------------------------------+----------------+

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}/position/{messagePosition}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String subName = "my-subscription";
int numMessages = 1;
admin.topics().peekMessages(topic, subName, numMessages);

ID로 메시지 가져오기 (Get message by ID)

주어진 레저 ID와 엔트리 ID로 메시지를 다음 방법으로 가져올 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics get-message-by-id \
    -l 10 -e 0 persistent://public/default/my-topic

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/ledger/{ledgerId}/entry/{entryId}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
long ledgerId = 10;
long entryId = 10;
admin.topics().getMessageById(topic, ledgerId, entryId);

메시지 검사 (Examine messages)

가장 오래된 또는 가장 최신 메시지에 대한 상대 위치로 토픽의 특정 메시지를 검사할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics examine-messages \
    -i latest -m 1 persistent://public/default/my-topic

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/examinemessage

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().examineMessage(topic, "latest", 1);

메시지 ID 가져오기 (Get message ID)

주어진 날짜시간에 게시되거나 그 직후의 메시지 ID를 가져올 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics get-message-id \
    persistent://public/default/my-topic \
    -d 2021-06-28T19:01:17Z

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/messageid/{timestamp}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
long timestamp = System.currentTimeMillis()
admin.topics().getMessageIdByTimestamp(topic, timestamp);

메시지 건너뛰기 (Skip messages)

주어진 토픽의 특정 구독에 대해 일정 수의 메시지를 건너뛸 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics skip \
    --count 10 --subscription my-subscription \
    persistent://test-tenant/ns1/tp1

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}/skip/{numMessages}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String subName = "my-subscription";
int numMessages = 1;
admin.topics().skipMessages(topic, subName, numMessages);

모든 메시지 건너뛰기 (Skip all messages)

주어진 토픽의 특정 구독에 대해 모든 오래된 메시지를 건너뛸 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics clear-backlog \
    --subscription my-subscription \
    persistent://test-tenant/ns1/tp1

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}/skip_all

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String subName = "my-subscription";
admin.topics().skipAllMessages(topic, subName);

커서 리셋 (Reset cursor)

구독 커서 위치를 X초(또는 100m, 3h, 2d, 5w 같은 다른 시간 단위) 이전에 기록된 위치로 리셋할 수 있어요. 본질적으로 X초 전의 커서 시간과 위치를 계산해 그 위치로 리셋하는 거예요. 커서를 다음 방법으로 리셋할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics reset-cursor \
    --subscription my-subscription --time 10 \
    persistent://test-tenant/ns1/tp1

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}/resetcursor/{timestamp}

Java


String topic = "persistent://my-tenant/my-namespace/my-topic";
String subName = "my-subscription";
long timestamp = 2342343L;
admin.topics().resetCursor(topic, subName, timestamp);

토픽의 소유 브로커 룩업 (Look up topic's owner broker)

주어진 토픽의 소유 브로커를 다음 방법으로 찾을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics lookup persistent://test-tenant/ns1/tp1

예시 출력:

"pulsar://broker1.org.com:4480"

REST API: GET /lookup/v2/topic/{topic-domain}/{tenant}/{namespace}/{topic}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.lookups().lookupDestination(topic);

파티션 토픽의 소유 브로커 룩업 (Look up partitioned topic's owner broker)

주어진 파티션 토픽의 소유 브로커를 다음 방법으로 찾을 수 있어요.

pulsar-admin · Java

pulsar-admin topics partitioned-lookup persistent://test-tenant/ns1/my-topic

예시 출력:

"persistent://test-tenant/ns1/my-topic-partition-0   pulsar://localhost:6650"
"persistent://test-tenant/ns1/my-topic-partition-1   pulsar://localhost:6650"
"persistent://test-tenant/ns1/my-topic-partition-2   pulsar://localhost:6650"
"persistent://test-tenant/ns1/my-topic-partition-3   pulsar://localhost:6650"

파티션 토픽을 브로커 URL로 정렬해 룩업해요.

pulsar-admin topics partitioned-lookup \
    persistent://test-tenant/ns1/my-topic --sort-by-broker

예시 출력:

pulsar://localhost:6650   [persistent://test-tenant/ns1/my-topic-partition-0, persistent://test-tenant/ns1/my-topic-partition-1, persistent://test-tenant/ns1/my-topic-partition-2, persistent://test-tenant/ns1/my-topic-partition-3]

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.lookups().lookupPartitionedTopic(topic);

번들 가져오기 (Get bundle)

주어진 토픽이 속하는 번들의 범위를 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics bundle-range persistent://test-tenant/ns1/tp1

예시 출력:

"0x00000000_0xffffffff"

REST API: GET /lookup/v2/topic/{topic-domain}/{tenant}/{namespace}/{topic}/bundle

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.lookups().getBundleRange(topic);

구독 가져오기 (Get subscriptions)

주어진 토픽의 모든 구독 이름을 다음 방법으로 확인할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics subscriptions persistent://test-tenant/ns1/tp1

예시 출력:

my-subscription

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscriptions

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getSubscriptions(topic);

마지막 메시지 ID (Last Message Id)

영속 토픽의 마지막 커밋된 메시지 ID를 얻을 수 있어요. 2.3.0 릴리스부터 사용 가능해요.

pulsar-admin · REST API · Java

pulsar-admin topics last-message-id topic-name

예시 출력:

{
  "ledgerId" : 97,
  "entryId" : 9,
  "partitionIndex" : -1
}

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/lastMessageId

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getLastMessage(topic);

백로그 크기 가져오기 (Get backlog size)

주어진 메시지 ID(바이트 단위)에 대해 단일 파티션 토픽 또는 비-파티션 토픽의 백로그 크기를 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics get-backlog-size \
  -m 1:1 \
  persistent://test-tenant/ns1/tp1-partition-0

REST API: PUT /admin/v2/persistent/{tenant}/{namespace}/{topic}/backlogSize

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
MessageId messageId = MessageId.earliest;
admin.topics().getBacklogSizeByMessageId(topic, messageId);

중복 제거 스냅샷 간격 구성 (Configure deduplication snapshot interval)

중복 제거 스냅샷 간격 가져오기 (Get deduplication snapshot interval)

토픽 수준의 중복 제거 스냅샷 간격을 얻으려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics get-deduplication-snapshot-interval my-topic

REST API: GET /admin/v2/namespaces/{tenant}/{namespace}/deduplicationSnapshotInterval

Java

admin.topics().getDeduplicationSnapshotInterval(topic)
중복 제거 스냅샷 간격 설정 (Set deduplication snapshot interval)

토픽 수준의 중복 제거 스냅샷 간격을 설정하려면 다음 방법 중 하나를 사용해요.

전제 조건: brokerDeduplicationEnabledtrue로 설정되어야 해요.

pulsar-admin · REST API · Java

pulsar-admin topics set-deduplication-snapshot-interval my-topic -i 1000

REST API: POST /admin/v2/namespaces/{tenant}/{namespace}/deduplicationSnapshotInterval

{
  "interval": 1000
}

Java

admin.topics().setDeduplicationSnapshotInterval(topic, 1000)
중복 제거 스냅샷 간격 제거 (Remove deduplication snapshot interval)

토픽 수준의 중복 제거 스냅샷 간격을 제거하려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics remove-deduplication-snapshot-interval my-topic

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}/deduplicationSnapshotInterval

Java

admin.topics().removeDeduplicationSnapshotInterval(topic)

비활성 토픽 정책 구성 (Configure inactive topic policies)

비활성 토픽 정책 가져오기 (Get inactive topic policies)

토픽 수준의 비활성 토픽 정책을 얻으려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics get-inactive-topic-policies my-topic

REST API: GET /admin/v2/namespaces/{tenant}/{namespace}/inactiveTopicPolicies

Java

admin.topics().getInactiveTopicPolicies(topic)
비활성 토픽 정책 설정 (Set inactive topic policies)

토픽 수준의 비활성 토픽 정책을 설정하려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics set-inactive-topic-policies my-topic

REST API: POST /admin/v2/namespaces/{tenant}/{namespace}/inactiveTopicPolicies

Java

admin.topics().setInactiveTopicPolicies(topic, inactiveTopicPolicies)
비활성 토픽 정책 제거 (Remove inactive topic policies)

토픽 수준의 비활성 토픽 정책을 제거하려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics remove-inactive-topic-policies my-topic

REST API: DELETE /admin/v2/namespaces/{tenant}/{namespace}/inactiveTopicPolicies

Java

admin.topics().removeInactiveTopicPolicies(topic)

오프로드 정책 구성 (Configure offload policies)

오프로드 정책 가져오기 (Get offload policies)

토픽 수준의 오프로드 정책을 얻으려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics get-offload-policies my-topic

REST API: GET /admin/v2/namespaces/{tenant}/{namespace}/offloadPolicies

Java

admin.topics().getOffloadPolicies(topic)
오프로드 정책 설정 (Set offload policies)

토픽 수준의 오프로드 정책을 설정하려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics set-offload-policies my-topic

REST API: POST /admin/v2/namespaces/{tenant}/{namespace}/offloadPolicies

Java

admin.topics().setOffloadPolicies(topic, offloadPolicies)
오프로드 정책 제거 (Remove offload policies)

토픽 수준의 오프로드 정책을 제거하려면 다음 방법 중 하나를 사용해요.

pulsar-admin · REST API · Java

pulsar-admin topics remove-offload-policies my-topic

REST API: DELETE /admin/v2/namespaces/{tenant}/{namespace}/removeOffloadPolicies

Java

admin.topics().removeOffloadPolicies(topic)

비-파티션 토픽 관리 (Manage non-partitioned topics)

Pulsar admin API로 비-파티션 토픽을 생성·삭제하고 상태를 확인할 수 있어요.

생성 (Create)

비-파티션 토픽은 명시적으로 만들어야 해요. 새 비-파티션 토픽을 만들 때 토픽 이름을 제공해야 해요.

기본적으로 생성 후 60초 동안 메시지가 없으면 토픽은 비활성으로 간주되어 쓰레기 데이터가 생기는 것을 막기 위해 자동으로 삭제돼요. 이 기능을 비활성화하려면 brokerDeleteInactiveTopicsEnabledfalse로 설정해요. 비활성 토픽 확인 빈도를 바꾸려면 brokerDeleteInactiveTopicsFrequencySeconds를 특정 값으로 설정해요.

두 파라미터에 대한 자세한 내용은 여기를 참고해요.

비-파티션 토픽을 다음 방법으로 만들 수 있어요.

pulsar-admin · REST API · Java

create 명령으로 비-파티션 토픽을 만들 때 인자로 토픽 이름을 지정해야 해요.

pulsar-admin topics create \
    persistent://my-tenant/my-namespace/my-topic

note

토픽 이름에 '-partition-' 접미사 뒤에 숫자 값('xyz-topic-partition-x'처럼)을 붙여 비-파티션 토픽을 만들 때, 같은 접미사('xyz-topic-partition-y')를 가진 파티션 토픽이 존재한다면 비-파티션 토픽의 숫자 값(x)이 파티션 토픽의 파티션 수(y)보다 커야 해요. 그렇지 않으면 그런 비-파티션 토픽을 만들 수 없어요.

REST API: PUT /admin/v2/persistent/{tenant}/{namespace}/{topic}

Java

String topicName = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().createNonPartitionedTopic(topicName);

삭제 (Delete)

비-파티션 토픽을 다음 방법으로 삭제할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics delete \
    persistent://my-tenant/my-namespace/my-topic

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}

Java

admin.topics().delete(topic);

나열 (List)

주어진 네임스페이스 아래의 토픽 목록을 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics list tenant/namespace

예시 출력:

persistent://tenant/namespace/topic1
persistent://tenant/namespace/topic2

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}

Java

admin.topics().getList(namespace);

통계 (Stats)

주어진 토픽과 연결된 프로듀서·컨슈머의 현재 통계를 다음 방법으로 확인할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics stats \
    persistent://test-tenant/namespace/topic \
    --get-precise-backlog

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/stats

Java

admin.topics().getStats(topic, false /* is precise backlog */);

다음은 예시예요. 토픽 통계 설명은 Pulsar statistics를 참고해요.

{
  "msgRateIn": 4641.528542257553,
  "msgThroughputIn": 44663039.74947473,
  "msgRateOut": 0,
  "msgThroughputOut": 0,
  "averageMsgSize": 1232439.816728665,
  "storageSize": 135532389160,
  "publishers": [
    {
      "msgRateIn": 57.855383881403576,
      "msgThroughputIn": 558994.7078932219,
      "averageMsgSize": 613135,
      "producerId": 0,
      "producerName": null,
      "address": null,
      "connectedSince": null
    }
  ],
  "subscriptions": {
    "my-topic_subscription": {
      "msgRateOut": 0,
      "msgThroughputOut": 0,
      "msgBacklog": 116632,
      "type": null,
      "msgRateExpired": 36.98245516804671,
      "consumers": []
    }
  },
  "replication": {}
}

내부 통계 (Internal stats)

토픽의 상세 통계를 확인할 수 있어요. 다음은 예시예요. 각 내부 토픽 통계 설명은 Pulsar statistics를 참고해요.

{
  "entriesAddedCounter": 20449518,
  "numberOfEntries": 3233,
  "totalSize": 331482,
  "currentLedgerEntries": 3233,
  "currentLedgerSize": 331482,
  "lastLedgerCreatedTimestamp": "2016-06-29 03:00:23.825",
  "lastLedgerCreationFailureTimestamp": null,
  "waitingCursorsCount": 1,
  "pendingAddEntriesCount": 0,
  "lastConfirmedEntry": "324711539:3232",
  "state": "LedgerOpened",
  "ledgers": [
    {
      "ledgerId": 324711539,
      "entries": 0,
      "size": 0
    }
  ],
  "cursors": {
    "my-subscription": {
      "markDeletePosition": "324711539:3133",
      "readPosition": "324711539:3233",
      "waitingReadOp": true,
      "pendingReadOps": 0,
      "messagesConsumedCounter": 20449501,
      "cursorLedger": 324702104,
      "cursorLedgerLastEntry": 21,
      "individuallyDeletedMessages": "[(324711539:3134‥324711539:3136], (324711539:3137‥324711539:3140], ]",
      "lastLedgerSwitchTimestamp": "2016-06-29 01:30:19.313",
      "state": "Open"
    }
  }
}

파티션 토픽의 내부 통계를 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics stats-internal \
    persistent://test-tenant/namespace/topic

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/internalStats

Java

admin.topics().getInternalStats(topic);

파티션 토픽 관리 (Manage partitioned topics)

Pulsar admin API로 파티션 토픽을 생성·업데이트·삭제하고 상태를 확인할 수 있어요.

생성 (Create)

새 파티션 토픽을 만들 때 토픽 이름과 파티션 수를 제공해야 해요.

note

기본적으로 생성 후 60초 동안 메시지가 없으면 토픽은 비활성으로 간주되어 쓰레기 데이터가 생기는 것을 막기 위해 자동으로 삭제돼요. 이 기능을 비활성화하려면 brokerDeleteInactiveTopicsEnabledfalse로 설정해요. 비활성 토픽 확인 빈도를 바꾸려면 brokerDeleteInactiveTopicsFrequencySeconds를 특정 값으로 설정해요.

두 파라미터에 대한 자세한 내용은 여기를 참고해요.

파티션 토픽을 다음 방법으로 만들 수 있어요.

pulsar-admin · REST API · Java

create-partitioned-topic 명령으로 파티션 토픽을 만들 때 인자로 토픽 이름을, -p 또는 --partitions 플래그로 파티션 수를 지정해야 해요.

pulsar-admin topics create-partitioned-topic \
    persistent://my-tenant/my-namespace/my-topic \
    --partitions 4

note

'-partition-' 접미사 뒤에 숫자 값을 붙인 비-파티션 토픽('xyz-topic-partition-10'처럼)이 있다면 'xyz-topic'이라는 이름의 파티션 토픽을 만들 수 없어요. 파티션 토픽의 파티션이 기존 비-파티션 토픽을 덮어쓸 수 있기 때문이에요. 그런 파티션 토픽을 만들려면 먼저 그 비-파티션 토픽을 삭제해야 해요.

REST API: PUT /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitions

Java

String topicName = "persistent://my-tenant/my-namespace/my-topic";
int numPartitions = 4;
admin.topics().createPartitionedTopic(topicName, numPartitions);

누락된 파티션 생성 (Create missed partitions)

토픽 자동 생성이 비활성화되어 있고 파티션이 없는 파티션 토픽이 있다면, create-missed-partitions 명령으로 토픽에 파티션을 만들 수 있어요.

pulsar-admin · REST API · Java

create-missed-partitions 명령으로 인자에 토픽 이름을 지정해 누락된 파티션을 만들 수 있어요.

pulsar-admin topics create-missed-partitions \
    persistent://my-tenant/my-namespace/my-topic

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/createMissedPartitions

Java

String topicName = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().createMissedPartitions(topicName);

메타데이터 가져오기 (Get metadata)

파티션 토픽은 메타데이터와 연관되어 있고, JSON 객체로 볼 수 있어요. 다음 메타데이터 필드를 사용할 수 있어요.

Field Description
partitions 토픽이 분할되는 파티션 수

pulsar-admin · REST API · Java

get-partitioned-topic-metadata 하위 명령으로 파티션 토픽의 파티션 수를 확인할 수 있어요.

pulsar-admin topics get-partitioned-topic-metadata \
    persistent://my-tenant/my-namespace/my-topic

예시 출력:

{
  "partitions" : 4,
  "deleted" : false
}

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitions

Java

String topicName = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getPartitionedTopicMetadata(topicName);

업데이트 (Update)

기존 파티션 토픽의 파티션 수를 업데이트할 수 있어요. 하지만 파티션 수는 늘릴 수만 있어요. 파티션 수를 줄이면 토픽을 삭제하는 것이므로 Pulsar에서 지원하지 않아요.

프로듀서와 컨슈머는 새로 생성된 파티션을 자동으로 찾을 수 있어요.

pulsar-admin · REST API · Java

update-partitioned-topic 명령으로 파티션 토픽을 업데이트할 수 있어요.

pulsar-admin topics update-partitioned-topic \
    persistent://my-tenant/my-namespace/my-topic \
    --partitions 8

REST API: POST /admin/v2/persistent/{tenant}/{namespace}/{topic}/partitions

Java

admin.topics().updatePartitionedTopic(topic, numPartitions);

삭제 (Delete)

delete-partitioned-topic 명령, REST API, Java로 파티션 토픽을 삭제할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics delete-partitioned-topic \
    persistent://my-tenant/my-namespace/my-topic

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}/partitions

Java

admin.topics().deletePartitionedTopic(topic);

나열 (List)

주어진 네임스페이스 아래의 파티션 토픽 목록을 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics list-partitioned-topics tenant/namespace

예시 출력:

persistent://tenant/namespace/topic1
persistent://tenant/namespace/topic2

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/partitioned

Java

admin.topics().getPartitionedTopicList(namespace);

통계 (Stats)

주어진 파티션 토픽과 연결된 프로듀서·컨슈머의 현재 통계를 다음 방법으로 확인할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics partitioned-stats \
    persistent://test-tenant/namespace/topic \
    --per-partition

REST API: GET /admin/v2/non-persistent/{tenant}/{namespace}/{topic}/partitioned-stats

Java

admin.topics().getPartitionedStats(topic, true /* per partition */, false /* is precise backlog */);

다음은 예시예요. 각 토픽 통계 설명은 Pulsar statistics를 참고해요.

구독 JSON 객체에서 chuckedMessageRate는 비권장(deprecated)이에요. chunkedMessageRate를 사용해주세요. 지금은 둘 다 JSON으로 보내져요.

{
  "msgRateIn" : 999.992947159793,
  "msgThroughputIn" : 1070918.4635439808,
  "msgRateOut" : 0.0,
  "msgThroughputOut" : 0.0,
  "bytesInCounter" : 270318763,
  "msgInCounter" : 252489,
  "bytesOutCounter" : 0,
  "msgOutCounter" : 0,
  "averageMsgSize" : 1070.926056966454,
  "msgChunkPublished" : false,
  "storageSize" : 270316646,
  "backlogSize" : 200921133,
  "publishers" : [ {
    "msgRateIn" : 999.992947159793,
    "msgThroughputIn" : 1070918.4635439808,
    "averageMsgSize" : 1070.3333333333333,
    "chunkedMessageRate" : 0.0,
    "producerId" : 0
  } ],
  "subscriptions" : {
    "test" : {
      "msgRateOut" : 0.0,
      "msgThroughputOut" : 0.0,
      "bytesOutCounter" : 0,
      "msgOutCounter" : 0,
      "msgRateRedeliver" : 0.0,
      "chuckedMessageRate" : 0,
      "chunkedMessageRate" : 0,
      "msgBacklog" : 144318,
      "msgBacklogNoDelayed" : 144318,
      "blockedSubscriptionOnUnackedMsgs" : false,
      "msgDelayed" : 0,
      "unackedMessages" : 0,
      "msgRateExpired" : 0.0,
      "lastExpireTimestamp" : 0,
      "lastConsumedFlowTimestamp" : 0,
      "lastConsumedTimestamp" : 0,
      "lastAckedTimestamp" : 0,
      "consumers" : [ ],
      "isDurable" : true,
      "isReplicated" : false
    }
  },
  "replication" : { },
  "metadata" : {
    "partitions" : 3
  },
  "partitions" : { }
}

내부 통계 (Internal stats)

파티션 토픽의 상세 통계를 확인할 수 있어요. 다음은 예시예요. 각 내부 토픽 통계 설명은 Pulsar statistics를 참고해요.

{
  "entriesAddedCounter": 20449518,
  "numberOfEntries": 3233,
  "totalSize": 331482,
  "currentLedgerEntries": 3233,
  "currentLedgerSize": 331482,
  "lastLedgerCreatedTimestamp": "2016-06-29 03:00:23.825",
  "lastLedgerCreationFailureTimestamp": null,
  "waitingCursorsCount": 1,
  "pendingAddEntriesCount": 0,
  "lastConfirmedEntry": "324711539:3232",
  "state": "LedgerOpened",
  "ledgers": [
    {
      "ledgerId": 324711539,
      "entries": 0,
      "size": 0
    }
  ],
  "cursors": {
    "my-subscription": {
      "markDeletePosition": "324711539:3133",
      "readPosition": "324711539:3233",
      "waitingReadOp": true,
      "pendingReadOps": 0,
      "messagesConsumedCounter": 20449501,
      "cursorLedger": 324702104,
      "cursorLedgerLastEntry": 21,
      "individuallyDeletedMessages": "[(324711539:3134‥324711539:3136], (324711539:3137‥324711539:3140], ]",
      "lastLedgerSwitchTimestamp": "2016-06-29 01:30:19.313",
      "state": "Open"
    }
  }
}

파티션 토픽의 내부 통계를 다음 방법으로 얻을 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics partitioned-stats-internal \
    persistent://test-tenant/namespace/topic

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/partitioned-internalStats

Java

admin.topics().getPartitionedInternalStats(topic);

구독 관리 (Manage subscriptions)

Pulsar admin API로 구독을 생성·확인·삭제할 수 있어요.

구독 생성 (Create subscription)

토픽에 대한 구독을 다음 방법 중 하나로 만들 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics create-subscription \
    --subscription my-subscription \
    persistent://test-tenant/ns1/tp1

REST API: PUT /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subscriptionName}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String subscriptionName = "my-subscription";
admin.topics().createSubscription(topic, subscriptionName, MessageId.latest);

구독 가져오기 (Get subscription)

주어진 토픽의 모든 구독 이름을 다음 방법 중 하나로 확인할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics subscriptions persistent://test-tenant/ns1/tp1

예시 출력:

my-subscription

REST API: GET /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscriptions

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
admin.topics().getSubscriptions(topic);

구독 해지 (Unsubscribe subscription)

구독이 더 이상 메시지를 처리하지 않으면 다음 방법 중 하나로 구독을 해지(unsubscribe)할 수 있어요.

pulsar-admin · REST API · Java

pulsar-admin topics unsubscribe \
    --subscription my-subscription \
    persistent://test-tenant/ns1/tp1

REST API: DELETE /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}

Java

String topic = "persistent://my-tenant/my-namespace/my-topic";
String subscriptionName = "my-subscription";
admin.topics().deleteSubscription(topic, subscriptionName);

더 알아보기 (Learn more)

  • 토픽 관련 전체 명령은 pulsar-admin의 topics 명령 참고서를 확인해요.
  • REST API 엔드포인트 상세는 /admin/v2/topics/lookup/v2 문서를 참고해요.
  • Java admin API의 topics 메서드는 PulsarAdmin 객체 문서에서 볼 수 있어요.
  • 토픽 통계 필드에 대한 상세 설명은 Pulsar statistics 문서를 참고해요.
  • 토픽·구독·파티션의 개념이 궁금하다면 핵심 개념 문서를 참고해요.