Amazon DynamoDB SQL 커넥터

Amazon DynamoDB SQL 커넥터

DynamoDB 커넥터는 Amazon DynamoDB 에 데이터를 쓸 수 있게 해줍니다.

출처: 문서

본문

이 커넥터는 SinkBatchStreaming Append & Upsert Mode 를 지원합니다.

DynamoDB 커넥터는 Amazon DynamoDB 에 데이터를 쓸 수 있게 해줍니다.

의존성

Flink 2.3 버전용 커넥터는 (아직) 사용할 수 없습니다.

DynamoDB 테이블 만들기

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

CREATE TABLE DynamoDbTable (
  `user_id` BIGINT,
  `item_id` BIGINT,
  `category_id` BIGINT,
  `behavior` STRING
)
WITH (
  'connector' = 'dynamodb',
  'table-name' = 'user_behavior',
  'aws.region' = 'us-east-2'
);

커넥터 옵션 (Connector Options)

옵션 필수 기본값 유형 설명
공통 옵션 (Common Options)
connector required (none) String 사용할 커넥터를 지정합니다. DynamoDB 의 경우 'dynamodb'.
table-name required (none) String 사용할 DynamoDB 테이블 이름.
aws.region required (none) String DynamoDB 테이블이 정의된 AWS 리전.
aws.endpoint optional (none) String DynamoDB 용 AWS 엔드포인트.
aws.trust.all.certificates optional false Boolean true 이면 모든 SSL 인증서를 수락합니다.
인증 옵션 (Authentication Options)
aws.credentials.provider optional AUTO String Kinesis 엔드포인트에 인증할 때 사용할 자격 증명 제공자(credentials provider). 자세한 내용은 인증 참고.
aws.credentials.basic.accesskeyid optional (none) String 자격 증명 제공자 유형이 BASIC 일 때 사용할 AWS 액세스 키 ID.
aws.credentials.basic.secretkey optional (none) String 자격 증명 제공자 유형이 BASIC 일 때 사용할 AWS 시크릿 키.
aws.credentials.profile.path optional (none) String 자격 증명 제공자 유형이 PROFILE 일 때 사용할 프로필 경로 선택 구성.
aws.credentials.profile.name optional (none) String 자격 증명 제공자 유형이 PROFILE 일 때 사용할 프로필 이름 선택 구성.
aws.credentials.role.arn optional (none) String 자격 증명 제공자 유형이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN 일 때 사용할 역할 ARN.
aws.credentials.role.sessionName optional (none) String 자격 증명 제공자 유형이 ASSUME_ROLE 또는 WEB_IDENTITY_TOKEN 일 때 사용할 역할 세션 이름.
aws.credentials.role.externalId optional (none) String 자격 증명 제공자 유형이 ASSUME_ROLE 일 때 사용할 외부 ID.
aws.credentials.role.provider optional (none) String 자격 증명 제공자 유형이 ASSUME_ROLE 일 때 역할 가정을 위한 자격 증명을 제공하는 자격 증명 제공자. 역할은 중첩될 수 있으므로 이 값은 다시 ASSUME_ROLE 로 설정될 수 있습니다.
aws.credentials.webIdentityToken.file optional (none) String 제공자 유형이 WEB_IDENTITY_TOKEN 일 때 사용해야 하는 웹 자격 증명 토큰 파일의 절대 경로.
aws.credentials.custom.class required only if credential provider is set to CUSTOM (none) String 자격 증명 제공자 유형이 CUSTOM 일 때 사용할 사용자 제공 클래스의 전체 경로(Java 패키지 표기법). 예: org.user_company.auth.CustomAwsCredentialsProvider.
Sink 옵션
sink.batch.max-size optional 25 Integer DynamoDB 에 쓸 요소의 최대 배치 크기.
sink.requests.max-inflight optional 50 Integer DynamoDB 에 대한 최대 병렬 배치 요청 수.
sink.requests.max-buffered optional 10000 String 업스트림 작업 그래프에 백프레셔를 적용하기 전 입력 버퍼 크기.
sink.flush-buffer.timeout optional 5000 Long 요소가 버퍼에 있다가 플러시되기까지의 임계 시간(ms).
sink.fail-on-error optional false Boolean 실패한 요청을 재시도하기 위한 플래그. 설정하면 요청 실패는 재시도되지 않고 작업이 실패합니다.
sink.ignore-nulls optional false Boolean sink 에서 null 값을 무시할지 여부를 결정합니다. true 로 설정하면 null 값은 처리에서 제외됩니다.
HTTP Client 옵션
sink.http-client.max-concurrency optional 10000 Integer HTTP 클라이언트가 허용하는 최대 동시 요청 수.
sink.http-client.read-timeout optional 360000 Integer 기반 소켓에 대한 각 읽기의 타임아웃.

