Amazon Kinesis Data Streams 커넥터

Amazon Kinesis Data Streams 커넥터 (Amazon Kinesis Data Streams Connector)

Kinesis 커넥터를 사용하면 Amazon Kinesis Data Streams에서 읽고 쓸 수 있어요.

출처: 문서

본문

의존성 (Dependency)

이 커넥터를 사용하려면 프로젝트에 아래 의존성을 추가해요. Flink 2.3 버전에는 아직 사용할 수 있는 커넥터가 없어요. PyFlink 작업에 사용하려면 다음 의존성을 사용해요. PyFlink 작업에서 JAR을 사용하는 방법은 Python dependency management를 참고해요.

flink-connector-aws-kinesis-streams Flink 2.3 버전에는 아직 SQL jar가 없음

Kinesis Streams Source

KinesisStreamsSource는 FLIP-27 소스 인터페이스에 기반한 exactly-once, 병렬 스트리밍 데이터 소스예요. 소스는 단일 Amazon Kinesis Data 스트림을 구독하고 이벤트를 읽으며, 특정 Kinesis partitionId 내에서 순서를 유지해요(Record ordering 참고). KinesisStreamsSource는 스트림의 샤드를 발견하고 연산자의 병렬도에 따라 각 적격 샤드에서 병렬로 읽기 시작해요. 적절한 병렬도 선택에 대한 자세한 내용은 Parallelism and Number of Shards를 참고해요. 또한 작업 실행 중 스트림의 재샤딩이 발생하면 Kinesis Data 스트림의 새 샤드를 투명하게 발견해 처리해요. 자세한 내용은 Shard Discovery 섹션을 참고해요.

참고: 데이터를 소비하기 전에 Kinesis Data Stream이 Amazon Kinesis Data Streams 콘솔에서 ACTIVE 상태로 생성되었는지 확인해요.

KinesisStreamsSource는 KinesisStreamsSource 인스턴스를 구성하기 위한 플루언트 빌더를 제공해요. 아래 코드 조각은 그 방법을 보여줘요.

Java:

// Configure the KinesisStreamsSource
Configuration sourceConfig = new Configuration();
sourceConfig.set(KinesisSourceConfigOptions.STREAM_INITIAL_POSITION, KinesisSourceConfigOptions.InitialPosition.TRIM_HORIZON); // This is optional, by default connector will read from LATEST

// Create a new KinesisStreamsSource to read from specified Kinesis Stream.
KinesisStreamsSource<String> kdsSource =
KinesisStreamsSource.<String>builder()
.setStreamArn("arn:aws:kinesis:us-east-1:123456789012:stream/test-stream")
.setSourceConfig(sourceConfig)
.setDeserializationSchema(new SimpleStringSchema())
.setKinesisShardAssigner(ShardAssignerFactory.uniformShardAssigner()) // This is optional, by default uniformShardAssigner will be used.
.build();

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

// Specify watermarking strategy and the name of the Kinesis Source operator.
// Specify return type using TypeInformation.
// Specify UID of operator in line with Flink best practice.
DataStream<String> kinesisRecordsWithEventTimeWatermarks = env.fromSource(kdsSource, WatermarkStrategy.<String>forMonotonousTimestamps().withIdleness(Duration.ofSeconds(1)), "Kinesis source")
.returns(TypeInformation.of(String.class))
.uid("custom-uid");

Scala:

val sourceConfig = new Configuration()
sourceConfig.set(KinesisSourceConfigOptions.STREAM_INITIAL_POSITION, KinesisSourceConfigOptions.InitialPosition.TRIM_HORIZON) // This is optional, by default connector will read from LATEST

val env = StreamExecutionEnvironment.getExecutionEnvironment()

val kdsSource = KinesisStreamsSource.builder[String]()
.setStreamArn("arn:aws:kinesis:us-east-1:123456789012:stream/test-stream")
.setSourceConfig(sourceConfig)
.setDeserializationSchema(new SimpleStringSchema())
.setKinesisShardAssigner(ShardAssignerFactory.uniformShardAssigner()) // This is optional, by default uniformShardAssigner will be used.
.build()

val kinesisEvents = env.fromSource(kdsSource, WatermarkStrategy.forMonotonousTimestamps().withIdleness(Duration.ofSeconds(1)), "Kinesis source")
.uid("custom-uid")

