Kinesis 커넥터

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

Kinesis 커넥터는 Amazon Kinesis Data Streams (KDS)에서 데이터를 읽고 쓸 수 있게 합니다. 스캔 소스는 무한(unbounded), 싱크는 배치/스트리밍 Append Mode로 동작합니다.

출처: 문서

본문

Kinesis 커넥터는 Amazon Kinesis Data Streams (KDS)에서 데이터를 읽고 쓸 수 있게 합니다.

의존성 (Dependencies)

Flink 버전 2.3용 커넥터는 아직 없습니다.

Kinesis 커넥터는 바이너리 배포판의 일부가 아닙니다. 클러스터 실행을 위해 연결하는 방법은 여기에서 확인하세요.

버저닝 (Versioning)

Kinesis 커넥터에는 두 가지 Table API 및 SQL 배포판이 있습니다. 이는 더 이상 사용되지 않는 SourceFunctionSinkFunction 인터페이스에서 새 SourceSink 인터페이스로의 진행 중인 마이그레이션에서 비롯되었습니다.

Flink의 Table API 및 SQL 인터페이스는 커넥터 식별자당 하나의 TableFactory만 허용합니다. 애플리케이션의 의존성에는 식별자 kinesis를 가진 TableFactory가 하나만 포함될 수 있습니다.

다음 표는 선택한 배포판에 따라 사용되는 기본 인터페이스를 명확히 합니다:

의존성 커넥터 버전 소스 커넥터 식별자 (인터페이스) 싱크 커넥터 식별자 (인터페이스)
flink-sql-connector-aws-kinesis-streams 5.x 이상 kinesis(Source) kinesis(Sink)
flink-sql-connector-aws-kinesis-streams 4.x 이하 N/A (소스 패키징 없음) kinesis(Sink)
flink-sql-connector-kinesis 5.x 이상 kinesis(Source), kinesis-legacy(SourceFunction) kinesis(Sink)
flink-sql-connector-kinesis 4.x 이하 kinesis(SourceFunction) kinesis(Sink)

flink-sql-connector-aws-kinesis-streams 또는 flink-sql-connector-kinesis 중 하나의 아티팩트만 포함하세요. 둘 다 포함하면 TableFactory 이름이 충돌합니다.

이 문서는 5.x 이상 버전을 대상으로 합니다. 주요 구성 섹션은 kinesis 식별자를 대상으로 합니다. 레거시 구성은 Configuration (kinesis-legacy) 문서를 참고하세요.

4.x와 5.x 사이에는 Table API와 SQL API 간 상태 호환성이 없습니다. 이는 기본 구현이 변경되었기 때문입니다.

v4.x kinesis 테이블로 작업을 중지한 시간보다 약간 이른 AT_TIMESTAMPsource.init.position으로 v5.x kinesis 테이블로 작업을 시작하는 것을 고려하세요. 이는 일부 레코드가 재처리될 수도 있음을 유의하세요.

Kinesis 데이터 스트림 테이블 만들기 (How to create a Kinesis data stream table)

Amazon KDS Developer Guide의 지침을 따라 Kinesis 스트림을 설정하세요. 다음 예시는 Kinesis 데이터 스트림으로 백업되는 테이블을 만드는 방법을 보여줍니다:

CREATE TABLE KinesisTable (
  `user_id` BIGINT,
  `item_id` BIGINT,
  `category_id` BIGINT,
  `behavior` STRING,
  `ts` TIMESTAMP(3)
)
PARTITIONED BY (user_id, item_id)
WITH (
  'connector' = 'kinesis',
  'stream.arn' = 'arn:aws:kinesis:us-east-1:012345678901:stream/my-stream-name',
  'aws.region' = 'us-east-1',
  'source.init.position' = 'LATEST',
  'format' = 'csv'
);

사용 가능한 메타데이터 (Available Metadata)

kinesis 테이블 Source에는 VIRTUAL 컬럼이 지원되지 않는 알려진 버그가 있습니다. 수정이 완료될 때까지 kinesis-legacy를 사용하세요.

다음 메타데이터는 테이블 정의에서 읽기 전용(VIRTUAL) 컬럼으로 노출될 수 있습니다. 이는 kinesis-legacy 커넥터에서만 사용할 수 있습니다.

데이터 타입 설명
timestamp TIMESTAMP_LTZ(3) NOT NULL 레코드가 스트림에 삽입된 대략적인 시간.
shard-id VARCHAR(128) NOT NULL 레코드가 읽힌 스트림 내 shard의 고유 식별자.
sequence-number VARCHAR(128) NOT NULL 해당 shard 내 레코드의 고유 식별자.

확장된 CREATE TABLE 예시는 이러한 메타데이터 필드를 노출하는 문법을 보여줍니다:

