트랜잭션 관리
트랜잭션 관리 (Manage transactions)
트랜잭션은 여러 토픽·파티션에 걸친 쓰기와 읽기를 원자적으로 묶어주는 기능이에요. 이 페이지에서는 느린 트랜잭션 가져오기, 트랜잭션 코디네이터 확장, 트랜잭션·코디네이터·pending ack·트랜잭션 버퍼 통계 조회를 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을 참고해요.
트랜잭션 리소스 (Transaction resources)
느린 트랜잭션 가져오기 (GetSlowTransactions)
프로덕션 환경에서는 완료되지 못한 오래 지속되는 트랜잭션이 있을 수 있어요. 특정 코디네이터 또는 모든 코디네이터에서 일정 시간 이상 생존한 느린 트랜잭션을 다음 방법으로 얻을 수 있어요.
pulsar-admin · REST API · Java
pulsar-admin transactions slow-transactions -c 1 -t 1s
REST API: GET /admin/v3/transactions/slowTransactions/{timeout}
Java
admin.transactions().getSlowTransactionsByCoordinatorId(coordinatorId, timeout, timeUnit)
//Or get slow transactions from all coordinators
admin.transactions().getSlowTransactions(timeout, timeUnit)
반환 값의 예시:
{
"(0,3)": {
"txnId": "(0,3)",
"status": "OPEN",
"openTimestamp": 1658120122474,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,2)": {
"txnId": "(0,2)",
"status": "OPEN",
"openTimestamp": 1658120122471,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,5)": {
"txnId": "(0,5)",
"status": "OPEN",
"openTimestamp": 1658120122478,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,4)": {
"txnId": "(0,4)",
"status": "OPEN",
"openTimestamp": 1658120122476,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,7)": {
"txnId": "(0,7)",
"status": "OPEN",
"openTimestamp": 1658120122482,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,10)": {
"txnId": "(0,10)",
"status": "OPEN",
"openTimestamp": 1658120122488,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,6)": {
"txnId": "(0,6)",
"status": "OPEN",
"openTimestamp": 1658120122480,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,9)": {
"txnId": "(0,9)",
"status": "OPEN",
"openTimestamp": 1658120122486,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,8)": {
"txnId": "(0,8)",
"status": "OPEN",
"openTimestamp": 1658120122484,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
},
"(0,11)": {
"txnId": "(0,11)",
"status": "OPEN",
"openTimestamp": 1658120122490,
"timeoutAt": 300000,
"producedPartitions": {},
"ackedPartitions": {}
}
}
트랜잭션 코디네이터 확장 (ScaleTransactionCoordinators)
트랜잭션 코디네이터 수가 부족해서 트랜잭션 성능이 병목에 도달하면, 다음 방법으로 트랜잭션 코디네이터 수를 확장할 수 있어요.
pulsar-admin · REST API · Java
pulsar-admin transactions scale-transactionCoordinators -r 17
REST API: POST /admin/v3/transactions/transactionCoordinator/replicas
Java
admin.transactions().scaleTransactionCoordinators(replicas);
트랜잭션 통계 (Transaction stats)
트랜잭션 메타데이터 가져오기 (Get transaction metadata)
가져올 수 있는 트랜잭션 메타데이터는 다음과 같아요.
txnId: 이 트랜잭션의 IDstatus: 이 트랜잭션의 상태openTimestamp: 이 트랜잭션의 열린 시간timeoutAt: 이 트랜잭션의 타임아웃producedPartitions: 이 트랜잭션으로 메시지가 전송된 파티션 또는 토픽ackedPartitions: 이 트랜잭션으로 메시지가 승인된 파티션 또는 토픽
다음 방법 중 하나로 트랜잭션 메타데이터를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions transaction-metadata -m 1 -l 1
REST API: GET /admin/v3/transactions/transactionMetadata/{mostSigBits}/{leastSigBits}
Java
admin.transactions().getTransactionMetadata(txnID);
반환 값의 예시:
{
"txnId" : "(1,18)",
"status" : "ABORTING",
"openTimestamp" : 1656592983374,
"timeoutAt" : 5000,
"producedPartitions" : {
"my-topic" : {
"startPosition" : "127:4959",
"aborted" : true
}
},
"ackedPartitions" : {
"my-topic" : {
"mysubName" : {
"cumulativeAckPosition" : null
}
}
}
}
트랜잭션 pending ack에서 트랜잭션 통계 가져오기 (Get transaction stats in transaction pending ack)
트랜잭션 pending ack에서 가져올 수 있는 트랜잭션 통계는 다음과 같아요.
cumulativeAckPosition: 이 트랜잭션이 이 구독에서 누적적으로 승인하는 위치
다음 방법 중 하나로 pending ack의 트랜잭션 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions transaction-in-pending-ack-stats -m 1 -l 1 -t my-topic -s mysubname
REST API: GET /admin/v3/transactions/transactionInPendingAckStats/{tenant}/{namespace}/{topic}/{subName}/{mostSigBits}/{leastSigBits}
Java
admin.transactions().getTransactionInPendingAckStats(txnID, topic, subname);
반환 값의 예시:
{
"cumulativeAckPosition" : "137:49959"
}
트랜잭션 버퍼에서 트랜잭션 통계 가져오기 (Get transaction stats in transaction buffer)
트랜잭션 버퍼에서 가져올 수 있는 트랜잭션 통계는 다음과 같아요.
startPosition: 트랜잭션 버퍼에서 이 트랜잭션의 시작 위치aborted: 이 트랜잭션이 중단되었는지 여부의 플래그
다음 방법 중 하나로 트랜잭션 버퍼의 트랜잭션 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions transaction-in-buffer-stats -m 1 -l 1 -t my-topic
REST API: GET /admin/v3/transactions/transactionInBufferStats/{tenant}/{namespace}/{topic}/{mostSigBits}/{leastSigBits}
Java
admin.transactions().getTransactionInBufferStatsAsync(txnID, topic);
반환 값의 예시:
{
"startPosition" : "137:49759",
"aborted" : false
}
트랜잭션 코디네이터 통계 (Transaction coordinator stats)
트랜잭션 코디네이터(TC)는 Pulsar 브로커 안의 모듈이에요. 트랜잭션의 전체 수명 주기를 유지하고 트랜잭션 타임아웃을 처리해요.
코디네이터 통계 가져오기 (Get coordinator stats)
가져올 수 있는 트랜잭션 코디네이터 통계는 다음과 같아요.
state: 이 트랜잭션 코디네이터의 상태leastSigBit:s: 이 트랜잭션 코디네이터의 시퀀스 IDlowWaterMark: 이 트랜잭션 코디네이터의 낮은 워터마크ongoingTxnSize: 이 트랜잭션 코디네이터에서 진행 중인 트랜잭션의 총 수recoverStartTime: 트랜잭션 코디네이터 복구의 시작 타임스탬프. 0L은 시작 없음을 의미recoverEndTime: 트랜잭션 코디네이터 복구의 종료 타임스탬프. 0L은 시작 없음을 의미
다음 방법 중 하나로 트랜잭션 코디네이터 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions coordinator-stats -c 1
REST API: GET /admin/v3/transactions/coordinatorStats
Java
admin.transactions().getCoordinatorStatsById(coordinatorId);
//Or get all coordinator stats.
admin.transactions().getCoordinatorStats();
반환 값의 예시:
{
"state" : "Ready",
"leastSigBits" : 1,
"lowWaterMark" : 0,
"ongoingTxnSize" : 0,
"recoverStartTime" : 1657021892377,
"recoverEndTime" : 1657021892378
}
코디네이터 내부 통계 가져오기 (Get coordinator internal stats)
가져올 수 있는 코디네이터의 내부 통계는 다음과 같아요.
transactionLogStats: 트랜잭션 코디네이터 로그의 통계managedLedgerName: 트랜잭션 코디네이터 로그가 저장되는 managed ledger의 이름managedLedgerInternalStats: 트랜잭션 코디네이터 로그가 저장되는 managed ledger의 내부 통계. 자세한 내용은 managedLedgerInternalStats를 참고해요
다음 방법 중 하나로 코디네이터의 내부 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions coordinator-internal-stats -c 1 -m
REST API: GET /admin/v3/transactions/coordinatorInternalStats/{coordinatorId}
Java
admin.transactions().getCoordinatorInternalStats(coordinatorId, metadata);
반환 값의 예시:
{
"transactionLogStats" : {
"managedLedgerName" : "pulsar/system/persistent/__transaction_log_1",
"managedLedgerInternalStats" : {
"entriesAddedCounter" : 3,
"numberOfEntries" : 3,
"totalSize" : 63,
"currentLedgerEntries" : 3,
"currentLedgerSize" : 63,
"lastLedgerCreatedTimestamp" : "2022-06-30T18:18:05.88+08:00",
"waitingCursorsCount" : 0,
"pendingAddEntriesCount" : 0,
"lastConfirmedEntry" : "13:2",
"state" : "LedgerOpened",
"ledgers" : [ {
"ledgerId" : 13,
"entries" : 0,
"size" : 0,
"offloaded" : false,
"metadata" : "LedgerMetadata{formatVersion=3, ensembleSize=1, writeQuorumSize=1, ackQuorumSize=1, state=CLOSED, length=63, lastEntryId=2, digestType=CRC32C, password=OMITTED, ensembles={0=[10.20.240.119:3181]}, customMetadata={component=base64:bWFuYWdlZC1sZWRnZXI=, pulsar/managed-ledger=base64:cHVsc2FyL3N5c3RlbS9wZXJzaXN0ZW50L19fdHJhbnNhY3Rpb25fbG9nXzE=, application=base64:cHVsc2Fy}}",
"underReplicated" : false
} ],
"cursors" : {
"transaction.subscription" : {
"markDeletePosition" : "13:2",
"readPosition" : "13:3",
"waitingReadOp" : false,
"pendingReadOps" : 0,
"messagesConsumedCounter" : 3,
"cursorLedger" : 22,
"cursorLedgerLastEntry" : 1,
"individuallyDeletedMessages" : "[]",
"lastLedgerSwitchTimestamp" : "2022-06-30T18:18:05.932+08:00",
"state" : "Open",
"numberOfEntriesSinceFirstNotAckedMessage" : 1,
"totalNonContiguousDeletedMessagesRange" : 0,
"subscriptionHavePendingRead" : false,
"subscriptionHavePendingReplayRead" : false,
"properties" : { }
}
}
}
}
}
트랜잭션 pending ack 통계 (Transaction pending ack stats)
Pending ack는 트랜잭션이 완료되기 전에 트랜잭션 내의 메시지 승인을 유지해요. 메시지가 pending acknowledge 상태에 있으면 그 메시지가 pending acknowledge 상태에서 제거될 때까지 다른 트랜잭션이 승인할 수 없어요.
트랜잭션 pending ack 통계 가져오기 (Get transaction pending ack stats)
가져올 수 있는 트랜잭션 pending ack 상태 통계는 다음과 같아요.
state: 이 트랜잭션 코디네이터의 상태lowWaterMark: 이 트랜잭션 코디네이터의 낮은 워터마크ongoingTxnSize: 이 트랜잭션 코디네이터에서 진행 중인 트랜잭션의 총 수recoverStartTime: 트랜잭션 pendingAck 복구의 시작 타임스탬프. 0L은 시작 없음을 의미recoverEndTime: 트랜잭션 pendingAck 복구의 종료 타임스탬프. 0L은 시작 없음을 의미
다음 방법 중 하나로 트랜잭션 pending ack 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin.transactions()s pending-ack-stats -t my-topic -s mysubName -l
REST API: GET /admin/v3/transactions/pendingAckStats/{tenant}/{namespace}/{topic}/{subName}
Java
admin.transactions().getPendingAckStats(topic, subName, lowWaterMarks)
반환 값의 예시:
{
"state" : "Ready",
"lowWaterMarks" : {
"1" : 0
},
"ongoingTxnSize" : 1,
"recoverStartTime" : 1657021899202,
"recoverEndTime" : 1657021899203
}
트랜잭션 pending ack 내부 통계 가져오기 (Get transaction pending ack internal stats)
가져올 수 있는 트랜잭션 pending ack 내부 통계는 다음과 같아요.
transactionLogStats: 트랜잭션 pending ack 로그의 통계managedLedgerName: 트랜잭션 pending ack 로그가 저장되는 managed ledger의 이름managedLedgerInternalStats: 트랜잭션 코디네이터 로그가 저장되는 managed ledger의 내부 통계. 자세한 내용은 managedLedgerInternalStats를 참고해요
다음 방법 중 하나로 트랜잭션 pending ack 내부 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions pending-ack-internal-stats -t my-topic -s mysubName -m
REST API: GET /admin/v3/transactions/pendingAckInternalStats/{tenant}/{namespace}/{topic}/{subName}
Java
admin.transactions().getPendingAckInternalStats(topic, subName, boolean metadata);
반환 값의 예시:
{
"pendingAckLogStats" : {
"managedLedgerName" : "public/default/persistent/my-topic-mysubName__transaction_pending_ack",
"managedLedgerInternalStats" : {
"entriesAddedCounter" : 2247,
"numberOfEntries" : 2247,
"totalSize" : 37212,
"currentLedgerEntries" : 104,
"currentLedgerSize" : 1732,
"lastLedgerCreatedTimestamp" : "2022-06-30T19:02:09.746+08:00",
"waitingCursorsCount" : 0,
"pendingAddEntriesCount" : 52,
"lastConfirmedEntry" : "64:51",
"state" : "LedgerOpened",
"ledgers" : [ {
"ledgerId" : 56,
"entries" : 2195,
"size" : 36346,
"offloaded" : false,
"metadata" : "LedgerMetadata{formatVersion=3, ensembleSize=1, writeQuorumSize=1, ackQuorumSize=1, state=CLOSED, length=36346, lastEntryId=2194, digestType=CRC32C, password=OMITTED, ensembles={0=[10.20.240.119:3181]}, customMetadata={component=base64:bWFuYWdlZC1sZWRnZXI=, pulsar/managed-ledger=base64:cHVibGljL2RlZmF1bHQvcGVyc2lzdGVudC9teS10b3BpYy1teXN1Yk5hbWVfX3RyYW5zYWN0aW9uX3BlbmRpbmdfYWNr, application=base64:cHVsc2Fy}}",
"underReplicated" : false
}, {
"ledgerId" : 64,
"entries" : 0,
"size" : 0,
"offloaded" : false,
"metadata" : "LedgerMetadata{formatVersion=3, ensembleSize=1, writeQuorumSize=1, ackQuorumSize=1, state=CLOSED, length=866, lastEntryId=51, digestType=CRC32C, password=OMITTED, ensembles={0=[10.20.240.119:3181]}, customMetadata={component=base64:bWFuYWdlZC1sZWRnZXI=, pulsar/managed-ledger=base64:cHVibGljL2RlZmF1bHQvcGVyc2lzdGVudC9teS10b3BpYy1teXN1Yk5hbWVfX3RyYW5zYWN0aW9uX3BlbmRpbmdfYWNr, application=base64:cHVsc2Fy}}",
"underReplicated" : false
} ],
"cursors" : {
"__pending_ack_state" : {
"markDeletePosition" : "56:-1",
"readPosition" : "56:0",
"waitingReadOp" : false,
"pendingReadOps" : 0,
"messagesConsumedCounter" : 0,
"cursorLedger" : 57,
"cursorLedgerLastEntry" : 0,
"individuallyDeletedMessages" : "[]",
"lastLedgerSwitchTimestamp" : "2022-06-30T18:55:26.842+08:00",
"state" : "Open",
"numberOfEntriesSinceFirstNotAckedMessage" : 1,
"totalNonContiguousDeletedMessagesRange" : 0,
"subscriptionHavePendingRead" : false,
"subscriptionHavePendingReplayRead" : false,
"properties" : { }
}
}
}
}
}
pending ack의 위치 통계 가져오기 (Get position stats in pending ack)
pending ack의 위치 통계는 다음과 같아요.
PendingAck: 위치가 pending ack 통계에 있음MarkDelete: 위치가 이미 승인됨NotInPendingAck: 위치가 트랜잭션 내에서 승인되지 않음PendingAckNotReady: pending ack가 아직 초기화되지 않음InvalidPosition: 위치가 유효하지 않음, 예를 들어 배치 인덱스 > 배치 크기
위치가 승인되었는지 알고 싶다면 다음 방법 중 하나로 pending ack 위치 통계를 얻을 수 있어요.
pulsar-admin · REST API · Java
pulsar-admin transactions position-stats-in-pending-ack -t my-topic -s mysubName -l 15 -e 6
REST API: GET /admin/v3/transactions/positionStatsInPendingAck/{tenant}/{namespace}/{topic}/{subName}/{ledgerId}/{entryId}
Java
admin.transactions().getPositionStatsInPendingAckAsync(topic, subName, ledgerId, entryId, lowWaterMarks);
반환 값의 예시:
{
"State" : "MarkDelete"
}
트랜잭션 버퍼 통계 (Transaction buffer stats)
트랜잭션 버퍼는 트랜잭션 내에서 토픽 파티션에 생산된 메시지를 처리해요. 트랜잭션 버퍼의 메시지는 트랜잭션이 커밋될 때까지 컨슈머에게 보이지 않아요. 트랜잭션이 중단(abort)되면 트랜잭션 버퍼의 메시지는 버려져요.
트랜잭션 버퍼 통계 가져오기 (Get transaction buffer stats)
가져올 수 있는 트랜잭션 버퍼 통계는 다음과 같아요.
state: 이 트랜잭션 버퍼의 상태maxReadPosition: 이 트랜잭션 버퍼의 최대 읽기 위치lastSnapshotTimestamps: 이 트랜잭션 버퍼의 마지막 스냅샷 타임스탬프lowWaterMarks(선택): 이 트랜잭션 버퍼의 낮은 워터마크 상세ongoingTxnSize: 이 트랜잭션 버퍼에서 진행 중인 트랜잭션의 총 수recoverStartTime: 트랜잭션 버퍼 복구의 시작 타임스탬프. 0L은 시작 없음을 의미recoverEndTime: 트랜잭션 버퍼 복구의 종료 타임스탬프. 0L은 시작 없음을 의미
다음 방법 중 하나로 트랜잭션 버퍼 통계를 얻어요.
pulsar-admin · REST API · Java
pulsar-admin transactions transaction-buffer-stats -t my-topic -l
REST API: GET /admin/v3/transactions/transactionBufferStats/{tenant}/{namespace}/{topic}
Java
admin.transactions().getTransactionBufferStats(topic, lowWaterMarks);
반환 값의 예시:
{
"state" : "Ready",
"maxReadPosition" : "38:101",
"lastSnapshotTimestamps" : 1657021903534,
"lowWaterMarks" : {
"1" : -1,
"2" : -1
},
"ongoingTxnSize" : 0,
"recoverStartTime" : 1657021892850,
"recoverEndTime" : 1657021893372
}
더 알아보기 (Learn more)
- 트랜잭션 관련 전체 명령은 pulsar-admin의 transactions 명령 참고서를 확인해요.
- REST API 엔드포인트 상세는
/admin/v3/transactions문서를 참고해요. - Java admin API의 transactions 메서드는
PulsarAdmin객체 문서에서 볼 수 있어요. - 트랜잭션의 개념과 사용법이 궁금하다면 트랜잭션 개념 문서를 참고해요.