위는 KinesisStreamsSource를 사용하는 간단한 예제예요.

  • 읽는 Kinesis 스트림은 Kinesis Stream ARN으로 지정돼요.
  • Source 구성은 Flink의 Configuration 클래스 인스턴스로 제공돼요. 구성 키는 AWSConfigOptions(AWS 특정 구성)와 KinesisSourceConfigOptions(Kinesis Source 구성)에서 가져올 수 있어요.
  • 예제는 시작 위치를 TRIM_HORIZON으로 지정해요(Configuring Starting Position 참고).
  • 역직렬화 포맷은 SimpleStringSchema예요(Deserialization Schema 참고).
  • 하위 작업 간 샤드 분배는 UniformShardAssigner로 제어돼요(Shard Assignment Strategy 참고).
  • 예제는 또한 증가하는 WatermarkStrategy를 지정하는데, 이는 각 레코드가 approximateArrivalTimestamp로 지정된 이벤트 시간으로 태그됨을 의미해요. 단조 증가 워터마크가 생성되며, 1초 후 레코드가 발행되지 않으면 하위 작업이 idle로 간주돼요.

IAM으로 Kinesis 접근 구성하기 (Configuring Access to Kinesis with IAM)

Kinesis 스트림 접근은 IAM identity로 제어돼요. Kinesis 스트림에서/으로 읽기·쓰기를 허용하는 적절한 IAM 정책을 만들어야 해요. 예제는 여기를 참고해요. 배포에 따라 적절한 AWS Credentials Provider를 선택할 수 있어요. 기본적으로 AUTO Credentials Provider가 사용돼요. access key ID와 secret key가 구성에 설정되면 BASIC provider가 사용돼요. 특정 Credentials Provider는 AWSConfigConstants.AWS_CREDENTIALS_PROVIDER 설정으로 선택적으로 설정할 수 있어요.

지원되는 Credential Provider:

  • AUTO(기본값): 다음 순서로 자격 증명을 찾는 기본 AWS Credentials Provider 체인 사용: ENV_VARS, SYS_PROPS, WEB_IDENTITY_TOKEN, PROFILE, EC2/ECS credentials provider.
  • BASIC: 구성으로 제공된 access key ID와 secret key 사용.
  • ENV_VAR: AWS_ACCESS_KEY_ID & AWS_SECRET_ACCESS_KEY 환경 변수 사용.
  • SYS_PROP: Java 시스템 프로퍼티 aws.accessKeyIdaws.secretKey 사용.
  • CUSTOM: 사용자 지정 클래스를 credential provider로 사용.
  • PROFILE: AWS credentials 프로파일 파일로 AWS 자격 증명 생성.
  • ASSUME_ROLE: 역할을 가정해(assume) AWS 자격 증명 생성. 역할을 가정하기 위한 자격 증명을 제공해야 해요.
  • WEB_IDENTITY_TOKEN: Web Identity Token으로 역할을 가정해 AWS 자격 증명 생성.

시작 위치 구성하기 (Configuring Starting Position)

KinesisStreamsSource가 Kinesis 스트림의 어디서 읽기 시작할지 지정하려면 구성에서 KinesisSourceConfigOptions.STREAM_INITIAL_POSITION을 설정할 수 있어요. 값(사용법)은 AWS Kinesis Data Streams 서비스의 명명을 따릅니다:

  • LATEST(기본값): 최신 레코드부터 스트림의 모든 샤드를 읽음.
  • TRIM_HORIZON: 가능한 가장 이른 레코드부터 스트림의 모든 샤드를 읽음(데이터는 Kinesis가 보존 설정에 따라 트리밍할 수 있음).
  • AT_TIMESTAMP: 지정된 타임스탬프부터 스트림의 모든 샤드를 읽음. 타임스탬프는 구성 프로퍼티에 KinesisSourceConfigOptions.STREAM_INITIAL_TIMESTAMP 값을 제공해 지정해야 하며, 다음 날짜 패턴 중 하나여야 해요:
    • Unix epoch 이후 경과한 초를 나타내는 음이 아닌 double 값(예: 1459799926.480).
    • KinesisSourceConfigOptions.STREAM_TIMESTAMP_DATE_FORMAT으로 제공된 유효한 SimpleDateFormat 패턴으로 지정된 사용자 정의 패턴. KinesisSourceConfigOptions.STREAM_TIMESTAMP_DATE_FORMAT이 정의되지 않으면 기본 패턴은 yyyy-MM-dd'T'HH:mm:ss.SSSXXX예요(예: 타임스탬프 값 2016-04-04이고 패턴이 사용자 제공 yyyy-MM-dd, 또는 패턴 없이 타임스탬프 값 2016-04-04T19:58:46.480-00:00).