CREATE TABLE KinesisTable (
  `user_id` BIGINT,
  `item_id` BIGINT,
  `category_id` BIGINT,
  `behavior` STRING,
  `ts` TIMESTAMP(3),
  `arrival_time` TIMESTAMP(3) METADATA FROM 'timestamp' VIRTUAL,
  `shard_id` VARCHAR(128) NOT NULL METADATA FROM 'shard-id' VIRTUAL,
  `sequence_number` VARCHAR(128) NOT NULL METADATA FROM 'sequence-number' VIRTUAL
)
PARTITIONED BY (user_id, item_id)
WITH (
  'connector' = 'kinesis-legacy',
  'stream' = 'user_behavior',
  'aws.region' = 'us-east-2',
  'scan.stream.initpos' = 'LATEST',
  'format' = 'csv'
);

커넥터 옵션 (Connector Options)

옵션 필수 전달 기본값 타입 설명
공통 옵션
connector required no (none) String 사용할 커넥터를 지정. Kinesis의 경우 'kinesis' 또는 'kinesis-legacy'. 자세한 내용은 Versioning 참고.
stream.arn required yes (none) String 이 테이블을 백업하는 Kinesis 데이터 스트림의 이름.
format required no (none) String Kinesis 데이터 스트림 레코드를 역직렬화/직렬화하는 데 사용되는 포맷. Data Type Mapping 참고.
aws.region required no (none) String 스트림이 정의된 AWS 리전.
aws.endpoint optional no (none) String Kinesis의 AWS 엔드포인트(설정되지 않으면 AWS 리전 설정에서 파생됨).
aws.trust.all.certificates optional no false Boolean true이면 모든 SSL 인증서를 수락. 프로덕션 환경에는 권장되지 않으며 테스트 목적으로만 사용해야 함.
인증 옵션
aws.credentials.provider optional no AUTO String Kinesis 엔드포인트에 인증할 때 사용할 자격증명 제공자. Authentication 참고.
aws.credentials.basic.accesskeyid optional no (none) String 자격증명 제공자 타입을 BASIC으로 설정할 때 사용할 AWS 액세스 키 ID.
aws.credentials.basic.secretkey optional no (none) String 자격증명 제공자 타입을 BASIC으로 설정할 때 사용할 AWS 시크릿 키.
aws.credentials.profile.path optional no (none) String 자격증명 제공자 타입이 PROFILE일 때 프로필 경로에 대한 선택적 구성.
aws.credentials.profile.name optional no (none) String 자격증명 제공자 타입이 PROFILE일 때 프로필 이름에 대한 선택적 구성.
aws.credentials.role.arn optional no (none) String 자격증명 제공자 타입이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN일 때 사용할 역할 ARN.
aws.credentials.role.sessionName optional no (none) String 자격증명 제공자 타입이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN일 때 사용할 역할 세션 이름.
aws.credentials.role.externalId optional no (none) String 자격증명 제공자 타입이 ASSUME_ROLE일 때 사용할 외부 ID.
aws.credentials.role.stsEndpoint optional no (none) String 자격증명 제공자 타입이 ASSUME_ROLE일 때 사용할 STS용 AWS 엔드포인트(AWS 리전 설정에서 파생됨).
aws.credentials.role.provider optional no (none) String 자격증명 제공자 타입이 ASSUME_ROLE일 때 역할을 가정하기 위한 자격증명을 제공하는 제공자. 역할은 중첩될 수 있으므로 이 값은 다시 ASSUME_ROLE로 설정될 수 있음.
aws.credentials.webIdentityToken.file optional no (none) String 제공자 타입이 WEB_IDENTITY_TOKEN일 때 사용해야 하는 웹 ID 토큰 파일의 절대 경로.
aws.credentials.custom.class required only if credential provider is set to CUSTOM no (none) String 자격증명 제공자 타입이 CUSTOM일 때 사용할 사용자 제공 클래스의 전체 경로(Java 패키지 표기법). 예: org.user_company.auth.CustomAwsCredentialsProvider.
소스 옵션
source.init.position optional no LATEST String 테이블에서 읽을 때 사용할 초기 위치. Start Reading Position 참고.
source.init.timestamp optional no (none) String Kinesis 스트림 읽기를 시작할 초기 타임스탬프(scan.stream.initpos가 AT_TIMESTAMP일 때). Start Reading Position 참고.
source.init.timestamp.format optional no yyyy-MM-dd'T'HH:mm:ss.SSSXXX String Kinesis 스트림 읽기를 시작할 초기 타임스탬프의 날짜 형식(scan.stream.initpos가 AT_TIMESTAMP일 때). Start Reading Position 참고.
source.shard.discovery.interval optional no 10 s Duration 새 shard를 발견하려는 각 시도 사이의 간격.
source.reader.type optional no POLLING String 소스에 사용할 ReaderType (POLLING|EFO).
source.shard.get-records.max-record-count optional no 10000 Integer POLLING ReaderType에만 적용. AWS Kinesis shard에서 레코드를 가져올 때마다 가져오려는 최대 레코드 수.
source.efo.consumer.name optional no (none) String EFO ReaderType에만 적용. KDS에 등록할 EFO consumer의 이름.
source.efo.lifecycle optional no JOB_MANAGED String EFO ReaderType에만 적용. EFO consumer가 Flink 작업에 의해 관리되는지 여부를 결정 (JOB_MANAGED|SELF_MANAGED).
source.efo.subscription.timeout optional no 60 s Duration EFO ReaderType에만 적용. EFO Consumer 구독 타임아웃.
source.efo.deregister.timeout optional no 10 s Duration EFO ReaderType에만 적용. consumer 등록 해제 타임아웃. 타임아웃 도달 시 코드는 정상대로 계속됨.
source.efo.describe.retry-strategy.attempts.max optional no 100 Integer EFO ReaderType에만 적용. DescribeStreamConsumer 호출 시 지수 백오프 재시도 전략의 최대 시도 횟수.
source.efo.describe.retry-strategy.delay.min optional no 2 s Duration EFO ReaderType에만 적용. DescribeStreamConsumer 호출 시 지수 백오프 재시도 전략의 기본 지연.
source.efo.describe.retry-strategy.delay.max optional no 60 s Duration EFO ReaderType에만 적용. DescribeStreamConsumer 호출 시 지수 백오프 재시도 전략의 최대 지연.
싱크 옵션
sink.partitioner optional yes random or row-based String Flink 파티션에서 Kinesis shard로의 선택적 출력 파티셔닝. Sink Partitioning 참고.
sink.partitioner-field-delimiter optional yes String PARTITION BY 절에서 파생된 fields 기반 파티셔너의 선택적 필드 구분자. Sink Partitioning 참고.
sink.producer.* optional no (none) 레거시 커넥터가 이전에 사용하던 디프리케이트된 옵션. KinesisStreamsSink에 해당 대안이 있는 옵션은 각각의 속성에 매칭됨. 지원되지 않는 옵션은 사용자에게 경고로 로그됨.

