Amazon Kinesis Data Firehose SQL Connector

Amazon Kinesis Data Firehose SQL Connector

Sink: Batch, Streaming Append Mode

Kinesis Data Firehose 커넥터는 Amazon Kinesis Data Firehose (KDF)로 데이터를 쓸 수 있게 합니다.

출처: 문서

본문

Kinesis Data Firehose 커넥터는 Amazon Kinesis Data Firehose (KDF)로 데이터를 쓸 수 있게 합니다.

의존성 (Dependencies)

Flink 2.3 버전용 커넥터는 (아직) 제공되지 않습니다.

Kinesis Data Firehose 테이블 만들기 (How to create a Kinesis Data Firehose table)

Amazon Kinesis Data Firehose Developer Guide의 지침에 따라 Kinesis Data Firehose delivery stream을 설정합니다. 다음 예제는 최소 필수 옵션으로 Kinesis Data Firehose delivery stream을 기반으로 하는 테이블을 만드는 방법을 보여줍니다:

CREATE TABLE FirehoseTable (
  `user_id` BIGINT,
  `item_id` BIGINT,
  `category_id` BIGINT,
  `behavior` STRING
)
WITH (
  'connector' = 'firehose',
  'delivery-stream' = 'user_behavior',
  'aws.region' = 'us-east-2',
  'format' = 'csv'
);

커넥터 옵션 (Connector Options)

Option Required Default Type Description
Common Options
connector required (none) String 사용할 커넥터를 지정합니다. Kinesis Data Firehose에는 'firehose'를 사용합니다.
delivery-stream required (none) String 이 테이블을 뒷받침하는 Kinesis Data Firehose delivery stream의 이름.
format required (none) String Kinesis Data Firehose 레코드를 역직렬화/직렬화하는 데 사용되는 포맷. 자세한 내용은 Data Type Mapping 참고.
aws.region required (none) String delivery stream이 정의된 AWS 리전. KinesisFirehoseSink 생성에 필요.
aws.endpoint optional (none) String Amazon Kinesis Data Firehose용 AWS 엔드포인트.
aws.trust.all.certificates optional false Boolean true면 모든 SSL 인증서를 수락합니다.
Authentication Options
aws.credentials.provider optional AUTO String Kinesis 엔드포인트 인증 시 사용할 자격 증명 프로바이더. Authentication 참고.
aws.credentials.basic.accesskeyid optional (none) String credential provider 유형을 BASIC으로 설정할 때 사용할 AWS access key ID.
aws.credentials.basic.secretkey optional (none) String credential provider 유형을 BASIC으로 설정할 때 사용할 AWS secret key.
aws.credentials.profile.path optional (none) String credential provider 유형이 PROFILE인 경우 프로파일 경로 설정.
aws.credentials.profile.name optional (none) String credential provider 유형이 PROFILE인 경우 프로파일 이름 설정.
aws.credentials.role.arn optional (none) String credential provider 유형이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN일 때 사용할 role ARN.
aws.credentials.role.sessionName optional (none) String credential provider 유형이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN일 때 사용할 role session name.
aws.credentials.role.externalId optional (none) String credential provider 유형이 ASSUME_ROLE일 때 사용할 external ID.
aws.credentials.role.stsEndpoint optional (none) String credential provider 유형이 ASSUME_ROLE일 때 사용할 STS용 AWS 엔드포인트(설정하지 않으면 AWS region 설정에서 파생).
aws.credentials.role.provider optional (none) String credential provider 유형이 ASSUME_ROLE일 때 role을 수임하기 위한 자격 증명을 제공하는 credentials provider. Role은 중첩될 수 있으므로 이 값도 다시 ASSUME_ROLE로 설정될 수 있음.
aws.credentials.webIdentityToken.file optional (none) String provider 유형이 WEB_IDENTITY_TOKEN일 때 사용할 web identity token 파일의 절대 경로.
aws.credentials.custom.class required only if credential provider is set to CUSTOM (none) String credential provider 유형이 CUSTOM일 때 사용할 사용자 제공 클래스의 전체 경로(Java 패키지 표기법). 예: org.user_company.auth.CustomAwsCredentialsProvider.
Sink Options
sink.http-client.max-concurrency optional 10000 Integer FirehoseAsyncClient가 delivery stream에 전달할 수 있는 최대 동시 요청 수.
sink.http-client.read-timeout optional 360000 Integer FirehoseAsyncClient가 실패 전까지 delivery stream에 요청을 보내는 최대 시간(ms).
sink.http-client.protocol.version optional HTTP2 String FirehoseAsyncClient가 사용하는 Http 버전.
sink.batch.max-size optional 500 Integer FirehoseAsyncClient에 전달되어 delivery stream에 하위로 쓰여질 요소의 최대 배치 크기.
sink.requests.max-inflight optional 16 Integer 새 쓰기 요청을 차단하기 전 FirehoseAsyncClient의 미완료 요청 임계값.
sink.requests.max-buffered optional 10000 String 새 쓰기 요청을 차단하기 전 FirehoseAsyncClient의 요청 버퍼 임계값.
sink.flush-buffer.size optional 5242880 Long 플러시 전 FirehoseAsyncClient의 writer 버퍼에 대한 임계값(바이트).
sink.flush-buffer.timeout optional 5000 Long 플러시 전 요소가 FirehoseAsyncClient의 버퍼에 있게 되는 임계 시간(ms).
sink.fail-on-error optional false Boolean 실패한 요청을 재시도하기 위한 플래그. 설정하면 어떤 요청 실패도 재시도되지 않고 job이 실패합니다.

