실시간 이벤트 처리를 위한 Kinesis Data Streams
Kinesis Data Streams로 높은 처리량의 이벤트 스트림을 생성하고 소비하며, 처리량을 위해 샤드를 관리하고 Lambda를 소비자로 사용합니다.
실시간 이벤트 처리를 위한 Kinesis Data Streams은(는) CoddyKit의 무료 AWS Solutions Architect 강의입니다. 이것은 4개 중 3번째 강의입니다. 아래에서 전체 강의를 무료로 읽을 수 있으며, 내장 코드 에디터와 24/7 AI 튜터와 함께 브라우저에서 직접 실습할 수 있습니다. 이 강의는 AWS Solutions Architect 학습 경로의 일부이며, 진행 상황이 웹과 CoddyKit 앱에 동기화됩니다. AWS Solutions Architect 강의에는 총 4개의 강의가 포함되어 있습니다.
Kinesis Data Streams 핵심 개념
Kinesis Data Streams (KDS)는 내구성이 높고 순서가 보장되는 실시간 데이터 스트리밍 서비스입니다. 데이터는 하나 이상의 샤드로 구성된 스트림에 정리됩니다. 각 샤드는 순서가 있는 데이터 레코드의 시퀀스입니다. 생산자는 레코드를 어느 샤드에 저장할지 결정하는 분할 키를 사용하여 샤드에 레코드를 저장합니다. 소비자는 샤드에서 레코드를 읽고, 각 샤드 안에서 레코드가 도착한 순서대로 처리합니다.
샤드 용량 및 처리량 제한
각 샤드는 초당 1MB 또는 초당 1,000개 레코드의 쓰기 처리량과 초당 2MB의 읽기 처리량을 지원합니다(해당 샤드의 모든 표준 소비자가 공유). 전체 스트림 용량은 샤드 수에 따라 선형적으로 확장됩니다. 다음 공식을 사용하세요: shards_needed = max(write_MB_per_s / 1, read_MB_per_s / 2). 생산자가 쓰기 제한에 도달하면 ProvisionedThroughputExceededException 오류가 나타납니다. 샤드를 분할하거나 분할 키를 더 고르게 분산하여 해결하세요.
# Calculate shards needed for a stream:
# - Ingest rate: 5 MB/s writes
# - Read rate: 3 consumers x 2 MB/s = 6 MB/s reads
# shards = max(5/1, 6/2) = max(5, 3) = 5 shards needed
aws kinesis create-stream \
--stream-name iot-telemetry \
--shard-count 5분할 키 및 데이터 분산
분할 키는 Kinesis가 해시(MD5)하여 레코드를 어느 샤드에 저장할지 결정하는 문자열입니다. 적절하게 선택한 분할 키는 레코드를 샤드 전체에 고르게 분산합니다(핫 샤드 방지). IoT에서는 장치 ID를 사용하세요. 클릭스트림에서는 세션 ID 또는 사용자 ID를 사용하세요. 카디널리티가 낮은 키(예: 값이 5개뿐인 국가 이름)는 피하세요. 이러한 키를 사용하면 한 샤드에 불균형적으로 많은 쓰기 트래픽이 몰리고 다른 샤드는 유휴 상태가 되는 핫 샤드가 발생합니다.
import boto3, json
client = boto3.client('kinesis', region_name='us-east-1')
# Good: use device_id as partition key for even distribution
event = {'deviceId': 'sensor-42', 'temp': 23.5, 'ts': '2024-01-15T10:00:00Z'}
client.put_record(
StreamName='iot-telemetry',
Data=json.dumps(event),
PartitionKey='sensor-42' # high-cardinality -> even distribution
)표준 소비자와 Enhanced Fan-Out
표준 소비자는 폴링 방식의 GetRecords를 사용하여 샤드당 초당 2MB의 읽기 처리량을 공유합니다. 한 샤드에 각각 초당 2MB가 필요한 소비자가 3개 있다면 총 초당 2MB를 나누어 사용하므로 처리량 제한이 발생합니다. Enhanced Fan-Out (EFO)은 등록된 각 소비자에게 지속적인 HTTP/2 푸시 연결(SubscribeToShard)을 통해 전용 초당 2MB 읽기 파이프를 제공합니다. EFO에는 소비자별·샤드별 시간당 비용이 추가되지만 읽기 경합을 완전히 제거합니다.
# Register an Enhanced Fan-Out consumer
aws kinesis register-stream-consumer \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--consumer-name real-time-analytics
# List registered consumers
aws kinesis list-stream-consumers \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetryKinesis 소비자로서의 Lambda
Lambda는 이벤트 소스 매핑을 통해 Kinesis Data Streams와 기본적으로 통합됩니다. Lambda는 스트림을 폴링하고 레코드 배치를 읽은 다음 함수를 호출합니다. BatchSize(1~10,000개 레코드), StartingPosition(가장 오래된 레코드에는 TRIM_HORIZON, 가장 최신 레코드에는 LATEST), 그리고 실패한 배치를 분할하는 BisectBatchOnFunctionError를 구성하세요. 병렬화 계수(1~10)를 사용하면 Lambda가 샤드마다 여러 동시 호출을 실행하여 빠르게 이동하는 스트림의 처리 속도를 따라갈 수 있습니다.
# Create a Lambda event source mapping for Kinesis
aws lambda create-event-source-mapping \
--function-name ProcessIoTEvents \
--event-source-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--batch-size 100 \
--starting-position LATEST \
--parallelization-factor 5 \
--bisect-batch-on-function-error true \
--destination-config '{
"OnFailure": {
"Destination": "arn:aws:sqs:us-east-1:123456789012:kinesis-dlq"
}
}'Kinesis Client Library (KCL)
Kinesis Client Library (KCL)는 견고한 Java(또는 MultiLangDaemon을 통한 다중 언어) Kinesis 소비자를 구축하기 위한 애플리케이션 프레임워크입니다. KCL은 샤드 열거, 리스 관리(작업자 인스턴스 간 샤드 분배), 진행 상황의 체크포인트 저장, 그리고 원활한 샤드 분할 및 병합을 처리합니다. 각 KCL 작업자는 하나 이상의 샤드를 처리하며, 작업자가 확장되거나 실패하면 KCL이 샤드를 자동으로 재분배합니다. 프로덕션 소비자 애플리케이션에서는 원시 SDK 폴링보다 KCL을 사용하는 것이 좋습니다.
# KCL stores checkpoints in a DynamoDB table automatically
# Each shard has one row tracking the last successfully processed sequence number
# KCL lease table structure:
# leaseKey (shardId) | checkpoint (sequenceNumber) | leaseOwner (workerId)
#
# To start a KCL application (pseudocode):
# KinesisClientLibConfiguration config = new KinesisClientLibConfiguration(
# 'iot-app', 'iot-telemetry', credentialsProvider, 'worker-1');
# Worker worker = new Worker.Builder().config(config).recordProcessorFactory(factory).build();
# worker.run();데이터 보존 및 재생
Kinesis Data Streams는 기본적으로 레코드를 24시간 동안 저장합니다(추가 비용을 지불하면 7일 또는 장기 보존을 통해 최대 365일까지 연장 가능). SQS와 달리 소비된 레코드는 소비 후 삭제되지 않습니다. 보존 기간이 만료될 때까지 계속 사용할 수 있습니다. 따라서 여러 소비자가 동일한 레코드를 서로 독립적으로 읽을 수 있으며, 소비자의 체크포인트를 이전 시퀀스 번호로 재설정하여 재생할 수도 있습니다. 이는 버그 수정이나 새 서비스로 데이터를 보충할 때 매우 유용합니다.
# Extend stream retention to 7 days
aws kinesis increase-stream-retention-period \
--stream-name iot-telemetry \
--retention-period-hours 168
# Get records from the oldest available record (replay)
SHARD_ITERATOR=$(aws kinesis get-shard-iterator \
--stream-name iot-telemetry \
--shard-id shardId-000000000000 \
--shard-iterator-type TRIM_HORIZON \
--query 'ShardIterator' --output text)
aws kinesis get-records --shard-iterator $SHARD_ITERATOR --limit 100온디맨드 용량 모드
Kinesis Data Streams는 두 가지 용량 모드를 지원합니다. 프로비저닝 모드: 샤드 수를 직접 관리하고 샤드 시간당 비용을 지불합니다. On-Demand 모드: Kinesis가 수신 처리량에 따라 샤드 용량을 자동으로 조정하며(기본적으로 쓰기 최대 초당 200MB, 읽기 최대 초당 400MB), 기록 및 검색한 데이터의 GB당 비용을 지불합니다. 샤드 확장을 직접 관리하고 싶지 않은 가변적이거나 예측하기 어려운 트래픽 패턴에는 온디맨드 모드가 적합합니다.
# Switch an existing stream to On-Demand mode
aws kinesis update-stream-mode \
--stream-arn arn:aws:kinesis:us-east-1:123456789012:stream/iot-telemetry \
--stream-mode-details StreamMode=ON_DEMAND샤드 내 순서 보장
Kinesis는 샤드 내 순서를 보장합니다. 동일한 분할 키를 가진 레코드는 항상 같은 샤드로 전송되며 기록된 순서대로 읽힙니다. 그러나 샤드 간 순서는 보장되지 않습니다. 모든 레코드에 대한 전역 순서가 애플리케이션에 필요하다면 단일 샤드를 사용하여 처리량을 초당 1MB로 제한하거나, 장치별 순서처럼 분할 키 그룹 내에서만 순서가 필요하도록 설계를 변경하세요. 이는 엄격한 중복 제거와 순서를 제공하는 SQS FIFO와 비교할 때 시험에서 자주 나오는 구분점입니다.
보안: 암호화 및 VPC
Kinesis Data Streams는 서버 측 암호화가 활성화된 경우 AWS KMS(CMK 또는 AWS 관리형 키)를 사용하여 레코드를 서버 측에서 저장 시 암호화합니다. 전송 중인 모든 데이터는 TLS로 암호화됩니다. 퍼블릭 인터넷을 통해 스트림 데이터를 전송해서는 안 되는 VPC 실행 애플리케이션에는 Kinesis용 VPC 인터페이스 엔드포인트(PrivateLink)를 사용하세요. 그러면 트래픽이 AWS 네트워크 백본 내부에만 머물러 규제 대상 워크로드에 중요합니다.
# Enable server-side encryption on a Kinesis stream
aws kinesis start-stream-encryption \
--stream-name iot-telemetry \
--encryption-type KMS \
--key-id arn:aws:kms:us-east-1:123456789012:key/mrk-abc123Kinesis Data Streams와 SQS 및 Kafka 비교
SAA-C03 시험에서는 Kinesis Data Streams를 대안 서비스와 비교해야 합니다. Kinesis와 SQS: Kinesis는 샤드 내 순서를 유지하고 여러 소비자가 동일한 데이터를 읽을 수 있도록 지원하지만, SQS는 소비 후 메시지를 삭제합니다. Kinesis와 MSK (Kafka): MSK는 관리형 Apache Kafka입니다. Kafka 프로토콜 호환성, 고급 주제 구성 또는 온프레미스 Kafka에서의 마이그레이션이 필요할 때 사용하세요. AWS 고유의 스트리밍과 Lambda, Firehose, Flink와의 긴밀한 통합이 필요하다면 Kinesis를 사용하세요. 문제에서 Kafka 또는 Kafka 호환성 요구 사항을 명시적으로 언급하지 않는 한 Kinesis를 선택하세요.
빠른 확인
이 레슨에서 다룬 AWS Solutions Architect (SAA-C03) 개념에 대한 이해도를 확인해 보세요.
레슨 요약
이 레슨에서는 Kinesis Data Streams가 24시간의 기본 보존 기간과 재생 기능을 제공하는 샤드 기반의 순서가 보장되는 내구성 높은 스트림을 제공한다는 점, Enhanced Fan-Out이 각 소비자에게 샤드별 초당 2MB의 전용 처리량을 제공하여 읽기 경합을 제거한다는 점, 그리고 On-Demand 모드가 예측하기 어려운 트래픽에 맞춰 샤드를 자동으로 확장한다는 점을 배웠습니다. 다음으로는 이벤트 기반 아키텍처에서 안무 방식과 오케스트레이션 방식의 패턴을 살펴보겠습니다.
자주 묻는 질문
“실시간 이벤트 처리를 위한 Kinesis Data Streams” 강의는 무료인가요?
네 — “실시간 이벤트 처리를 위한 Kinesis Data Streams” 전체 내용을 이 웹사이트에서 무료로 읽을 수 있습니다. 인터랙티브하게 실습하려면(내장 코드 에디터와 24/7 AI 튜터), CoddyKit PRO로 업그레이드하면 AWS Solutions Architect 강의 전체를 잠금 해제할 수 있습니다. AWS Solutions Architect 강의에는 총 4개의 강의가 포함되어 있습니다.
“실시간 이벤트 처리를 위한 Kinesis Data Streams”에서 뭘 배우나요?
Kinesis Data Streams로 높은 처리량의 이벤트 스트림을 생성하고 소비하며, 처리량을 위해 샤드를 관리하고 Lambda를 소비자로 사용합니다. 브라우저에서 직접 실행하는 실습 코드로 AWS Solutions Architect을(를) 배우며, 24/7 AI 튜터가 강의를 진행하면서 질문에 답변해줍니다.
AWS Solutions Architect을(를) 시작하는 데 경험이 필요한가요?
사전 경험은 필요하지 않습니다. CoddyKit의 AWS Solutions Architect은(는) 초급자부터 고급 학습자까지를 위해 구성되어 있으므로, 여기서 시작하거나 처음부터 시작할 수 있으며 자신의 속도대로 진행할 수 있습니다. 이것은 4개 중 3번째 강의입니다.
“실시간 이벤트 처리를 위한 Kinesis Data Streams” 강의는 얼마나 걸리나요?
대부분의 CoddyKit 강의는 약 5~10분이 소요됩니다. 각 강의는 간결하고 인터랙티브하여 꾸준한 진행이 가능하며, 웹과 앱에서 중단한 부분부터 바로 시작할 수 있습니다.
이 AWS Solutions Architect 강의에서 코드를 작성하고 실행할 수 있나요?
네. 모든 AWS Solutions Architect 강의에는 내장 코드 에디터가 포함되어 있으므로, 브라우저에서 바로 실제 코드를 작성하고 실행한 후 즉시 AI 피드백을 받을 수 있습니다 — 로컬 설정이 필요 없습니다.
이 강의의 모든 강의
- EventBridge: 이벤트 버스 및 규칙
- Step Functions: 서버리스 워크플로 오케스트레이션
- 실시간 이벤트 처리를 위한 Kinesis Data Streams
- 안무 방식과 오케스트레이션 패턴 비교