Amazon Kinesis에서 수집
Amazon Kinesis에서 수집 (Ingest from Amazon Kinesis)
Amazon Kinesis 스트림의 레코드를 Pinot 테이블로 수집하는 방법을 보여 주는 가이드예요.
본문
Amazon Kinesis 스트림의 이벤트를 Pinot으로 수집하려면 테이블 설정에 다음 설정을 넣어요:
{
"tableName": "kinesisTable",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "kinesis",
"stream.kinesis.topic.name": "<your kinesis stream name>",
"region": "<your region>",
"accessKey": "<your access key>",
"secretKey": "<your secret key>",
"shardIteratorType": "AFTER_SEQUENCE_NUMBER",
"stream.kinesis.fetch.timeout.millis": "30000",
"stream.kinesis.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.kinesis.consumer.factory.class.name": "org.apache.pinot.plugin.stream.kinesis.KinesisConsumerFactory",
"realtime.segment.flush.threshold.rows": "1000000",
"realtime.segment.flush.threshold.time": "6h"
}
},
"metadata": {
"customConfigs": {}
}
}
여기서 Kinesis 특정 속성은:
| 속성 (Property) | 설명 (Description) |
|---|---|
| streamType | "kinesis"로 설정해야 함 |
| stream.kinesis.topic.name | Kinesis 스트림 이름 |
| region | Kinesis 리전 예: us-west-1 |
| accessKey | Kinesis 액세스 키 |
| secretKey | Kinesis 시크릿 키 |
| shardIteratorType | Pinot이 샤드 이터레이터를 열 때 사용하는 AWS 샤드 이터레이터 타입. Pinot은 이 값을 Kinesis 클라이언트에 전달. 기본값: LATEST |
| maxRecordsToFetch | Kinesis에 대한 단일 getRecords API 호출에서 검색할 최대 레코드 수 지정. 데이터 검색의 배치 크기를 제어. 1~10,000 사이로 설정 가능 (Kinesis API 한도). 큰 값은 필요한 API 호출 수를 줄이지만 배치당 지연과 메모리 사용을 늘릴 수 있음. 기본값은 최대 10000. 메모리 제약이 있을 때만 낮추기 |
| requests_per_second_limit | Pinot이 샤드당 시도하는 최대 Kinesis 읽기 요청 수를 초당 제어. Pinot은 이 예산을 getRecords와 getShardIterator 읽기 모두에 적용하며, 0.25 같은 분수 값을 허용하고 기본값은 1.0입니다. Kinesis의 getRecords 샤드당 5 요청/초 한도를 Pinot 복제본, Pinot 테이블, 샤드를 공유하는 비-Pinot 소비자에 걸쳐 나누는 것부터 시작하세요. |
Kinesis가 여전히 ProvisionedThroughputExceededException을 반환하면 Pinot은 빈 배치를 즉시 반환하는 대신 페치 타임아웃까지 백오프하고 재시도해요. 이는 여러 소비자가 같은 샤드 예산을 공유해야 할 때 분수 requests_per_second_limit 값을 실용적으로 만들어 줘요.
Kinesis는 DefaultCredentialsProviderChain을 사용한 인증을 지원해요. 크레덴셜 프로바이더는 다음 순서로 크레덴셜을 찾아요:
- 환경 변수 -
AWS_ACCESS_KEY_ID와AWS_SECRET_ACCESS_KEY(.NET 제외 모든 AWS SDK와 CLI에서 인식되므로 RECOMMENDED), 또는AWS_ACCESS_KEY와AWS_SECRET_KEY(Java SDK만 인식) - Java 시스템 속성 -
aws.accessKeyId와aws.secretKey - 환경 또는 컨테이너의 Web Identity Token 크레덴셜
- 기본 위치
(~/.aws/credentials)의 크레덴셜 프로파일 파일 (모든 AWS SDK와 AWS CLI가 공유) AWS_CONTAINER_CREDENTIALS_RELATIVE_URI환경 변수가 설정되고 보안 매니저가 그 변수에 접근 권한이 있을 때 Amazon EC2 컨테이너 서비스를 통해 전달된 크레덴셜- Amazon EC2 메타데이터 서비스를 통해 전달된 인스턴스 프로파일 크레덴셜
{% hint style="info" %}
Pinot이 AWS Kinesis 데이터 스트림과 함께 작동하려면 모든 read access level 권한을 제공해야 합니다. 자세한 내용은 AWS 문서 참고.
{% endhint %}
위 속성에서 accessKey와 secretKey를 지정할 수도 있지만, 이 불안전한 방법은 권장하지 않아요. 프로덕션이 아닌 개념 검증(POC) 설정에만 사용을 권장해요. AWS\_SESSION\_TOKEN 같은 다른 AWS 필드도 환경 변수와 설정으로 지정할 수 있고 동작해요.
리샤딩 (Resharding)
Kinesis에서 스트림을 리샤딩할 때마다 샤드의 분할(split) 또는 병합(merge) 연산으로 이루어져요. 샤드를 분할하면 그 샤드는 닫히고 2개의 새 자식 샤드를 만들어요. shard0에서 시작해 분할하면 shard1과 shard2가 생겨요. 마찬가지로 2개 샤드를 병합하면 둘 다 닫히고 자식 샤드를 만들어요. 같은 예시에서 샤드 1과 2를 병합하면 shard3이 활성 샤드가 되고, shard0, shard1, shard2는 영원히 닫힌 상태로 남아요.
레시피는 여기를 참고하세요: https://dev.startree.ai/docs/pinot/recipes/github-events-stream-kinesis#resharding-kinesis-stream
Pinot에서 스트림의 리샤딩은 주기적 태스크 RealtimeValidationManager: docs로 감지돼요. 이것은 매시간 실행돼요. 리샤딩하면 다음 중 하나가 아니면 새 샤드가 감지되지 않아요:
- 부모 샤드에서 수집을 완전히 끝낼 때
- 1 다음에 RealtimeValidationManager가 실행될 때
부모가 자연스럽게 수집을 마치고 자식에서 수집을 시작하기 위해 RealtimeValidationManager를 기다리는 동안 이상 상태(ideal state)가 모든 세그먼트를 ONLINE으로 보여 주는 기간이 보일 거예요.
수집을 더 빨리 하려면 RealtimeValidationManager를 수동으로 호출할 수 있어요: docs
제한 사항
ShardID는 "shardId-000000000001" 형식이에요. 우리는 숫자 부분을partitionId로 사용해요. 우리의partitionId변수는 정수예요. shardId가Integer.MAX\_VALUE를 넘어 커지면 partitionId 공간으로 오버플로돼요.- 세그먼트 크기 기반 세그먼트 완료 임계값은 동작하지 않아요. 파티션 "0"이 항상 존재한다고 가정하기 때문이에요. 하지만 shard 0이 분할/병합되면 더 이상 파티션 0이 없어져요.