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 배포판이 있습니다. 이는 더 이상 사용되지 않는 SourceFunction 및 SinkFunction 인터페이스에서 새 Source 및 Sink 인터페이스로의 진행 중인 마이그레이션에서 비롯되었습니다.
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_TIMESTAMP의 source.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 | 실패한 요청을 재시도하는 플래그. 설정되면 요청 실패가 재시도되지 않고 작업이 실패함. |