Amazon Kinesis 수집
Amazon Kinesis 수집
Kinesis indexing service를 활성화하면 Overlord에서 supervisor를 구성해서 Kinesis indexing task의 생성과 수명을 관리할 수 있어요. Kinesis indexing task는 Kinesis shard와 시퀀스 번호(sequence number) 메커니즘으로 이벤트를 읽어 정확히 한 번(Exactly-once) 수집을 보장합니다. supervisor는 indexing task의 상태를 감독하면서 hand-off를 조정하고, 실패를 관리하며, 확장성·복제 요구사항이 유지되도록 보장해요.
이 주제는 Apache Druid의 Kinesis indexing service supervisor에 대한 구성 정보를 담고 있습니다.
출처: 문서
본문
설정 (Setup)
Kinesis indexing service를 사용하려면 먼저 Overlord와 Middle Manager 양쪽에 druid-kinesis-indexing-service 코어 확장을 로드해야 해요. 자세한 내용은 Loading extensions를 참고하세요.
druid-kinesis-indexing-service 확장을 프로덕션에 배포하기 전에 Known issues를 검토하세요.
Supervisor spec 구성
이 섹션은 Amazon Kinesis 스트리밍 수집 방법에 특화된 구성 속성을 다룹니다. Druid가 지원하는 모든 스트리밍 수집 방법에 공통으로 적용되는 속성은 Supervisor spec을 참고해 주세요.
다음 예시는 KinesisStream이라는 이름의 스트림에 대한 supervisor spec이에요.
예시 보기
{ "type": "kinesis", "spec": { "ioConfig": { "type": "kinesis", "stream": "KinesisStream", "inputFormat": { "type": "json" }, "useEarliestSequenceNumber": true }, "tuningConfig": { "type": "kinesis" }, "dataSchema": { "dataSource": "KinesisStream", "timestampSpec": { "column": "timestamp", "format": "iso" }, "dimensionsSpec": { "dimensions": [ "isRobot", "channel", "flags", "isUnpatrolled", "page", "diffUrl", { "type": "long", "name": "added" }, "comment", { "type": "long", "name": "commentLength" }, "isNew", "isMinor", { "type": "long", "name": "delta" }, "isAnonymous", "user", { "type": "long", "name": "deltaBucket" }, { "type": "long", "name": "deleted" }, "namespace", "cityName", "countryName", "regionIsoCode", "metroCode", "countryIsoCode", "regionName" ] }, "granularitySpec": { "queryGranularity": "none", "rollup": false, "segmentGranularity": "hour" } } }}
I/O 구성
다음 표는 Kinesis에 특화된 ioConfig 구성 속성을 정리한 것이에요. 모든 스트리밍 수집 방법에 공통인 속성은 Supervisor I/O configuration을 참고하세요.
| Property | Type | Description | Required | Default |
| stream | String | 읽을 Kinesis 스트림 | Yes | |
| endpoint | String | 특정 리전의 AWS Kinesis 스트림 엔드포인트. 엔드포인트 목록은 AWS service endpoints 문서에서 확인 가능 | No | kinesis.us-east-1.amazonaws.com |
| useEarliestSequenceNumber | Boolean | supervisor가 datasource를 처음 관리할 때 Kinesis로부터 시작 시퀀스 번호 집합을 얻음. 이 플래그는 supervisor가 Kinesis에서 가장 이른 시퀀스 번호를 가져올지 최신 번호를 가져올지 결정. 정상 상황에서는 이후 태스크들이 이전 세그먼트가 끝난 위치에서 시작하므로 이 플래그는 첫 실행에서만 사용됨 | No | false |
| fetchDelayMillis | Integer | Kinesis에서 레코드를 가져오는 연속 호출 사이에 대기할 시간(밀리초). Determine fetch settings 참고 | No | 0 |
| awsAssumedRoleArn | String | 추가 권한에 사용할 AWS assumed role | No | |
| awsExternalId | String | 추가 권한에 사용할 AWS external ID | No | |
데이터 형식
Kinesis indexing service는 inputFormat을 지원해요. 자세한 내용은 Source input formats을 참고하세요.
Kinesis indexing service는 inputFormat에 대해 다음 값을 지원합니다.
kinesiscsvtvsjsonavro_streamprotobufthrift
Tuning 구성
다음 표는 Kinesis에 특화된 tuningConfig 구성 속성을 정리한 것이에요. 모든 스트리밍 수집 방법에 공통인 속성은 Supervisor tuning configuration을 참고하세요.
| Property | Type | Description | Required | Default |
| skipSequenceNumberAvailabilityCheck | Boolean | 특정 Kinesis shard에서 현재 시퀀스 번호가 여전히 사용 가능한지 확인을 활성화할지 여부. false면 indexing task는 resetOffsetAutomatically 값에 따라 현재 시퀀스 번호 재설정을 시도. resetOffsetAutomatically 속성에 대한 자세한 내용은 Supervisor tuning configuration 참고 | No | false |
| recordBufferSizeBytes | Integer | Druid가 Kinesis fetch 스레드와 메인 수집 스레드 사이에서 사용하는 버퍼의 크기(힙 메모리 바이트) | No | 기본값은 Determine fetch settings 참고 |
| recordBufferOfferTimeout | Integer | 타임아웃 전에 버퍼에 공간이 생기기를 기다리는 밀리초 수 | No | 5000 |
| recordBufferFullWait | Integer | Druid가 Kinesis에서 레코드를 다시 가져오기 전에 버퍼가 비워지기를 기다리는 밀리초 수 | No | 5000 |
| fetchThreads | Integer | Kinesis에서 데이터를 가져오는 스레드 풀 크기. Kinesis shard보다 많은 스레드가 있어도 이점은 없음 | No | procs * 2, 여기서 procs는 태스크가 사용할 수 있는 프로세서 수 |
| maxBytesPerPoll | Integer | 폴당 버퍼에서 가져오는 최대 바이트 수. 이 구성과 무관하게 폴당 최소 하나의 레코드가 폴링됨 | No | 1000000 bytes |
| repartitionTransitionDuration | ISO 8601 period | shard가 분할·병합될 때 supervisor는 shard와 태스크 그룹 매핑을 다시 계산함. supervisor는 또한 이전 매핑 아래 생성된 실행 중인 태스크들에 현재 시간 + repartitionTransitionDuration 시점에 조기 종료하도록 신호를 보냄. 태스크를 조기 종료하면 Druid가 새 shard에서 더 빨리 읽기 시작할 수 있음. 이 속성으로 제어되는 repartition 전환 대기 시간은 분할·병합 후 스트림이 새 shard에 레코드를 쓸 추가 시간을 주어, empty shard handling 문제를 피하는 데 도움 | No | PT2M |
| useListShards | Boolean | AWS Kinesis SDK의 listShards API를 사용해서 수집 중 LimitExceededException을 방지할 수 있는지 여부. 필요한 IAM 권한을 반드시 설정해야 함 | No | false |
AWS 인증
Druid는 AWS 액세스 키와 시크릿 키로 Kinesis API 요청을 인증해요. 이 정보를 Druid에 제공하는 방법은 몇 가지가 있습니다.
- 역할 또는 단기 자격 증명(short-term credentials) 사용: Druid는 환경 변수, Web Identity Token, 기본 프로필 구성 파일, EC2 인스턴스 프로필 제공자(이 순서대로)에서 자격 증명을 찾습니다.
- 장기 보안 자격 증명(long-term security credentials) 사용:
아래 예시처럼
common.runtime.properties파일에 AWS 액세스 키와 시크릿 키를 직접 제공할 수 있어요.
druid.kinesis.accessKey=AKIAWxxxxxxxxxx4NCKSdruid.kinesis.secretKey=Jbytxxxxxxxxxxx2+555
정보
AWS는 보안 위험이 될 수 있으므로 구성 파일에 장기 보안 자격 증명을 제공하는 것을 권장하지 않아요. 이 방식을 사용하면 다른 모든 자격 증명 제공 방법보다 우선합니다.
Kinesis에서 데이터를 수집하려면 IAM 역할에 부착된 정책에 필요한 권한이 포함되어 있는지 확인하세요. 필요한 권한은 useListShards 값에 따라 달라집니다.
useListShards 플래그가 true로 설정되면 다음 권한이 필요해요.
ListStreams— 데이터 스트림 나열Get*—GetShardIterator에 필요GetRecords— 데이터 스트림의 shard에서 데이터 레코드 가져오기ListShards— 관심 있는 스트림의 shard 가져오기
정책 예시는 다음과 같아요.
[ { "Effect": "Allow", "Action": ["kinesis:List*"], "Resource": ["*"] }, { "Effect": "Allow", "Action": ["kinesis:Get*"], "Resource": [<ARN for shards to be ingested>] }]
useListShards 플래그가 false로 설정되면 다음 권한이 필요해요.
ListStreams— 데이터 스트림 나열Get*—GetShardIterator에 필요GetRecords— 데이터 스트림의 shard에서 데이터 레코드 가져오기DescribeStream— 지정된 데이터 스트림 설명
정책 예시는 다음과 같아요.
[ { "Effect": "Allow", "Action": ["kinesis:ListStreams"], "Resource": ["*"] }, { "Effect": "Allow", "Action": ["kinesis:DescribeStream"], "Resource": ["*"] }, { "Effect": "Allow", "Action": ["kinesis:Get*"], "Resource": [<ARN for shards to be ingested>] }]
Shard와 세그먼트 hand-off
각 Kinesis indexing task는 Kinesis shard에서 소비한 이벤트를 세그먼트 granularity interval에 대해 단일 세그먼트로 쓰다가, maxRowsPerSegment, maxTotalRows, intermediateHandoffPeriod 중 하나에 도달합니다. 이 시점에 태스크는 이후 이벤트를 담기 위해 이 세그먼트 granularity에 대해 새 shard를 만듭니다.
Kinesis indexing task는 또한 증분 hand-off(값 hand-off)를 수행해서, 태스크가 만든 세그먼트가 태스크 duration이 끝날 때까지 붙잡혀 있지 않도록 해요. 태스크가 maxRowsPerSegment, maxTotalRows, intermediateHandoffPeriod 중 하나의 한도에 도달하면 모든 세그먼트를 hand-off 하고 이후 이벤트를 위해 새로운 세그먼트 집합을 만듭니다. 이를 통해 태스크가 Middle Manager 서비스에 오래된 세그먼트를 로컬로 쌓아 두지 않고도 더 긴 duration으로 실행될 수 있습니다.
Kinesis indexing service는 여전히 작은 세그먼트를 일부 만들 수 있어요. 예를 들어 다음 시나리오를 고려해 보세요.
- 태스크 duration이 4시간
- 세그먼트 granularity가 HOUR로 설정
- supervisor가 9:10에 시작
4시간 뒤인 13:10에 Druid는 새 태스크 집합을 시작합니다. 13:00 ~ 14:00 interval의 이벤트는 기존 태스크와 새 태스크 집합에 나뉘어 들어가면서 작은 세그먼트가 생길 수 있어요. 이들을 이상적인 크기(세그먼트당 약 500~700 MB 범위)의 새 세그먼트로 병합하려면, 선택적으로 다른 세그먼트 granularity로 re-indexing 태스크를 예약할 수 있습니다.
세그먼트 크기 최적화 방법은 Segment size optimization을 참고하세요.
Fetch 설정 결정 (Determine fetch settings)
Kinesis indexing task는 fetchThreads 스레드로 레코드를 가져와요. fetchThreads가 Kinesis shard 수보다 많으면 초과 스레드는 사용되지 않습니다. 각 fetch 스레드는 Kinesis shard에서 한 번에 최대 10 MB의 레코드를 가져오며, fetch 사이에 fetchDelayMillis만큼의 지연이 있어요. 각 스레드가 가져온 레코드는 크기 recordBufferSizeBytes의 공유 큐로 들어갑니다.
이 파라미터들의 기본값은 다음과 같아요.
fetchThreads: 태스크가 사용할 수 있는 프로세서 수의 두 배. 태스크가 사용할 수 있는 프로세서 수는 서버의 총 프로세서 수를druid.worker.capacity(해당 서버의 태스크 슬롯 수)로 나눈 값. 또한 각 스레드가 한 번에 10 MB의 레코드를 가져온다고 가정할 때, 특정 시점에 가져온 총 데이터 레코드가 구성된 최대 힙의 5%를 초과하지 않도록 이 값은 더 제한됨. 이 구성에 지정된 값이 이 한도보다 높아도 실패는 발생하지 않지만 경고가 로깅되고, 값은 이 제약이 허용하는 최대치로 암묵적으로 낮춰짐.fetchDelayMillis: 0 (fetch 사이 지연 없음).recordBufferSizeBytes: 100 MB 또는 사용 가능한 힙의 대략 10% 중 더 작은 값.maxBytesPerPoll: 1000000.
Kinesis는 레코드 fetch 호출에 다음 제한을 둡니다.
- 각 데이터 레코드는 최대 1 MB 크기
- 각 shard는 읽기에 대해 초당 최대 5 트랜잭션 지원
- 각 shard는 초당 최대 2 MB 읽기 가능
- GetRecords가 반환할 수 있는 데이터 최대 크기는 10 MB
위 한도를 초과하면 Kinesis는 ProvisionedThroughputExceededException 오류를 던집니다. 이런 일이 발생하면 Druid Kinesis task는 fetchDelayMillis 또는 3초 중 더 큰 값만큼 멈췄다가 다시 호출을 시도합니다.
대부분의 경우 fetch 파라미터의 기본 설정은 과도한 메모리 사용 없이 좋은 성능을 내기에 충분해요. 다만 어떤 경우에는 fetch 속도와 메모리 사용을 더 세밀하게 제어하기 위해 이 파라미터들을 조정해야 할 수도 있습니다. 최적 값은 레코드의 평균 크기와 특정 shard에서 읽는 소비자 수(replicas 값, Kinesis 스트림을 다른 소비자가 함께 읽고 있지 않다면)에 따라 달라집니다.
역집계 (Deaggregation)
Kinesis indexing service는 단일 Kinesis Data Streams 레코드 안에 저장된 여러 행의 역집계(de-aggregation)를 지원해서 더 효율적인 데이터 전송을 가능하게 해요.
재샤딩 (Resharding)
재샤딩(resharding)은 스트림의 shard 수를 조정해서 스트림을 흐르는 데이터 속도 변화에 적응시키는 고급 작업이에요.
Kinesis 스트림의 shard 수를 변경할 때는 Kinesis 수집 태스크의 조기 종료와 가능한 태스크 실패가 있는 재샤딩 작업 주변의 시간 창이 있습니다.
조기 종료와 태스크 실패는 예상된 동작이에요. shard가 닫히고 완전히 읽히면 supervisor가 shard와 태스크 그룹 매핑을 갱신하기 때문에 발생합니다. 이렇게 하면 완전히 읽힌 닫힌 shard 할당으로 태스크가 실행되지 않고, 활성 shard의 분포가 태스크들 사이에 균형 있게 유지됩니다.
이 조기 태스크 종료·가능한 실패 창은 다음 때 끝납니다.
- 모든 닫힌 shard가 완전히 읽히고 Kinesis 수집 태스크가 그 shard들의 데이터를 게시해서 "closed" 상태를 메타데이터 스토리지에 커밋
- 할당에 비활성 shard가 있던 남은 태스크들이 종료됨. 이 태스크들은 닫힌 shard가 완전히 소진되기 전에 생성된 것들
참고로 supervisor가 실행 중에 새 파티션을 감지하면, 태스크는 useEarliestSequence 설정과 무관하게 새 파티션을 가장 이른 시퀀스 번호부터 읽어요. 이는 새 shard가 즉시 발견되므로 lag가 생길 가능성이 낮기 때문입니다.
supervisor가 일시 중지(suspended)된 상태에서 재샤딩이 발생하고 useEarliestSequence가 false로 설정되어 있으면, supervisor를 재개했을 때 태스크가 새 shard를 최신 시퀀스부터 읽습니다. 이는 supervisor가 일시 중지된 동안 쌓인 lag를 소비자가 빠르게 따라잡도록 하기 위한 설계예요.
알려진 이슈 (Known issues)
druid-kinesis-indexing-service 확장을 프로덕션에 배포하기 전에 다음 알려진 이슈를 고려하세요.
- Kinesis는 shard당 읽기 처리량 한도를 부과합니다. 여러 supervisor가 같은 Kinesis 스트림을 읽는다면, 모든 supervisor가 충분한 읽기 처리량을 확보하도록 shard를 더 추가하는 것을 고려하세요.
- Kinesis supervisor는 때때로 체크포인트 시퀀스 번호를 스트림의 보존 윈도우(양)와 비교해서 뒤처졌는지 확인할 수 있어요. 이런 검사는 Kinesis의 가장 이른 시퀀스 번호를 가져와서 AWS CloudWatch에서
IteratorAgeMilliseconds가 매우 높아질 수 있습니다.
더 알아보기 (Learn more)
관련 주제는 다음을 참고하세요.
- Supervisor API — API로 supervisor를 관리·모니터링하는 방법
- Supervisor — supervisor 상태와 용량 계획