Kinesis 통합
Kinesis 통합 (Streaming Kinesis Integration)
Amazon Kinesis는 대규모 스트리밍 데이터를 실시간으로 처리하기 위한 완전 관리형 서비스예요. Kinesis 리시버는 Amazon이 Amazon Software License(ASL)로 제공하는 Kinesis Client Library(KCL)를 이용해 입력 DStream을 만들어요. KCL은 Apache 2.0 라이선스의 AWS Java SDK 위에 구축되며, Worker, Checkpoint, Shard Lease 개념을 통해 로드 밸런싱, 내결함성, 체크포인팅을 제공해요. 여기서는 Spark Streaming이 Kinesis에서 데이터를 받도록 구성하는 방법을 설명할게요.
본문
Kinesis 구성하기
Kinesis 스트림은 유효한 Kinesis 엔드포인트 중 한 곳에서 이 가이드에 따라 1개 이상의 샤드로 설정할 수 있어요.
Spark Streaming 애플리케이션 구성하기
Linking: SBT/Maven 프로젝트 정의를 쓰는 Scala/Java 애플리케이션이라면, 스트리밍 애플리케이션을 다음 아티팩트에 연결해요(자세한 내용은 메인 프로그래밍 가이드의 Linking 섹션 참고).
groupId = org.apache.spark
artifactId = spark-streaming-kinesis-asl_2.13
version = 4.2.0
Python 애플리케이션이라면 애플리케이션을 배포할 때 위 라이브러리와 그 의존성을 추가해야 해요. 아래의 Deploying 하위 섹션을 보세요. 이 라이브러리에 연결하면 애플리케이션에 ASL-라이선스 코드가 포함된다는 점에 유의하세요.
Programming: 스트리밍 애플리케이션 코드에서 KinesisInputDStream을 import하고 다음과 같이 byte 배열의 입력 DStream을 만들어요.
from pyspark.streaming.kinesis import KinesisUtils, InitialPositionInStream
kinesisStream = KinesisUtils.createStream(
streamingContext, [Kinesis app name], [Kinesis stream name], [endpoint URL],
[region name], [initial position], [checkpoint interval], [metricsLevel.DETAILED],
StorageLevel.MEMORY_AND_DISK_2)
API 문서와 예제를 보세요. 예제 실행 방법은 Running the Example 하위 섹션을 참고해요.
- CloudWatch 지표 수준과 차원. 자세한 내용은 KCL 모니터링에 관한 AWS 문서를 보세요. 기본값은
MetricsLevel.DETAILED예요.
import org.apache.spark.storage.StorageLevel
import org.apache.spark.streaming.kinesis.KinesisInputDStream
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.kinesis.KinesisInitialPositions
val kinesisStream = KinesisInputDStream.builder
.streamingContext(streamingContext)
.endpointUrl([endpoint URL])
.regionName([region name])
.streamName([streamName])
.initialPosition([initial position])
.checkpointAppName([Kinesis app name])
.checkpointInterval([checkpoint interval])
.metricsLevel([metricsLevel.DETAILED])
.storageLevel(StorageLevel.MEMORY_AND_DISK_2)
.build()
API 문서와 예제를 보세요. 예제 실행 방법은 Running the Example 하위 섹션을 참고해요.
import org.apache.spark.storage.StorageLevel;
import org.apache.spark.streaming.kinesis.KinesisInputDStream;
import org.apache.spark.streaming.Seconds;
import org.apache.spark.streaming.StreamingContext;
import org.apache.spark.streaming.kinesis.KinesisInitialPositions;
KinesisInputDStream<byte[]> kinesisStream = KinesisInputDStream.builder()
.streamingContext(streamingContext)
.endpointUrl([endpoint URL])
.regionName([region name])
.streamName([streamName])
.initialPosition([initial position])
.checkpointAppName([Kinesis app name])
.checkpointInterval([checkpoint interval])
.metricsLevel([metricsLevel.DETAILED])
.storageLevel(StorageLevel.MEMORY_AND_DISK_2)
.build();
API 문서와 예제를 보세요. 예제 실행 방법은 Running the Example 하위 섹션을 참고해요.
다음 설정도 제공할 수 있어요. 이 기능은 현재 Scala와 Java에서만 지원돼요.
- Kinesis
KinesisClientRecord를 받아 제네릭 객체T를 반환하는 "메시지 핸들러 함수". 파티션 키처럼Record에 포함된 다른 데이터를 쓰고 싶을 때 유용해요.
import org.apache.spark.storage.StorageLevel
import org.apache.spark.streaming.kinesis.KinesisInputDStream
import org.apache.spark.streaming.{Seconds, StreamingContext}
import org.apache.spark.streaming.kinesis.KinesisInitialPositions
import software.amazon.kinesis.metrics.{MetricsLevel, MetricsUtil}
val kinesisStream = KinesisInputDStream.builder
.streamingContext(streamingContext)
.endpointUrl([endpoint URL])
.regionName([region name])
.streamName([streamName])
.initialPosition([initial position])
.checkpointAppName([Kinesis app name])
.checkpointInterval([checkpoint interval])
.storageLevel(StorageLevel.MEMORY_AND_DISK_2)
.metricsLevel(MetricsLevel.DETAILED)
.metricsEnabledDimensions(
Set(MetricsUtil.OPERATION_DIMENSION_NAME, MetricsUtil.SHARD_ID_DIMENSION_NAME))
.buildWithMessageHandler([message handler])
import java.util.Set;
import scala.jdk.javaapi.CollectionConverters;
import org.apache.spark.storage.StorageLevel;
import org.apache.spark.streaming.kinesis.KinesisInputDStream;
import org.apache.spark.streaming.Seconds;
import org.apache.spark.streaming.StreamingContext;
import org.apache.spark.streaming.kinesis.KinesisInitialPositions;
import software.amazon.kinesis.metrics.MetricsLevel;
import software.amazon.kinesis.metrics.MetricsUtil;
KinesisInputDStream<byte[]> kinesisStream = KinesisInputDStream.builder()
.streamingContext(streamingContext)
.endpointUrl([endpoint URL])
.regionName([region name])
.streamName([streamName])
.initialPosition([initial position])
.checkpointAppName([Kinesis app name])
.checkpointInterval([checkpoint interval])
.storageLevel(StorageLevel.MEMORY_AND_DISK_2)
.metricsLevel(MetricsLevel.DETAILED)
.metricsEnabledDimensions(
CollectionConverters.asScala(
Set.of(
MetricsUtil.OPERATION_DIMENSION_NAME,
MetricsUtil.SHARD_ID_DIMENSION_NAME)).toSet())
.buildWithMessageHandler([message handler]);
streamingContext: 이 Kinesis 애플리케이션을 Kinesis 스트림에 연결하는 데 쓰는 애플리케이션 이름을 담은 StreamingContext[Kinesis app name]: DynamoDB 테이블에 Kinesis 시퀀스 번호를 체크포인트하는 데 쓰는 애플리케이션 이름- 애플리케이션 이름은 주어진 계정과 리전에서 고유해야 해요.
- 테이블이 존재하는데 체크포인트 정보가 잘못됐다면(다른 스트림이거나 오래된 만료 시퀀스 번호), 일시적인 오류가 생길 수 있어요.
[Kinesis stream name]: 이 스트리밍 애플리케이션이 데이터를 가져올 Kinesis 스트림[endpoint URL]: 유효한 Kinesis 엔드포인트 URL은 여기에서 찾을 수 있어요.[region name]: 유효한 Kinesis 리전 이름은 여기에서 찾을 수 있어요.[checkpoint interval]: Kinesis Client Library가 스트림에서 자신의 위치를 저장하는 간격(예:Duration(2000)= 2초). 처음에는 스트리밍 애플리케이션의 배치 간격과 같게 설정해요.[initial position]:KinesisInitialPositions.TrimHorizon,KinesisInitialPositions.Latest,KinesisInitialPositions.AtTimestamp중 하나(자세한 내용은Kinesis Checkpointing섹션과Amazon Kinesis API 문서참고).[message handler]: KinesisKinesisClientRecord를 받아 제네릭T를 출력하는 함수
API의 다른 버전에서는 AWS 액세스 키와 시크릿 키를 직접 지정할 수도 있어요.
Deploying: 다른 Spark 애플리케이션과 마찬가지로 spark-submit으로 애플리케이션을 실행해요. 다만 Scala/Java 애플리케이션과 Python 애플리케이션은 세부 사항이 조금 달라요.
Scala와 Java 애플리케이션의 경우, 프로젝트 관리에 SBT나 Maven을 쓴다면 spark-streaming-kinesis-asl_2.13과 그 의존성을 애플리케이션 JAR에 패키징해요. spark-core_2.13과 spark-streaming_2.13은 Spark 설치에 이미 있으므로 provided 의존성으로 표시해요. 그런 다음 spark-submit으로 애플리케이션을 실행해요(메인 프로그래밍 가이드의 Deploying 섹션 참고).
SBT/Maven 프로젝트 관리가 없는 Python 애플리케이션의 경우, spark-streaming-kinesis-asl_2.13과 그 의존성을 --packages로 spark-submit에 직접 추가할 수 있어요(Application Submission Guide 참고). 즉,
./bin/spark-submit --packages org.apache.spark:spark-streaming-kinesis-asl_2.13:4.2.0 ...
또는 Maven 아티팩트 spark-streaming-kinesis-asl-assembly의 JAR을 Maven 저장소에서 다운로드해 --jars로 spark-submit에 추가할 수도 있어요.
런타임에 기억할 점:
- Kinesis 데이터 처리는 파티션별로 순서가 보장되고 메시지마다 최소 한 번(at-least once) 처리돼요.
- 여러 애플리케이션이 같은 Kinesis 스트림에서 읽을 수 있어요. Kinesis는 애플리케이션별 샤드와 체크포인트 정보를 DynamoDB에 유지해요.
- 하나의 Kinesis 스트림 샤드는 한 번에 하나의 입력 DStream으로 처리돼요.
- 하나의 Kinesis 입력 DStream은 여러
KinesisRecordProcessor스레드를 만들어 Kinesis 스트림의 여러 샤드에서 읽을 수 있어요. - 별도의 프로세스/인스턴스에서 실행되는 여러 입력 DStream이 Kinesis 스트림에서 읽을 수 있어요.
- 각 입력 DStream은 한 샤드를 처리하는
KinesisRecordProcessor스레드를 최소 하나 만드므로, Kinesis 스트림 샤드 수보다 많은 Kinesis 입력 DStream이 필요하지 않아요. - 수평 확장은 Kinesis 입력 DStream을 추가/제거해서(단일 프로세스 내 또는 여러 프로세스/인스턴스에 걸쳐) 이루어져요 — 앞 항목에 따라 Kinesis 스트림 샤드 수까지 가능해요.
- Kinesis 입력 DStream은 모든 DStream 사이에서 로드를 균형 있게 분배해요 — 프로세스/인스턴스 간에도 그렇죠.
- Kinesis 입력 DStream은 로드 변화 때문에 생기는 re-shard 이벤트(병합과 분할) 동안에도 로드를 균형 있게 분배해요.
- 모범 사례로, 가능하면 과잉 프로비저닝(over-provisioning)으로 re-shard 지터를 피하는 게 좋아요.
- 각 Kinesis 입력 DStream은 자신의 체크포인트 정보를 유지해요. 자세한 내용은 Kinesis Checkpointing 섹션을 보세요.
- Kinesis 스트림 샤드 수와 입력 DStream 처리 중 Spark 클러스터에 만들어지는 RDD 파티션/샤드 수 사이에는 상관관계가 없어요. 이 둘은 독립적인 두 개의 파티셔닝 체계예요.
예제 실행하기
예제를 실행하려면,
- 다운로드 사이트에서 Spark 바이너리를 내려받아요.
- AWS에서 Kinesis 스트림을 설정해요(앞 섹션 참고). Kinesis 스트림 이름과 스트림이 만들어진 리전에 해당하는 엔드포인트 URL을 적어두세요.
- AWS 자격 증명으로 환경 변수
AWS_ACCESS_KEY_ID와AWS_SECRET_ACCESS_KEY를 설정해요. - Spark 루트 디렉터리에서 예제를 다음과 같이 실행해요.
./bin/spark-submit --jars 'connector/kinesis-asl-assembly/target/spark-streaming-kinesis-asl-assembly_*.jar' \
connector/kinesis-asl/src/main/python/examples/streaming/kinesis_wordcount_asl.py \
[Kinesis app name] [Kinesis stream name] [endpoint URL] [region name]
./bin/run-example --packages org.apache.spark:spark-streaming-kinesis-asl_2.13:4.2.0 streaming.KinesisWordCountASL [Kinesis app name] [Kinesis stream name] [endpoint URL]
./bin/run-example --packages org.apache.spark:spark-streaming-kinesis-asl_2.13:4.2.0 streaming.JavaKinesisWordCountASL [Kinesis app name] [Kinesis stream name] [endpoint URL]
이렇게 하면 Kinesis 스트림에서 데이터를 받을 때까지 기다려요.
- Kinesis 스트림에 랜덤 문자열 데이터를 넣으려면, 다른 터미널에서 관련 Kinesis 데이터 프로듀서를 실행해요.
./bin/run-example streaming.KinesisWordProducerASL [Kinesis stream name] [endpoint URL] 1000 10
이렇게 하면 줄당 랜덤 숫자 10개로 된 줄을 초당 1000개씩 Kinesis 스트림에 넣어요. 그러면 실행 중인 예제가 이 데이터를 받아 처리하게 돼요.
레코드 역집계(Record De-aggregation)
데이터가 Kinesis Producer Library(KPL)로 생성되면 비용 절감을 위해 메시지가 집계(aggregated)될 수 있어요. Spark Streaming은 소비 중에 레코드를 자동으로 역집계해요.
Kinesis 체크포인팅
- 각 Kinesis 입력 DStream은 주기적으로 스트림의 현재 위치를 백킹 DynamoDB 테이블에 저장해요. 이렇게 하면 시스템이 실패에서 복구해 DStream이 중단된 지점부터 처리를 계속할 수 있어요.
- 너무 자주 체크포인트하면 AWS 체크포인트 저장 계층에 과부하가 걸리고 AWS 스로틀링으로 이어질 수 있어요. 제공된 예제는 랜덤 백오프 재시도 전략으로 이 스로틀링을 처리해요.
- 입력 DStream이 시작될 때 Kinesis 체크포인트 정보가 없으면, 가장 오래된 기록(
KinesisInitialPositions.TrimHorizon), 최신 끝(KinesisInitialPositions.Latest), 또는 (Python 제외) 주어진 UTC 타임스탬프가 가리키는 위치(KinesisInitialPositions.AtTimestamp(Date timestamp)) 중 하나에서 시작해요. 이건 설정 가능해요. KinesisInitialPositions.Latest는 입력 DStream이 실행되지 않는 동안(그리고 체크포인트 정보가 저장되지 않는 동안) 데이터가 스트림에 추가되면 레코드를 놓칠 수 있어요.KinesisInitialPositions.TrimHorizon은 레코드를 중복 처리하게 할 수 있는데, 그 영향은 체크포인트 빈도와 처리 멱등성에 따라 달라져요.
Kinesis 재시도 구성
spark.streaming.kinesis.retry.waitTime: 지속 시간 문자열로 된 Kinesis 재시도 사이의 대기 시간. Amazon Kinesis에서 읽을 때 초당 5건 이상 소비하거나 최대 읽기 속도인 초당 2MiB를 초과하면ProvisionedThroughputExceededException이 발생할 수 있어요. 이 구성은 가져오기가 실패했을 때 가져오기 사이의 대기 시간을 늘려 그런 예외를 줄이는 데 쓸 수 있어요. 기본값은 "100ms"예요.spark.streaming.kinesis.retry.maxAttempts: Kinesis 가져오기의 최대 재시도 횟수. 위에서 언급한 시나리오에서 KinesisProvisionedThroughputExceededException에 대처하는 데도 쓸 수 있어요. Kinesis 읽기의 재시도 횟수를 늘리려면 이 값을 높여요. 기본값은 3이에요.
더 알아보기 (Learn more)
- Spark Streaming 프로그래밍 가이드 — DStream과 스트리밍 애플리케이션의 전반적 내용.
- Amazon Kinesis — AWS 공식 서비스 문서.