Amazon Kinesis에서 수집

Amazon Kinesis에서 수집 (Ingest from Amazon Kinesis)

Amazon Kinesis 스트림의 레코드를 Pinot 테이블로 수집하는 방법을 보여 주는 가이드예요.

출처: Ingest from Amazon Kinesis

본문

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. 부모 샤드에서 수집을 완전히 끝낼 때
  2. 1 다음에 RealtimeValidationManager가 실행될 때

부모가 자연스럽게 수집을 마치고 자식에서 수집을 시작하기 위해 RealtimeValidationManager를 기다리는 동안 이상 상태(ideal state)가 모든 세그먼트를 ONLINE으로 보여 주는 기간이 보일 거예요.

수집을 더 빨리 하려면 RealtimeValidationManager를 수동으로 호출할 수 있어요: docs

제한 사항

  1. ShardID는 "shardId-000000000001" 형식이에요. 우리는 숫자 부분을 partitionId로 사용해요. 우리의 partitionId 변수는 정수예요. shardId가 Integer.MAX\_VALUE를 넘어 커지면 partitionId 공간으로 오버플로돼요.
  2. 세그먼트 크기 기반 세그먼트 완료 임계값은 동작하지 않아요. 파티션 "0"이 항상 존재한다고 가정하기 때문이에요. 하지만 shard 0이 분할/병합되면 더 이상 파티션 0이 없어져요.

더 알아보기 (Learn more)