구성된 시작 위치는 애플리케이션이 체크포인트 또는 savepoint에서 재시작하면 무시돼요. 자세한 내용은 Fault tolerance 섹션을 참고해요.

Exactly-Once 사용자 정의 상태 갱신 의미론을 위한 장애 허용 (Fault Tolerance for Exactly-Once User-Defined State Update Semantics)

Flink의 체크포인팅이 활성화되면 KinesisStreamsSource는 Kinesis 스트림의 샤드에서 레코드를 소비하고 각 샤드의 진행 상황을 주기적으로 체크포인트해요. 작업 실패 시 Flink는 스트리밍 프로그램을 최신 완료 체크포인트의 상태로 복원하고 체크포인트에 저장된 진행 상황부터 Kinesis 샤드의 레코드를 다시 소비해요.

체크포인트나 savepoint에서 복원할 때 구성된 시작 위치는 무시된다는 점에 주의해요. KinesisStreamsSource는 체크포인트나 savepoint에서 중단된 지점부터 읽기를 진행해요. 복원된 체크포인트나 savepoint가 오래된 경우(예: 저장된 샤드가 만료되어 Kinesis 스트림의 보존 기간을 지남) 소스는 실패하지 않고 가능한 가장 이른 이벤트부터 읽기 시작해요(사실상 TRIM_HORIZON).

기존 체크포인트나 savepoint에서 Flink 작업을 복원하되 스트림의 구성된 시작 위치를 존중하고 싶다면, KinesisStreamsSource 연산자의 uid를 바꿔 상태 없이 이 연산자를 효과적으로 복원할 수 있어요. 이는 Flink 모범 사례에 부합해요.

샤드 할당 전략 (Shard Assignment Strategy)

대부분의 사용 사례에서 사용자는 병렬 하위 작업 간 레코드의 균일한 분배를 선호해요. 이는 데이터가 Kinesis Data Stream에 고르게 분포돼 있으면 데이터 치우침을 방지해요. 이는 기본 샤드 할당 전략인 UniformShardAssigner로 달성돼요. 사용자는 KinesisShardAssigner 인터페이스를 구현해 자신만의 커스텀 전략을 구현할 수 있어요.

스트림이 재샤딩되었을 때 병렬 하위 작업 간 샤드를 균일하게 분배하는 것은 사소하지 않아요. Amazon Kinesis Data Streams는 partitionId를 주어진 스트림의 전체 HashKeyRange에 걸쳐 고르게 분배하며, UNIFORM_SCALING을 사용하면 이 범위들이 모든 열린(opened) 샤드에 고르게 분배돼요. 하지만 Kinesis Data Stream에는 Open과 Closed 샤드가 섞여 있고, 각 샤드의 상태는 재조정 연산 중에 바뀔 수 있어요.

각 병렬 하위 작업 간 partitionId의 균일한 분배를 보장하기 위해 UniformShardAssigner는 각 샤드의 HashKeyRange를 사용해 발견된 샤드를 어떤 병렬 하위 작업이 읽을지 결정해요.

레코드 순서 (Record ordering)

Kinesis는 Kinesis 스트림 내 partitionId당 레코드의 쓰기 순서를 유지해요. KinesisStreamsSource는 재샤딩 연산을 통해서도 주어진 partitionId 내에서 같은 순서로 레코드를 읽어요. 주어진 샤드를 읽기 전에 먼저 그 샤드의 부모(최대 2개 샤드)가 완전히 읽혔는지 확인하는 방식으로 이 작업을 해요.

역직렬화 스키마 (Deserialization Schema)

KinesisStreamsSource는 Kinesis Data Stream에서 이진 데이터를 가져오며, 이를 Java 객체로 변환하려면 스키마가 필요해요. Flink의 DeserializationSchema와 커스텀 KinesisDeserializationSchema 모두 KinesisStreamsSource가 받아들여요. KinesisDeserializationSchema는 레코드당 추가 Kinesis 특정 메타데이터를 제공해 사용자가 그 메타데이터를 기반으로 직렬화 결정을 내릴 수 있게 해요.