권한 부여 (Authorization)

Kinesis Data Firehose delivery stream에 대한 읽기/쓰기를 허용하도록 적절한 IAM 정책을 생성하세요.

인증 (Authentication)

배포 방식에 따라 Kinesis Data Firehose에 접근을 허용할 다른 Credentials Provider를 선택할 수 있습니다. 기본적으로 AUTO Credentials Provider가 사용됩니다. 배포 구성에 access key ID와 secret key가 설정되어 있으면 BASIC 프로바이더가 사용됩니다.

특정 AWSCredentialsProvideraws.credentials.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 사용.
  • PROFILE - AWS 자격 증명을 만들기 위해 AWS credentials 프로파일 사용.
  • ASSUME_ROLE - role을 수임하여 AWS 자격 증명 생성. role 수임을 위한 자격 증명을 제공해야 함.
  • WEB_IDENTITY_TOKEN - Web Identity Token을 사용하여 role을 수임해 AWS 자격 증명 생성.
  • CUSTOM - AWSCredentialsProvider 인터페이스를 구현하고 MyCustomClass(java.util.Properties config) 생성자를 가지는 커스텀 클래스 제공. 모든 커넥터 프로퍼티는 생성자를 통해 이 커스텀 credential provider 클래스에 전달됩니다.

데이터 유형 매핑 (Data Type Mapping)

Kinesis Data Firehose는 레코드를 Base64로 인코딩된 바이너리 데이터 객체로 저장하므로 내부 레코드 구조라는 개념이 없습니다. 대신 Kinesis Data Firehose 레코드는 'avro', 'csv', 'json' 같은 포맷에 의해 역직렬화/직렬화됩니다. Kinesis Data Firehose 기반 테이블의 메시지 데이터 유형을 결정하려면 format 키워드로 적절한 Flink 포맷을 선택하세요. 자세한 내용은 Formats 페이지를 참고하세요.

주의 (Notice)

Kinesis Data Firehose SQL 커넥터의 현재 구현은 Kinesis Data Firehose 기반 sink만 지원하며 source 질의에 대한 구현은 제공하지 않습니다. 다음과 같은 질의를:

SELECT * FROM FirehoseTable;

실행하면 다음과 유사한 오류가 발생해야 합니다:

Connector firehose can only be used as a sink. It cannot be used as a source.

더 알아보기 (Learn more)