자세한 옵션(소스/싱크 시작 위치, 메타데이터, Data Type Mapping 등)은 공식 문서를 참고하세요. 레거시 kinesis-legacy 커넥터 옵션(scan.*, sink.* 레거시 옵션, shard.consumer.error.recoverable[0].exception, scan.watermark.* 등)도 동일 페이지에 문서화되어 있습니다.

레거시 커넥터 소스 옵션의 예:

옵션 필수 전달 기본값 타입 설명
scan.shard.discovery.intervalmillis optional no 10000 Integer 새 shard를 발견하려는 각 시도 사이의 간격.
scan.shard.adaptivereads optional no false Boolean shard에서 적응형 읽기를 켜는 구성. AdaptivePollingRecordPublisher 문서 참고.
scan.shard.idle.interval optional no -1 Long 워터마크 생성을 위해 shard를 유휴로 간주할 기준이 되는 간격(밀리초). 양수 값은 일부 shard가 새 레코드를 받지 않아도 워터마크가 진행되게 함.
scan.watermark.sync.interval optional no 30000 Long 공유 워터마크 상태를 주기적으로 동기화하는 간격(밀리초).
scan.watermark.lookahead.millis optional no 0 Long 리더가 공유 전역 워터마크보다 앞서 진행할 수 있는 최대 델타(밀리초).
scan.watermark.sync.queue.capacity optional no 100 Integer shard 소비를 중단하기 전에 버퍼링되는 최대 레코드 수.

싱크 옵션(레거시)의 예:

옵션 필수 전달 기본값 타입 설명
sink.http-client.max-concurrency optional no 10000 Integer KinesisAsyncClient가 허용하는 최대 동시 요청 수.
sink.http-client.read-timeout optional no 360000 Integer KinesisAsyncClient가 요청을 보내는 최대 시간(ms).
sink.http-client.protocol.version optional no HTTP2 String Kinesis Client가 사용하는 Http 버전.
sink.batch.max-size optional yes 500 Integer 다운스트림 쓰기를 위해 KinesisAsyncClient에 전달되는 요소의 최대 배치 크기.
sink.requests.max-inflight optional yes 16 Integer 새 쓰기 요청을 차단하고 백프레셔를 적용하기 전에 KinesisAsyncClient의 미완료 요청에 대한 임계값.
sink.requests.max-buffered optional yes 10000 String 새 쓰기 요청을 차단하고 백프레셔를 적용하기 전에 KinesisAsyncClient의 버퍼링된 요청에 대한 임계값.
sink.flush-buffer.size optional yes 5242880 Long 플러시 전에 KinesisAsyncClient의 writer 버퍼에 대한 임계값(바이트).
sink.flush-buffer.timeout optional yes 5000 Long 플러시 전에 요소가 KinesisAsyncClient 버퍼에 있을 수 있는 최대 시간(밀리초).
sink.fail-on-error optional yes false Boolean 실패한 요청을 재시도하는 플래그. 설정되면 요청 실패가 재시도되지 않고 작업이 실패함.

더 알아보기 (Learn more)