편의를 위해 Flink는 다음 스키마를 기본 제공해요:

  • SimpleStringSchema와 JsonSerializationSchema.
  • TypeInformationSerializationSchema: Flink의 TypeInformation 기반 스키마를 만들어요. 데이터를 Flink가 쓰고 읽는 경우 유용해요. 이 스키마는 다른 일반 직렬화 접근의 성능 좋은 Flink 특정 대안이에요.
  • GlueSchemaRegistryJsonDeserializationSchema: AWS Glue Schema Registry에서 writer의 스키마(레코드를 쓰는 데 사용된 스키마)를 조회하는 기능을 제공해요. 이를 사용하면 AWS Glue Schema Registry에서 검색한 스키마로 역직렬화 레코드를 읽고, 수동 제공 스키마가 있는 일반 레코드를 나타내는 com.amazonaws.services.schemaregistry.serializers.json.JsonDataWithSchema 또는 mbknor-jackson-jsonSchema로 생성된 JAVA POJO로 변환해요. 이 역직렬화 스키마를 사용하려면 추가 의존성(GlueSchemaRegistryJsonDeserializationSchema)을 추가해야 해요. Flink 2.3 버전에는 아직 사용할 수 있는 커넥터가 없어요.
  • AvroDeserializationSchema: 정적으로 제공된 스키마로 Avro 포맷 직렬화 데이터를 읽어요. Avro 생성 클래스에서 스키마를 추론할 수 있거나(AvroDeserializationSchema.forSpecific(...)) 수동 제공 스키마로 GenericRecords와 동작할 수 있어요(AvroDeserializationSchema.forGeneric(...)). Avro GenericRecords를 역직렬화할 때 역직렬화 스키마는 레코드에 임베디드 스키마가 포함되지 않는다고 기대해요.
  • AWS Glue Schema Registry를 사용해 writer의 스키마를 검색할 수 있어요. 마찬가지로 AWS Glue Schema Registry의 스키마로 역직렬화 레코드를 읽고(GlueSchemaRegistryAvroDeserializationSchema.forGeneric(...) 또는 GlueSchemaRegistryAvroDeserializationSchema.forSpecific(...) 중 하나로) 변환해요. AWS Glue Schema Registry와 Apache Flink 통합에 대한 자세한 내용은 Use Case: Amazon Kinesis Data Analytics for Apache Flink를 참고해요. 이 역직렬화 스키마를 사용하려면 추가 의존성(AvroDeserializationSchema)을 추가해야 해요.

AvroDeserializationSchema:

<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-avro</artifactId>
<version>2.3.0</version>
</dependency>

GlueSchemaRegistryAvroDeserializationSchema: Flink 2.3 버전에는 아직 사용할 수 있는 커넥터가 없어요.

병렬도와 샤드 수 (Parallelism and Number of Shards)

KinesisStreamsSource의 구성된 병렬도는 Kinesis 스트림의 총 샤드 수와 독립적이에요.

  • KinesisStreamsSource의 병렬도가 총 샤드 수보다 작으면 단일 병렬 하위 작업이 여러 샤드를 처리해요.
  • KinesisStreamsSource의 병렬도가 총 샤드 수보다 많으면 어떤 샤드도 읽지 않는 병렬 하위 작업이 있을 거예요. 이 경우 사용자는 WatermarkStrategy에 withIdleness를 설정해야 해요. 그렇게 하지 않으면 idle 하위 작업 때문에 워터마크 생성이 차단돼요.

소스의 워터마크 처리 (Watermark Handling in the source)

KinesisStreamsSource는 Kinesis가 제공하는 approximateArrivalTimestamp를 읽은 각 레코드에 연관된 이벤트 시간으로 제공해요. Flink의 이벤트 시간 처리에 대한 자세한 내용은 Event time을 참고해요. 이 타임스탬프는 일반적으로 Kinesis 서버 측 타임스탬프라고 하며, 정확성이나 순서 정확성에 대한 보장은 없어요(즉, 타임스탬프가 항상 오름차순이 아닐 수 있음).

샤드 리더의 이벤트 시간 정렬 (Event Time Alignment for Shard Readers)

샤드가 병렬로 소비되고 있으므로 일부 샤드가 다른 샤드보다 훨씬 앞선 레코드를 읽고 있을 수 있어요. Flink는 split 특정 워터마크 정렬로 이 읽기 속도 간의 동기화를 지원해요. KinesisStreamsSource는 split 특정 워터마크 정렬을 지원하며, 특정 샤드의 워터마크가 다른 샤드보다 너무 앞서 있으면 그 샤드에서 읽기를 일시 중지해요. 다른 샤드가 따라잡으면 그 split에서 읽기를 재개해요.

스레딩 모델 (Threading Model)

KinesisStreamsSource는 샤드 발견과 데이터 소비에 여러 스레드를 사용해요.

샤드 발견 (Shard Discovery)

샤드 발견을 위해 SplitEnumerator는 JobManager에서 실행되어 ListShard API로 주기적으로 새 샤드를 발견해요(shard discovery 참고). 새 샤드가 발견되면 부모 샤드가 완료되었는지 확인해요. 모든 부모가 완료되면 샤드는 읽기 위해 TaskManager의 SplitReader에 할당돼요.