권한 (Authorization)

DynamoDB 테이블에 쓰기 권한을 허용하도록 적절한 IAM 정책을 만드세요.

인증 (Authentication)

배포 방식에 따라 DynamoDB 에 접근을 허용할 적절한 Credentials Provider 를 선택합니다. 기본적으로 AUTO Credentials Provider 가 사용됩니다. 배포 구성에 액세스 키 ID 와 시크릿 키가 설정되어 있으면 BASIC 제공자가 사용됩니다.

특정 AWSCredentialsProvideraws.credentials.provider 설정을 사용해 선택적으로 설정할 수 있습니다. 지원되는 값:

  • AUTO - 다음 순서로 자격 증명을 검색하는 기본 AWS Credentials Provider 체인을 사용합니다: ENV_VARS, SYS_PROPS, WEB_IDENTITY_TOKEN, PROFILE, EC2/ECS credentials provider.
  • BASIC - 구성으로 제공된 액세스 키 ID 와 시크릿 키를 사용합니다.
  • ENV_VAR - AWS_ACCESS_KEY_ID & AWS_SECRET_ACCESS_KEY 환경 변수를 사용합니다.
  • SYS_PROP - Java 시스템 속성 aws.accessKeyIdaws.secretKey 를 사용합니다.
  • PROFILE - AWS 자격 증명 프로필을 사용해 AWS 자격 증명을 만듭니다.
  • ASSUME_ROLE - 역할을 가정해 AWS 자격 증명을 만듭니다. 역할 가정을 위한 자격 증명을 제공해야 합니다.
  • WEB_IDENTITY_TOKEN - Web Identity Token 을 사용해 역할을 가정해 AWS 자격 증명을 만듭니다.
  • CUSTOM - 인터페이스 AWSCredentialsProvider 를 구현하고 생성자 MyCustomClass(java.util.Properties config) 를 가진 사용자 정의 클래스를 제공합니다. 모든 커넥터 속성은 생성자를 통해 이 사용자 정의 자격 증명 제공자 클래스에 전달됩니다.

Sink 파티셔닝 (Sink Partitioning)

DynamoDB sink 는 PARTITIONED BY 절을 통해 데이터의 클라이언트 측 중복 제거를 지원합니다. 파티션 키 목록을 지정할 수 있으며, sink 는 배치 내에서 각 복합 키의 최신 레코드만 보냅니다. 예:

CREATE TABLE DynamoDbTable (
  `user_id` BIGINT,
  `item_id` BIGINT,
  `category_id` BIGINT,
  `behavior` STRING
) PARTITIONED BY ( user_id )
WITH (
  'connector' = 'dynamodb',
  'table-name' = 'user_behavior',
  'aws.region' = 'us-east-2'
);

공지 (Notice)

DynamoDB SQL 커넥터의 현재 구현은 쓰기 전용이며 소스 쿼리용 구현을 제공하지 않습니다. 다음과 같은 쿼리는:

SELECT * FROM DynamoDbTable;

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

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

더 알아보기 (Learn more)