토픽 관리
토픽 관리 (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)
토픽 수준의 중복 제거 스냅샷 간격을 설정하려면 다음 방법 중 하나를 사용해요.
전제 조건: brokerDeduplicationEnabled가 true로 설정되어야 해요.
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초 동안 메시지가 없으면 토픽은 비활성으로 간주되어 쓰레기 데이터가 생기는 것을 막기 위해 자동으로 삭제돼요. 이 기능을 비활성화하려면 brokerDeleteInactiveTopicsEnabled를 false로 설정해요. 비활성 토픽 확인 빈도를 바꾸려면 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초 동안 메시지가 없으면 토픽은 비활성으로 간주되어 쓰레기 데이터가 생기는 것을 막기 위해 자동으로 삭제돼요. 이 기능을 비활성화하려면
brokerDeleteInactiveTopicsEnabled를false로 설정해요. 비활성 토픽 확인 빈도를 바꾸려면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 문서를 참고해요.
- 토픽·구독·파티션의 개념이 궁금하다면 핵심 개념 문서를 참고해요.