폴링(기본) Split Reader (Polling (default) Split Reader)

POLLING 데이터 소비의 경우 병렬 하위 작업당 단일 스레드가 생성되어 할당된 샤드를 소비해요. 이는 열린 스레드 수가 Flink 연산자의 병렬도에 따라 확장됨을 의미해요.

Enhanced Fan-Out Split Reader

EFO 데이터 소비의 스레딩 모델은 POLLING과 같아요 — 병렬 하위 작업당 스레드 하나. 하지만 Kinesis와의 비동기 통신을 처리하는 추가 스레드 풀이 있어요. AWS SDK v2.x KinesisAsyncClient는 IO와 비동기 응답을 처리하는 데 Netty용 추가 스레드를 사용해요. 각 병렬 하위 작업은 KinesisAsyncClient의 자체 인스턴스를 가져요. 즉, consumer를 병렬도 10으로 실행하면 총 10개의 KinesisAsyncClient 인스턴스가 있어요. 스트림 consumer를 등록·등록 해제할 때 별도 클라이언트가 생성되고 이후 파괴돼요.

Enhanced Fan-Out 사용 (Using Enhanced Fan-Out)

Enhanced Fan-Out (EFO)는 Kinesis 스트림당 최대 동시 consumer 수를 늘려요. EFO 없이 모든 동시 consumer는 샤드당 단일 읽기 할당량을 공유해요. EFO를 사용하면 각 consumer가 샤드당 별개의 전용 읽기 할당량을 얻어 consumer 수에 따라 읽기 처리량이 확장될 수 있어요. EFO 사용은 추가 비용이 발생해요.

EFO를 활성화하려면 두 개의 추가 구성 파라미터가 필요해요:

  • READER_TYPE: EFO 또는 POLLING을 사용할지 결정. 기본 ReaderType은 POLLING.
  • EFO_CONSUMER_NAME: consumer를 식별하는 이름. 주어진 Kinesis data stream에 대해 각 consumer는 고유한 이름을 가져야 해요. 그러나 consumer 이름은 data stream 간에 고유할 필요는 없어요. consumer 이름을 재사용하면 기존 구독이 종료돼요.

아래 코드 조각은 EFO consumer를 구성하는 간단한 예제예요.

Java:

Configuration sourceConfig = new Configuration();
sourceConfig.set(KinesisSourceConfigOptions.READER_TYPE, KinesisSourceConfigOptions.ReaderType.EFO);
sourceConfig.set(KinesisSourceConfigOptions.EFO_CONSUMER_NAME, "my-flink-efo-consumer");

Scala:

val sourceConfig = new Configuration()
sourceConfig.set(KinesisSourceConfigOptions.READER_TYPE, KinesisSourceConfigOptions.ReaderType.EFO)
sourceConfig.set(KinesisSourceConfigOptions.EFO_CONSUMER_NAME, "my-flink-efo-consumer")

EFO 스트림 consumer 수명주기 관리 (EFO Stream Consumer Lifecycle Management)

EFO를 사용하려면 소비할 스트림에 대해 stream consumer를 등록해야 해요. 기본적으로 KinesisStreamsSource는 Flink 작업이 시작/중지될 때 stream consumer의 수명주기를 자동으로 관리해요. stream consumer는 EFO_CONSUMER_NAME 구성이 제공한 이름으로 등록되고, 작업이 정상적으로 중지되면 등록 해제돼요. KinesisStreamsSource는 KinesisSourceConfigOptions.EFO_CONSUMER_LIFECYCLE에 대한 두 가지 수명주기 옵션을 제공해요:

  • JOB_MANAGED(기본값): stream consumer는 Flink 작업이 실행되기 시작할 때 등록돼요. stream consumer가 이미 존재하면 재사용돼요. 작업이 정상적으로 중지되면 consumer가 등록 해제돼요. 대부분의 애플리케이션에 선호되는 전략이에요.
  • SELF_MANAGED: stream consumer 등록/등록 해제는 KinesisStreamsSource가 수행하지 않아요. 등록은 AWS CLI 또는 SDKRegisterStreamConsumer를 호출해 외부에서 수행해야 해요. stream consumer ARN은 consumer 구성을 통해 작업에 제공해야 해요.

내부적으로 사용되는 Kinesis API (Internally Used Kinesis APIs)

