스트림 수집 가이드
스트림 수집 가이드 (Stream Ingestion Guide)
레코드 스트림을 Pinot 테이블로 수집하는 방법을 단계별로 보여 주는 가이드예요.
본문
이 가이드는 레코드 스트림을 Pinot 테이블로 수집하는 방법을 보여 줘요.
Apache Pinot은 사용자가 스트림에서 데이터를 소비해 데이터베이스에 바로 밀어 넣을 수 있게 해요. 이 과정을 스트림 수집(stream ingestion)이라고 해요. 스트림 수집은 발행 후 수 초 내에 데이터를 쿼리 가능하게 만들어 줘요.
스트림 수집은 데이터 손실을 막기 위한 체크포인트(checkpoint)를 지원해요.
스트림 수집을 설정하려면 다음 단계를 수행해요 (이 페이지에서 자세히 설명):
- 스키마 설정 생성
- 테이블 설정 생성
- 수집 설정 생성
- 테이블과 스키마 스펙 업로드
여기 수집할 데이터가 다음과 같은 형식이라고 가정하는 예시예요:
{"studentID":205,"firstName":"Natalie","lastName":"Jones","gender":"Female","subject":"Maths","score":3.8,"timestamp":1571900400000}
{"studentID":205,"firstName":"Natalie","lastName":"Jones","gender":"Female","subject":"History","score":3.5,"timestamp":1571900400000}
{"studentID":207,"firstName":"Bob","lastName":"Lewis","gender":"Male","subject":"Maths","score":3.2,"timestamp":1571900400000}
{"studentID":207,"firstName":"Bob","lastName":"Lewis","gender":"Male","subject":"Chemistry","score":3.6,"timestamp":1572418800000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Geography","score":3.8,"timestamp":1572505200000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"English","score":3.5,"timestamp":1572505200000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Maths","score":3.2,"timestamp":1572678000000}
{"studentID":209,"firstName":"Jane","lastName":"Doe","gender":"Female","subject":"Physics","score":3.6,"timestamp":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"Maths","score":3.8,"timestamp":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"English","score":3.5,"timestamp":1572678000000}
{"studentID":211,"firstName":"John","lastName":"Doe","gender":"Male","subject":"History","score":3.2,"timestamp":1572854400000}
{"studentID":212,"firstName":"Nick","lastName":"Young","gender":"Male","subject":"History","score":3.6,"timestamp":1572854400000}
스키마 설정 생성
스키마는 필드와 데이터 타입을 정의해요. 스키마는 또한 필드가 dimension, metric, 또는 timestamp인지도 정의해요. 스키마 설정에 대한 자세한 내용은 첫 테이블과 스키마 (First table and schema) 참고.
샘플 데이터에 대한 스키마 설정은 다음과 같아요:
{
"schemaName": "transcript",
"dimensionFieldSpecs": [
{
"name": "studentID",
"dataType": "INT"
},
{
"name": "firstName",
"dataType": "STRING"
},
{
"name": "lastName",
"dataType": "STRING"
},
{
"name": "gender",
"dataType": "STRING"
},
{
"name": "subject",
"dataType": "STRING"
}
],
"metricFieldSpecs": [
{
"name": "score",
"dataType": "FLOAT"
}
],
"dateTimeFieldSpecs": [{
"name": "timestamp",
"dataType": "LONG",
"format" : "1:MILLISECONDS:EPOCH",
"granularity": "1:MILLISECONDS"
}]
}
수집 설정이 포함된 테이블 설정 생성
다음 단계는 수집된 모든 데이터가 흘러들어 쿼리될 수 있는 테이블을 만드는 거예요. 각 테이블 구성 요소에 대한 자세한 내용은 table 레퍼런스 참고.
테이블 설정에는 스트리밍 데이터를 Pinot으로 수집하는 방법을 지정하는 수집 설정(ingestionConfig)이 포함돼요. 자세한 내용은 수집 설정 (ingestion configuration) 레퍼런스 참고.
ingestionConfig가 포함된 예시 테이블 설정
샘플 데이터와 스키마에 대한 테이블 설정은 다음과 같아요:
{
"tableName": "transcript",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"timeType": "MILLISECONDS",
"schemaName": "transcript",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
},
"metadata": {
"customConfigs": {}
},
"ingestionConfig": {
"streamIngestionConfig": {
"streamConfigMaps": [
{
"realtime.segment.flush.threshold.rows": "0",
"stream.kafka.decoder.prop.format": "JSON",
"key.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"streamType": "kafka",
"value.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"realtime.segment.flush.threshold.segment.rows": "50000",
"stream.kafka.broker.list": "localhost:9876",
"realtime.segment.flush.threshold.time": "3600000",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"stream.kafka.topic.name": "transcript-topic"
}
]
},
"transformConfigs": [],
"continueOnError": true,
"rowTimeValueCheck": true,
"segmentTimeValueCheck": false
},
"isDimTable": false
}
}
여러 스트림 설정이 포함된 예시 ingestionConfig
⚠️ 중요 버그 수정: 이전 구현에서 몇 가지 문제가 확인되었으며 PR #17953과 PR #17217로 수정되었어요. 이 수정은 Pinot 1.4.0에는 없어요. 이 기능을 사용하려면 이 PR들을 체리픽하거나 Pinot 1.5.0 릴리스를 사용하세요.
Pinot은 streamConfigMaps에 하나 이상의 스트림 설정을 포함할 수 있어요. 아래 예시에서는 샘플 데이터가 transcript-topic1과 transcript-topic2 두 Kafka 토픽에 복제되었고, 테이블이 두 토픽 모두에서 수집한다고 가정해요:
{
"tableName": "transcript",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"timeType": "MILLISECONDS",
"schemaName": "transcript",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
},
"metadata": {
"customConfigs": {}
},
"ingestionConfig": {
"streamIngestionConfig": {
"streamConfigMaps": [
{
"realtime.segment.flush.threshold.rows": "0",
"stream.kafka.decoder.prop.format": "JSON",
"key.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"streamType": "kafka",
"value.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"realtime.segment.flush.threshold.segment.rows": "50000",
"stream.kafka.broker.list": "localhost:9876",
"realtime.segment.flush.threshold.time": "3600000",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"stream.kafka.topic.name": "transcript-topic1"
},
{
"realtime.segment.flush.threshold.rows": "0",
"stream.kafka.decoder.prop.format": "JSON",
"key.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"stream.kafka.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"streamType": "kafka",
"value.serializer": "org.apache.kafka.common.serialization.ByteArraySerializer",
"realtime.segment.flush.threshold.segment.rows": "50000",
"stream.kafka.broker.list": "localhost:9876",
"realtime.segment.flush.threshold.time": "3600000",
"stream.kafka.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kafka30.KafkaConsumerFactory",
"stream.kafka.consumer.prop.auto.offset.reset": "smallest",
"stream.kafka.topic.name": "transcript-topic2"
}
]
},
"transformConfigs": [],
"continueOnError": true,
"rowTimeValueCheck": true,
"segmentTimeValueCheck": false
},
"isDimTable": false
}
}
streamConfigMaps에 여러 항목이 있으면 Pinot은 테이블 생성·업데이트를 허용하기 전에 결합된 설정을 검증해요:
- 각 항목은 그 자체로 여전히 유효한 스트림 설정이어야 함.
- 모든 항목은 같은
streamType을 사용해야 함. - 세그먼트 플러시 설정이 모든 항목에서 일치해야 함:
realtime.segment.flush.threshold.rowsrealtime.segment.flush.threshold.timerealtime.segment.flush.threshold.variance.fractionrealtime.segment.flush.threshold.segment.sizerealtime.segment.flush.threshold.segment.rows
- 토픽 이름은 항목 전체에서 고유해야 함.
- 여러 스트림 설정은
pauselessConsumptionEnabled=true와 함께 지원되지 않음. - 여러 스트림 설정은 업서트 테이블에서 지원되지 않음.
검증 후 Pinot은 나머지 수집 설정을 정상적으로 적용해요:
- 변환 함수는 설정된 모든 스트림의 레코드에 적용됨.
- 기존 인스턴스 배정 전략은 평소처럼 계속 동작.
- 파티션 변경은 단일 스트림 설정과 같은 방식으로 처리됨.
- 기본 수집은 여전히
LOWLEVEL모드에서 동작하며:transcript-topic1세그먼트는transcript__0__0__20250101T0000Z처럼 이름이 붙음transcript-topic2세그먼트는transcript__10000__0__20250101T0000Z처럼 이름이 붙음
스키마와 테이블 설정 업로드
이제 테이블과 스키마 설정이 있으니 Pinot 클러스터에 업로드해요. 설정이 업로드되는 즉시 Pinot이 토픽에서 사용 가능한 레코드 수집을 시작해요.
Docker:
docker run \
--network=pinot-demo \
-v /tmp/pinot-quick-start:/tmp/pinot-quick-start \
--name pinot-streaming-table-creation \
apachepinot/pinot:latest AddTable \
-schemaFile /tmp/pinot-quick-start/transcript-schema.json \
-tableConfigFile /tmp/pinot-quick-start/transcript-table-realtime.json \
-controllerHost pinot-quickstart \
-controllerPort 9000 \
-exec
Launcher 스크립트:
bin/pinot-admin.sh AddTable \
-schemaFile /path/to/transcript-schema.json \
-tableConfigFile /path/to/transcript-table-realtime.json \
-exec
스트림 설정 튜닝
파티션이 너무 많이 뒤쳐지면 건너뛰기
실시간 파티션이 너무 많이 뒤쳐지면 Pinot은 세그먼트 커밋 시점에 백로그를 건너뛰고 이전 세그먼트의 nextOffset 대신 최신 스트림 오프셋에서 다음 세그먼트를 시작할 수 있어요.
streamConfigMaps 안의 다음 스트림 설정을 사용해요:
{
"realtime.segment.offsetAutoReset.enable": "true",
"realtime.segment.offsetAutoReset.offsetThreshold": "1000000",
"realtime.segment.offsetAutoReset.timeThresholdSeconds": "3600"
}
Pinot은 소비 중인 세그먼트를 씰링(seal)할 때 임계값을 확인해요:
latestOffset - nextOffset가realtime.segment.offsetAutoReset.offsetThreshold를 초과하면 최신 오프셋으로 건너뜀.- 다음 오프셋이
realtime.segment.offsetAutoReset.timeThresholdSeconds보다 오래됐으면 최신 오프셋으로 건너뜀. - 임계값 중 하나라도 양수로 설정. 둘 다 미설정이거나 양수가 아니면 Pinot은 정상
nextOffset을 유지.
수집 백로그를 재생 대신 버리고 싶을 때 유용해요.
스트림 소비 중 디코드 오류 처리
기본적으로 Pinot은 디코드 오류를 기록하고 문제가 있는 행을 조용히 버려 수집이 중단 없이 계속되게 해요. 하지만 일부 시나리오에서는 디코드 오류 발생 시 수집을 즉시 중지해 조사·해결하고 싶을 수 있어요.
스트림 설정의 stopOnDecodeError 구성 파라미터로 이 동작을 제어할 수 있어요:
{
"tableName": "transcript",
"tableType": "REALTIME",
...
"ingestionConfig": {
"streamIngestionConfig": {
"streamConfigMaps": [{
"streamType": "kafka",
"stream.kafka.topic.name": "transcript-topic",
"stopOnDecodeError": "true",
...
}]
}
},
...
}
설정 옵션:
"stopOnDecodeError": "true"- 디코드 오류 발생 시 소비가 즉시 중지되고 전체 스택 트레이스로 오류가 기록돼요. 이를 통해 이어서 조사하고 재개하기 전에 서버 로그에서 근본 원인을 확인할 수 있어요."stopOnDecodeError": "false"(기본) - 첫 디코드 오류는 기록되고 이후 오류는 조용히 버려져요. 디코드 오류가 있는 행은 건너뛰고 수집은 정상적으로 계속돼요.
"stopOnDecodeError": "true"를 사용할 때:
- 데이터 품질이 중요하고 디코드 오류를 즉시 잡아야 할 때
- 새 디코더 구현을 개발하거나 테스트 중일 때
- 예상치 못한 데이터 형식 문제를 조사하고 싶을 때
"stopOnDecodeError": "false"(기본)를 사용할 때:
- 스트림에 가끔 잘못된 메시지가 나타날 것으로 예상할 때
- 개별 불량 레코드에도 불구하고 수집이 탄력적으로 계속되길 원할 때
- 몇 개 버려진 행의 데이터 손실이 허용될 때
스트림 소비 조절 (Throttle)
입력 스트림의 메시지율이 버스트로 몰리면 Pinot 서버에 긴 GC 일시정지를 일으키거나 같은 서버의 다른 실시간 테이블 수집율에 영향을 줄 수 있는 시나리오가 있어요. 이런 경우 스트림 수집 중 소비율을 조절해 전반적인 성능을 더 잘 관리할 수 있어요.
두 가지 독립적인 조절 메커니즘이 있어요:
- 메시지율 기반 조절 (테이블 수준, records/sec)
- 바이트율 기반 조절 (서버 수준, bytes/sec)
두 메커니즘을 동시에 활성화할 수 있어요.
메시지율 기반 조절 (테이블 수준)
스트림 소비 조절은 파티션별 상한을 위한 partition.consumption.rate.limit 또는 토픽 전체 상한(Pinot이 이를 토픽 파티션 수로 나눔)을 위한 topic.consumption.rate.limit 스트림 설정으로 조정할 수 있어요.
소비 조절을 설정하는 샘플 설정:
{
"tableName": "transcript",
"tableType": "REALTIME",
...
"ingestionConfig": {
"streamIngestionConfig":,
"streamConfigMaps": {
"streamType": "kafka",
"stream.kafka.topic.name": "transcript-topic",
...
"topic.consumption.rate.limit": 1000
}
},
...
이 설정을 조정할 때 유의할 점:
-
topic.consumption.rate.limit을 사용하면 Pinot이 토픽 전체 속도를 토픽의 파티션 수로 나누고 결과를 각 파티션 소비자에 적용해요. 복제 팩터는 고려하지 않아요. 예시 topic.consumption.rate.limit - 1000 num partitions in Kafka topic - 4 replication factor in table - 3Pinot은 각 파티션에 고정 상한 1000 / 4 = 초당 250 records를 부과해요.
-
멀티 테넌트 배포(같은 서버 인스턴스에 테이블이 1개 이상)에서는 한 테이블의 속도 제한이 다른 테이블의 속도 제한을 방해/기아시키지 않도록 해야 해요. 같은 서버에 테이블이 1개 이상 있을 때(대부분 그렇지만) 모든 스트리밍 테이블의 조절 임계값을 재조정해야 할 수 있어요.
-
pinot.server.consumption.rate.limit설정은 테이블 설정이 아니라 서버의 인스턴스 설정에 구성해야 해요. 이 서버 전체의 rows/sec 상한은 테이블 수준partition.consumption.rate.limit또는topic.consumption.rate.limit상한에 추가로 적용돼요.
테이블에 대해 조절이 활성화되면 다음과 비슷한 로그를 찾아 확인할 수 있어요:
A consumption rate limiter is set up for topic <topic_name> in table <tableName> with rate limit: <rate_limit> (topic rate limit: <topic_rate_limit>, partition count: <partition_count>)
또한 유효한 파티션별 설정 상한의 CONSUMPTION_RATE_LIMIT (consumptionRateLimit)과 파티션별 사용률의 CONSUMPTION_QUOTA_UTILIZATION (consumptionQuotaUtilization)을 모니터링할 수 있어요.
스트림 설정의 topic.consumption.rate.limit 변경은 즉시 적용되지 않습니다. 새 설정은 다음 소비 중인 세그먼트부터 반영돼요. 새 설정을 강제하려면 forceCommit API를 트리거해야 해요. 자세한 내용은 Pause Stream Ingestion 참고.
$ curl -X POST {controllerHost}/tables/{tableName}/forceCommit
바이트율 기반 조절 (서버 수준)
메시지율 조절에 더해 Pinot은 서버 수준의 바이트 기반 스트림 소비 조절을 지원해요.
이 조절 메커니즘은 Pinot 서버가 해당 서버에서 호스팅하는 모든 실시간 테이블과 파티션에 걸쳐 초당 소비하는 총 바이트 수를 제한해요.
바이트 기반 조절을 쓸 때
바이트 기반 조절은 특히 다음과 같은 경우에 유용해요:
- 메시지 크기가 크게 다양할 때
- 수집 압력이 레코드 수보다 페이로드 크기에 의해 결정될 때
- 서버 수준에서 네트워크, 다이렉트 메모리, 또는 디스크 IO 사용을 상한선으로 잡고 싶을 때
- 여러 실시간 테이블이 같은 서버에 공존할 때
설정
바이트 기반 조절은 테이블이나 스트림 설정이 아니라 클러스터 설정으로 구성돼요.
설정 키
pinot.server.consumption.rate.limit.bytes
값은 초당 바이트로 지정돼요.
설정 업데이트
설정은 클러스터 설정 API를 사용해 동적으로 업데이트할 수 있어요.
이것은 각 Pinot 서버가 모든 실시간 테이블에 걸쳐 초당 최대 3,000,000 bytes (~3 MB/sec)를 소비하도록 제한해요.
curl 예시
curl -X POST
'{controllerHost}/cluster/configs'
-H 'Content-Type: application/json'
-d '{
"pinot.server.consumption.rate.limit.bytes": "3000000"
}'
바이트 기반 조절 동작 방식
- 바이트 속도 제한은 서버별로 강제됨
- 제한은 해당 서버에서 호스팅되는 모든 소비 파티션과 테이블에 집합적으로 적용됨
- 이 조절은 테이블 수준 메시지율 조절과 독립적임
메시지율 조절과의 상호작용
두 조절이 모두 활성화되면:
- 테이블 수준
topic.consumption.rate.limit이 테이블별 records/sec를 제어 - 서버 수준
pinot.server.consumption.rate.limit.bytes가 서버별 bytes/sec를 제어 - Pinot이 두 제한 모두 적용
- 두 제한 중 하나라도 도달하면 즉시 소비가 조절됨
이를 통해 메시지 개수와 페이로드 크기 모두가 중요할 때 정밀한 제어가 가능해요.
동적 업데이트와 전파
- 바이트 기반 조절은 클러스터 설정 변경 리스너를 통해 동적으로 업데이트됨
- 서버 재시작 불필요
- 서버가 업데이트된 클러스터 설정을 받으면 변경이 자동으로 적용됨
조절 확인
활성화되면 Pinot은 서버 수준 바이트 소비 리미터가 적용되었음을 나타내는 메시지를 기록해요.
다음 메트릭으로 조절 동작을 모니터링할 수도 있어요:
SERVER_CONSUMPTION_RATE_LIMIT(serverConsumptionRateLimit)는 설정된 서버 전체 상한을 보고. 서버 수준 속도 제한이 비활성화되면 Pinot은 이를-1로 설정해요.SERVER_CONSUMPTION_QUOTA_UTILIZATION(serverConsumptionQuotaUtilization)은 서버 전체 사용률 백분율을 보고.
서버 전체 값을 consumptionQuotaUtilization{table="realtimeRowsConsumed"}에서 읽었다면, 대시보드와 알림을 serverConsumptionQuotaUtilization으로 전환하세요.
커스텀 수집 지원
사용 중인 플랫폼이 기본 지원되지 않으면 수집 플러그인을 직접 작성할 수도 있어요. 워크스루는 Stream Ingestion Plugin 참고.
스트림 수집 일시정지
테이블이 쿼리 가능한 상태에서 실시간 수집을 잠시 멈추고 싶은 시나리오가 있어요. 예를 들어 스트림 수집에 문제가 있고, 문제를 해결하는 동안 이미 수집된 데이터에 대한 쿼리는 계속 실행하고 싶은 경우예요. 이런 경우 먼저 컨트롤러 호스트에 Pause 요청을 보내요. 스트림 문제 해결이 끝나면 컨트롤러에 다시 요청을 보내 소비를 재개할 수 있어요.
$ curl -X POST {controllerHost}/tables/{tableName}/pauseConsumption
$ curl -X POST "{controllerHost}/tables/{tableName}/resumeConsumption?comment=maintenance-complete"
소비 중인 세그먼트가 많은 테이블에서는 일시정지 요청을 배치 처리해 컨트롤러가 한 번에 커밋하는 세그먼트 수를 줄일 수 있어요:
$ curl -X POST "{controllerHost}/tables/{tableName}/pauseConsumption?batchSize=50&batchStatusCheckIntervalSec=5&batchStatusCheckTimeoutSec=180"
batchSize는 Pinot이 한 배치에서 커밋하는 소비 중인 세그먼트 수를 제한해요. batchStatusCheckIntervalSec과 batchStatusCheckTimeoutSec은 컨트롤러가 각 배치가 끝나기까지 기다리는 주기와 시간을, 그리고 실패하거나 pause 요청을 실패시키기 전까지 제어해요.
Pause 요청이 발행되면 컨트롤러는 테이블을 호스팅하는 실시간 서버에 소비 중인 세그먼트를 즉시 커밋하도록 지시해요. 하지만 커밋 프로세스는 완료까지 시간이 걸릴 수 있어요. Pause와 Resume 요청은 비동기라는 점에 유의하세요. OK 응답은 일시정지/재개 지시가 실시간 서버에 성공적으로 전송되었음을 의미해요. 소비가 실제로 중지되었거나 재개되었는지 알고 싶으면 pause 상태 요청을 발행해요.
$ curl -X GET {controllerHost}/tables/{tableName}/pauseStatus
일시정지 상태 응답에는 현재 pauseFlag, consumingSegments 집합, 저장된 reasonCode, 지속된 comment, 그리고 현재 일시정지 상태의 컨트롤러 측 timestamp가 포함돼요. 토픽 수준 일시정지 API가 일부 스트림 토픽만 일시정지했을 때 Pinot은 비활성 상태로 남은 0-기반 토픽 인덱스와 함께 indexOfInactiveTopics도 반환해요. 명시적 주석이 저장되지 않았으면 Pinot은 기본적인 일시정지/미일시정지 메시지를 반환해요.
실시간 서버의 소비 중인 세그먼트는 휘발성 메모리에 저장되고, 그 리소스는 소비 중인 세그먼트가 처음 생성될 때 할당된다는 점에 유의해야 해요. 소비 도중에 소비 파라미터를 바꾸면 이 리소스를 변경할 수 없어요. 이런 변경이 효과를 내기까지 몇 시간이 걸릴 수 있어요. 게다가 파라미터를 비호환 방식으로 바꾸면(예: 완전히 새로운 오프셋 집합으로 기본 스트림을 바꾸거나, 메시지를 소비할 스트림 엔드포인트 변경) 테이블이 오류 상태에 빠질 수 있어요.
pause와 resume 기능은 이런 경우에 유용해요. 운영자가 pause 요청을 발행하면 새 변경 가능(mutable) 세그먼트를 시작하지 않고 소비 중인 세그먼트를 커밋해요. 대신 resume 요청이 발행될 때만 새 변경 가능 세그먼트가 시작돼요. 이 메커니즘은 운영자와 개발자 모두에게 더 큰 유연성을 제공해요. 또한 Pinot이 기본 스트림이 부과하는 운영·기능 제약에 더 탄력적으로 대응하게 해요.
pause와 resume 기능의 기본 요소를 활용하는 또 다른 기능인 Force Commit이 있어요. 운영자가 force commit 요청을 발행하면 현재 변경 가능 세그먼트를 커밋하고 즉시 새 것을 시작해요. 운영자는 이제 이 기능을 사용해 호환되는 모든 테이블 설정 파라미터 변경을 즉시 적용할 수 있어요.
$ curl -X POST {controllerHost}/tables/{tableName}/forceCommit
선택적 필터와 배치(partitions와 segments는 함께 쓰지 마세요):
$ curl -X POST "{controllerHost}/tables/{tableName}/forceCommit?partitions=0,1&batchSize=50&batchStatusCheckIntervalSec=5&batchStatusCheckTimeoutSec=180"
$ curl -X POST "{controllerHost}/tables/{tableName}/forceCommit?segments=table__0__12__20250610T2140Z,table__1__12__20250610T2140Z"
실시간 테이블은 Pinot 자체에 의해 자동으로 일시정지될 수도 있어요. 특히 테이블이 quota.storage를 초과하면 컨트롤러가 주기적 검증 중에 그 테이블을 STORAGE_QUOTA_EXCEEDED 이유 코드로 paused로 표시하고 새 소비 중인 세그먼트 생성을 중지해요. 테이블이 쿼터 안으로 돌아오면 Pinot이 그 pause 상태를 해제하고 세그먼트 생성을 재개해요. 쿼터 문제를 고친 후 더 빨리 재개하고 싶다면 resumeConsumption을 수동으로 호출할 수도 있어요.
(v 0.12.0+) 제출되면 forceCommit API는 forceCommit 작업의 현재 진행 상황을 조회하는 데 쓸 수 있는 jobId를 반환해요. 샘플 응답과 상태 API 호출:
$ curl -X POST {controllerHost}/tables/{tableName}/forceCommit
{
"forceCommitJobId": "6757284f-b75b-45ce-91d8-a277bdbc06ae",
"forceCommitStatus": "SUCCESS",
"jobMetaZKWriteStatus": "SUCCESS"
}
$ curl -X GET {controllerHost}/tables/forceCommitStatus/6757284f-b75b-45ce-91d8-a277bdbc06ae
{
"jobId": "6757284f-b75b-45ce-91d8-a277bdbc06ae",
"segmentsForceCommitted": "[\"airlineStats__0__0__20230119T0700Z\",\"airlineStats__1__0__20230119T0700Z\",\"airlineStats__2__0__20230119T0700Z\"]",
"submissionTimeMs": "1674111682977",
"numberOfSegmentsYetToBeCommitted": 0,
"jobType": "FORCE_COMMIT",
"segmentsYetToBeCommitted": [],
"tableName": "airlineStats_REALTIME"
}
{% hint style="info" %} forceCommit 요청은 소비 중인 세그먼트가 종료 기준에 도달하기 전에 정규 커밋을 트리거할 뿐이므로, 일반 커밋과 같은 메커니즘을 따릅니다. 일회성 요청이며 실패 시 자동 재시도되지 않습니다. 필요하다면 성공할 때까지 계속 발행해도 될 만큼 멱등(idempotent)합니다.
HTTP 200은 비동기 승인이지 커밋 완료의 증거가 아닙니다. forceCommitStatus=SUCCESS는 컨트롤러가 작업을 시작했음을 의미하며, 상태 API에서 numberOfSegmentsYetToBeCommitted이 0이 될 때까지 기다리세요. jobMetaZKWriteStatus=FAILED는 커밋이 여전히 실행될 수 있지만 Pinot이 추적 가능한 job id를 지속하지 못했다는 뜻입니다.
ZK 상태 항목은 제출 시간과 포함된 소비 중인 세그먼트를 기록합니다. 진행 상황은 그 목록을 최신 IdealState/메타데이터와 비교해 파생됩니다. 상태 항목은 성공 또는 실패 시 삭제되지 않아 오래될 수 있으며, Pinot은 ZK에 force-commit 잡의 수를 제한합니다(기본 100 via controller.force.commit.maxJobsInZK).
전체 파라미터 테이블, partitions vs segments의 상호 배제, 실패 모드: Force commit API.
{% endhint %}
비호환 파라미터 변경으로는 완전히 새로운 오프셋 집합을 처리하는 옵션이 resume 요청에 추가돼 있어요. 운영자는 이제 3단계 절차를 따를 수 있어요: 첫째, pause 요청 발행. 둘째, 소비 파라미터 변경. 마지막으로, 적절한 옵션으로 resume 요청 발행. 이 단계들은 이전 데이터를 보존하고 새 데이터를 즉시 소비할 수 있게 해줘요. 작업 내내 쿼리는 계속 서빙돼요.
$ curl -X POST {controllerHost}/tables/{tableName}/resumeConsumption?consumeFrom=smallest
$ curl -X POST {controllerHost}/tables/{tableName}/resumeConsumption?consumeFrom=largest
스트림의 파티션 변경 처리
Pinot 테이블이 Low Level (파티션 기반) 스트림 타입으로 소비하도록 설정되어 있다면, 테이블의 파티션이 시간에 따라 바뀔 수 있어요. 예를 들어 Kafka에서는 파티션 수가 늘어날 수 있어요. Kinesis에서는 파티션 수가 늘거나 줄 수 있어요 -- 일부 파티션이 병합되어 새 것을 만들거나, 기존 파티션이 분할되어 새 것을 만드는 식.
Pinot은 RealtimeSegmentValidationManager라는 주기적 태스크를 실행해 이런 변경을 모니터링하고 필요에 따라 새 파티션에서 소비를 시작(또는 이전 파티션에서 소비 중지)해요. 이는 컨트롤러 주기적 태스크이므로, Pinot이 새 파티션을 인식하고 소비를 시작하는 데 시간이 걸릴 수 있어요. 이로 인해 새 파티션의 데이터가 pinot이 반환하는 결과에 나타나는 것이 지연될 수 있어요.
새 파티션을 더 빨리 인식하고 싶다면 수동으로 주기적 태스크를 트리거해 그러한 데이터를 즉시 인식하세요.
Kafka 저수준 소비자의 경우, 같은 태스크가 업스트림에 여전히 존재하는 Kafka 파티션에 대해 Pinot이 CONSUMING 세그먼트를 놓치고 있을 때 파티션을 복구할 수도 있어요. 운영자는 다음 RealtimeSegmentValidationManager 실행을 기다리거나 수동으로 트리거할 수 있어요. Pinot이 누락된 소비 중인 세그먼트를 다시 만들 때, 검증이 선택한 복구 오프셋을 사용해요. Pinot이 그 파티션의 LLC 메타데이터를 아직 갖고 있으면 저장된 끝 오프셋에서 재개할 수 있고, 그렇지 않으면 보통 Kafka의 가장 오래된 보존 오프셋에서 시작하므로 복구는 Kafka 보존에 의해 제한돼요.
실시간 테이블의 수집 상태 추론
실시간 테이블의 데이터 수집 속도를 이해하는 것이 중요할 때가 많아요. 보통 소비자의 소비 랙(consumption lag)을 보고 이를 파악해요. 랙 자체는 여러 차원에서 관찰할 수 있어요. Pinot은 가능할 때마다(커넥터 특성에 따라 다름) 오프셋 차원과 시간 차원으로 소비 랙을 관찰할 수 있게 지원해요.
커넥터의 수집 상태는 /consumingSegmentsInfo API 또는 테이블의 /debug API를 조회해 관찰할 수 있어요:
# GET /tables/{tableName}/consumingSegmentsInfo
curl -X GET "http://<controller_url:controller_admin_port>/tables/meetupRsvp/consumingSegmentsInfo" -H "accept: application/json"
# GET /debug/tables/{tableName}
curl -X GET "http://localhost:9000/debug/tables/meetupRsvp?type=REALTIME&verbosity=1" -H "accept: application/json"
Kafka 기반 실시간 테이블의 샘플 응답은 다음과 같아요. 수집 상태는 테이블의 각 CONSUMING 세그먼트에 대해 표시돼요.
{
"_segmentToConsumingInfoMap": {
"meetupRsvp__0__0__20221019T0639Z": [
{
"serverName": "Server_192.168.0.103_7000",
"consumerState": "CONSUMING",
"lastConsumedTimestamp": 1666161593904,
"partitionToOffsetMap": { // <<-- Deprecated. See currentOffsetsMap for same info
"0": "6"
},
"partitionOffsetInfo": {
"currentOffsetsMap": {
"0": "6" // <-- Current consumer position
},
"latestUpstreamOffsetMap": {
"0": "6" // <-- Upstream latest position
},
"recordsLagMap": {
"0": "0" // <-- Lag, in terms of #records behind latest
},
"recordsAvailabilityLagMap": {
"0": "2" // <-- Lag, in terms of time
}
}
}
],
| 용어 (Term) | 설명 (Description) |
|---|---|
| currentOffsetsMap | 파티션별 현재 소비 오프셋 위치 |
| latestUpstreamOffsetMap | (해당 시) 업스트림 토픽 파티션에서 발견된 최신 오프셋 |
| recordsLagMap | (가능 시) 현재 레코드의 오프셋/포인터가 업스트림 최신 레코드보다 얼마나 뒤처졌는지 정의. 랙 계산 요청 시점에 파티션의 latestUpstreamOffset와 currentOffset 차이로 계산 |
| recordsAvailabilityLagMap | (가능 시) 레코드 수집 직후 Pinot이 그 레코드를 소비했는지. 레코드가 소비된 시간과 업스트림에서 수집된 시간의 차이로 계산 |
Pinot이 랙 값을 계산할 수 없으면 숫자 자리표시자 대신 문자열 센티널 NOT_CALCULATED를 반환해요. 예를 들어 스트림 커넥터가 파티션의 최신 업스트림 오프셋을 제공할 수 없거나, 마지막으로 소비된 레코드에 유효한 업스트림 수집 타임스탬프가 없을 때 발생해요. Pinot은 /consumingSegmentsInfo와 /debug 응답 모두에서 같은 센티널을 사용해요.
실시간 수집 모니터링
실시간 수집은 메시지 처리의 3단계인 Decode, Transform, Index를 포함해요.
각 단계에서 실패가 발생할 수 있고, 수집 실패로 이어질 수도 있고 아닐 수도 있어요. 수집 문제를 조사하는 데 다음 메트릭을 사용할 수 있어요:
- Decode 단계 -> 여기서 오류는
INVALID_REALTIME_ROWS_DROPPED로 기록됨 - Transform 단계 -> 가능한 오류:
- FILTER 변환으로 메시지가 버려지면
REALTIME_ROWS_FILTERED로 기록됨 - 변환 파이프라인이 메시지에
$INCOMPLETE_RECORD_KEY$키를 설정하면,continueOnError설정이 활성화된 경우에만INCOMPLETE_REALTIME_ROWS_CONSUMED로 기록됨.continueOnError가 활성화되지 않으면 수집이 실패함
- FILTER 변환으로 메시지가 버려지면
- Index 단계 -> 이 단계에서 실패가 있으면 수집은 보통 중지되고 파티션을 ERROR로 표시
ROWS_WITH_ERROR라는 또 다른 메트릭도 있는데, 위 3개 단계의 모든 오류 수의 합이에요.
또한 소비 중 일시적/영구적 스트림 예외가 보일 때마다 REALTIME_CONSUMPTION_EXCEPTIONS 메트릭이 증가해요.
이 메트릭들은 서버 로그를 파고들기 전에 특정 테이블 파티션에서 수집이 왜 실패했는지 이해하는 데 쓸 수 있어요.