KinesisStreamsSource는 내부적으로 AWS v2 Java SDK를 사용해 샤드 발견과 데이터 소비에 Kinesis API를 호출해요. Amazon의 Kinesis Streams 서비스 한도 때문에 KinesisStreamsSource는 사용자가 실행 중일 수 있는 다른 비-Flink 소비 애플리케이션과 경쟁해요. 아래는 소비자가 호출하는 API 목록과 consumer가 API를 사용하는 방법에 대한 설명, 그리고 KinesisStreamsSource가 이 서비스 한도로 인해 가질 수 있는 오류나 경고를 처리하는 방법에 대한 정보예요.

재시도 전략 (Retry Strategy)

AWS SDK 클라이언트가 사용하는 재시도 전략은 다음 구성 옵션으로 튜닝할 수 있어요:

  • AWSConfigOptions.RETRY_STRATEGY_MAX_ATTEMPTS_OPTION: Flink 작업을 재시작하기 전에 재시도 가능한 오류에 대한 최대 API 재시도 수.
  • AWSConfigOptions.RETRY_STRATEGY_MIN_DELAY_OPTION: 지수 백오프 계산에 사용되는 기본 지연.
  • AWSConfigOptions.RETRY_STRATEGY_MAX_DELAY_OPTION: 지수 백오프의 최대 지연.

재시도 전략은 DescribeStreamConsumer를 제외한 모든 API 요청에 사용돼요. DescribeStreamConsumer용 재시도 전략을 구성하려면 대신 KinesisSourceConfigOptions.EFO_DESCRIBE_CONSUMER_RETRY_STRATEGY* 옵션을 사용해요. 이 API 호출은 작업 시작 중 EFO consumer 등록을 검증할 때만 호출되므로 별도 재시도 전략을 가져요. 정상 동작 중 사용되는 다른 API 호출과는 다른 재시도 전략이 필요할 수 있어요.

샤드 발견 (Shard Discovery)

  • ListShards: SplitEnumerator가 Flink 작업당 하나씩 주기적으로 호출해 스트림 재샤딩 결과로 새 샤드를 발견해요. 기본적으로 SplitEnumerator는 10초 간격으로 샤드 발견을 수행해요. 이것이 다른 비-Flink 소비 애플리케이션과 간섭하면 제공된 Configuration에서 KinesisSourceConfigOptions.SHARD_DISCOVERY_INTERVAL 값을 설정해 이 API 호출을 늦출 수 있어요. 이는 발견 간격을 다른 값으로 설정해요. 이 설정은 간격 동안 샤드가 발견되지 않으므로 새 샤드 발견과 소비 시작의 최대 지연에 직접 영향을 준다는 점에 주의해요.

폴링(기본) Split Reader (Polling (default) Split Reader)

  • GetShardIterator: 샤드당 한 번 호출돼요. 이 API의 속도 한도는 샤드당(스트림당이 아니라)이므로 split reader 자체가 한도를 초과하지 않아야 해요.
  • GetRecords: Kinesis에서 레코드를 가져오기 위해 지속적으로 호출돼요. 샤드에 여러 동시 consumer가 있으면(다른 비-Flink 소비 애플리케이션이 실행 중일 때) 샤드당 속도 한도를 초과할 수 있어요. 기본적으로 이 API를 호출할 때마다 Kinesis가 API의 데이터 크기/트랜잭션 한도 초과를 보고하면 consumer가 재시도해요.

Enhanced Fan-Out Split Reader

  • SubscribeToShard: 샤드 구독을 얻기 위해 샤드당 호출돼요. 샤드 구독은 보통 5분 동안 활성 상태이지만, 복구 가능한 오류가 발생하면 구독이 다시 획득돼요. 구독을 획득하면 consumer는 SubscribeToShardEvents 스트림을 받아요.
  • DescribeStreamConsumer: 애플리케이션 시작 중에만 연산자의 병렬 하위 작업당 호출돼요. 이는 스트림에 연결된 ACTIVE consumer의 consumerArn을 검색하기 위한 것이에요. 재시도 전략은 KinesisSourceConfigOptions.EFO_DESCRIBE_CONSUMER_RETRY_STRATEGY* 옵션으로 구성할 수 있어요.
  • RegisterStreamConsumer: SELF_MANAGED consumer 수명주기가 구성되지 않는 한 stream consumer 등록 중 스트림당 한 번 호출돼요.
  • DeregisterStreamConsumer: SELF_MANAGED 등록 전략이 구성되지 않는 한 stream consumer 등록 해제 중 스트림당 한 번 호출돼요.

Kinesis Consumer

이전 Kinesis 소스 org.apache.flink.streaming.connectors.kinesis.FlinkKinesisConsumer는 deprecated이며 향후 Flink 릴리스에서 제거될 수 있어요. Kinesis Source를 대신 사용해주세요. FlinkKinesisConsumer와 KinesisStreamsSource 사이에는 상태 호환성이 없다는 점에 주의해요. 자세한 내용은 migration 섹션을 참고해요.

기존 작업을 새로운 Kinesis Streams Source로 마이그레이션 (Migrating existing jobs to new Kinesis Streams Source from Kinesis Consumer)

FlinkKinesisConsumer와 KinesisStreamsSource 사이에는 상태 호환성이 없어요. 이는 FlinkKinesisConsumer에서 KinesisStreamsSource로 마이그레이션할 때 시작 위치가 손실됨을 의미해요. FlinkKinesisConsumer가 중지된 시간보다 약간 앞선 AT_TIMESTAMP에서 KinesisStreamsSource를 시작하는 것을 고려해요. 일부 레코드가 다시 처리될 수 있다는 점에 주의해요.

모범 사례를 따라 소스 연산자 uid를 지정하고 있다면, savepoint에서 Flink 작업을 복원할 때 uid를 바꾸고 allowNonRestoredState를 활성화해야 해요.

Kinesis Streams Sink

Kinesis Streams 싱크(이하 "Kinesis sink")는 AWS v2 Java SDK를 사용해 Flink 스트림의 데이터를 Kinesis 스트림에 써요. Kinesis 스트림에 데이터를 쓰려면 스트림이 Amazon Kinesis Data Stream 콘솔에서 "ACTIVE"로 표시되는지 확인해요. 모니터링이 동작하려면 스트림에 접근하는 사용자가 CloudWatch 서비스에 접근할 수 있어야 해요.

Java:

Properties sinkProperties = new Properties();
// Required
sinkProperties.put(AWSConfigConstants.AWS_REGION, "us-east-1");

// Optional, provide via alternative routes e.g. environment variables
sinkProperties.put(AWSConfigConstants.AWS_ACCESS_KEY_ID, "aws_access_key_id");
sinkProperties.put(AWSConfigConstants.AWS_SECRET_ACCESS_KEY, "aws_secret_access_key");

KinesisStreamsSink<String> kdsSink =
KinesisStreamsSink.<String>builder()
.setKinesisClientProperties(sinkProperties) // Required
.setSerializationSchema(new SimpleStringSchema()) // Required
.setPartitionKeyGenerator(element -> String.valueOf(element.hashCode())) // Required
.setStreamName("your-stream-name") // Required
.setFailOnError(false) // Optional
.setMaxBatchSize(500) // Optional
.setMaxInFlightRequests(50) // Optional
.setMaxBufferedRequests(10_000) // Optional
.setMaxBatchSizeInBytes(5 * 1024 * 1024) // Optional
.setMaxTimeInBufferMS(5000) // Optional
.setMaxRecordSizeInBytes(1 * 1024 * 1024) // Optional
.build();

DataStream<String> simpleStringStream = ...;
simpleStringStream.sinkTo(kdsSink);

Python:

# Required
sink_properties = {
# Required
'aws.region': 'us-east-1',
# Optional, provide via alternative routes e.g. environment variables
'aws.credentials.provider.basic.accesskeyid': 'aws_access_key_id',
'aws.credentials.provider.basic.secretkey': 'aws_secret_access_key',
'aws.endpoint': 'http://localhost:4567'
}

kds_sink = KinesisStreamsSink.builder() \
.set_kinesis_client_properties(sink_properties) \ # Required
.set_serialization_schema(SimpleStringSchema()) \ # Required
.set_partition_key_generator(PartitionKeyGenerator.fixed()) \ # Required
.set_stream_name("your-stream-name") \ # Required
.set_fail_on_error(False) \ # Optional
.set_max_batch_size(500) \ # Optional
.set_max_in_flight_requests(50) \ # Optional
.set_max_buffered_requests(10000) \ # Optional
.set_max_batch_size_in_bytes(5 * 1024 * 1024) \ # Optional
.set_max_time_in_buffer_ms(5000) \ # Optional
.set_max_record_size_in_bytes(1 * 1024 * 1024) \ # Optional
.build()

simple_string_stream = ...
simple_string_stream.sink_to(kds_sink)

위는 Kinesis sink를 사용하는 간단한 예제예요. AWS_REGION, AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY가 구성된 java.util.Properties 인스턴스를 만드는 것으로 시작해요. 그런 다음 빌더로 sink를 구성할 수 있어요. 선택적 구성의 기본값은 위에 나와 있어요. 이 값들 중 일부는 KDS 구성의 결과로 설정됐어요.

항상 직렬화 스키마와 레코드에서 파티션 키를 생성하는 로직을 지정해야 해요.

한 요청의 일부 또는 모든 레코드는 여러 이유로 Kinesis Data Streams가 영속화하지 못할 수 있어요. failOnError가 켜져 있으면 런타임 예외가 발생해요. 그렇지 않으면 해당 레코드는 재시도를 위해 버퍼에 다시 대기열에 들어가요.

Kinesis Sink는 Flink의 metrics system을 통해 커넥터의 동작을 분석할 수 있는 일부 메트릭을 제공해요. 노출된 모든 메트릭 목록은 여기에서 찾을 수 있어요.

sink의 기본 최대 레코드 크기는 1MB이고 최대 배치 크기는 5MB로 Kinesis Data Streams의 최대값과 일치해요. 이 최대값을 자세히 설명하는 AWS 문서는 여기에서 찾을 수 있어요.

Kinesis Sink와 장애 허용 (Kinesis Sinks and Fault Tolerance)

sink는 at-least-once 처리 보장을 제공하도록 Flink의 체크포인팅에 참여하도록 설계됐어요. 체크포인트를 찍는 동안 진행 중인 요청을 완료하는 방식으로 이를 수행해요. 이는 체크포인트 이전에 트리거된 모든 요청이 더 많은 레코드를 처리하기 전에 Kinesis Data Streams에 성공적으로 전달되었음을 효과적으로 보장해요.

Flink가 체크포인트(또는 savepoint)에서 복원해야 하면 체크포인트 이후 쓰여진 데이터가 Kinesis에 다시 쓰여 스트림에 중복이 생겨요. 또한 sink는 내부적으로 PutRecords API 호출을 사용하는데, 이는 이벤트 순서를 유지한다는 보장은 없어요.

백프레셔 (Backpressure)

sink의 백프레셔는 sink 버퍼가 가득 차고 sink에 대한 쓰기가 차단 동작을 보이기 시작할 때 발생해요. Kinesis Data Streams의 속도 제한에 대한 자세한 정보는 Quotas and Limits에서 찾을 수 있어요. 일반적으로 내부 큐의 크기를 늘려 백프레셔를 줄여요:

Java:

KinesisStreamsSink<String> kdsSink =
KinesisStreamsSink.<String>builder()
...
.setMaxBufferedRequests(10_000)
...

Python:

kds_sink = KinesisStreamsSink.builder() \
.set_max_buffered_requests(10000) \
.build()

Kinesis Producer

이전 Kinesis sink org.apache.flink.streaming.connectors.kinesis.FlinkKinesisProducer는 deprecated이며 향후 Flink 릴리스에서 제거될 수 있어요. Kinesis Sink를 대신 사용해주세요. 새 sink는 AWS v2 Java SDK를 사용하는 반면, 이전 sink는 Kinesis Producer Library를 사용해요. 이 때문에 새 Kinesis sink는 aggregation을 지원하지 않아요.

사용자 정의 Kinesis 엔드포인트 사용 (Using Custom Kinesis Endpoints)

Flink가 Kinesis VPC 엔드포인트나 Kinesalite 같은 비-AWS Kinesis 엔드포인트에 대해 소스 또는 싱크로 동작하게 하는 것이 바람직할 때가 있어요. 이는 Flink 애플리케이션의 기능 테스트를 수행할 때 특히 유용해요. Flink 구성에 설정된 AWS region으로 보통 추론되는 AWS 엔드포인트는 구성 프로퍼티로 오버라이드해야 해요.

AWS 엔드포인트를 오버라이드하려면 AWSConfigConstants.AWS_ENDPOINTAWSConfigConstants.AWS_REGION 프로퍼티를 설정해요. region은 엔드포인트 URL에 서명하는 데 사용돼요.

Java:

Properties config = new Properties();
config.put(AWSConfigConstants.AWS_REGION, "us-east-1");
config.put(AWSConfigConstants.AWS_ACCESS_KEY_ID, "aws_access_key_id");
config.put(AWSConfigConstants.AWS_SECRET_ACCESS_KEY, "aws_secret_access_key");
config.put(AWSConfigConstants.AWS_ENDPOINT, "http://localhost:4567");

Python:

config = {
'aws.region': 'us-east-1',
'aws.credentials.provider.basic.accesskeyid': 'aws_access_key_id',
'aws.credentials.provider.basic.secretkey': 'aws_secret_access_key',
'aws.endpoint': 'http://localhost:4567'
}

더 알아보기 (Learn more)