Amazon DynamoDB 커넥터
Amazon DynamoDB 커넥터 (Amazon DynamoDB Connector)
DynamoDB 커넥터는 사용자가 Amazon DynamoDB에서 읽기/쓰기를 할 수 있게 해줍니다.
소스로서 커넥터는 Amazon DynamoDB Streams를 사용해 DynamoDB 테이블에서 변경 데이터 캡처(change data capture) 스트림을 읽을 수 있게 해줍니다.
싱크로서 커넥터는 BatchWriteItem API를 사용해 Amazon DynamoDB 테이블에 직접 쓸 수 있게 해줍니다.
출처: 문서
본문
의존성 (Dependency)
Apache Flink는 이 커넥터를 사용자가 활용할 수 있도록 제공합니다. 커넥터를 사용하려면 프로젝트에 다음 Maven 의존성을 추가하세요:
현재 Flink 2.3 버전용 커넥터는 아직 없습니다.
Amazon DynamoDB Streams 소스 (Amazon DynamoDB Streams Source)
DynamoDB Streams 소스는 AWS v2 SDK for Java를 사용해 Amazon DynamoDB Streams에서 읽습니다. 변경 데이터 캡처 스트림을 설정하고 구성하려면 AWS 문서의 지침을 따르세요.
DynamoDB Stream으로 스트리밍되는 실제 이벤트는 DynamoDB Stream 자체가 지정한 StreamViewType에 따라 달라집니다. 자세한 내용은 AWS 문서를 참고하세요.
사용법 (Usage)
DynamoDbStreamsSource는 DynamoDbStreamsSource 인스턴스를 구성하기 위한 유창한(fluent) 빌더를 제공합니다. 아래 코드 스니펫은 그 방법을 보여줍니다.
Java
// Configure the DynamodbStreamsSource
Configuration sourceConfig = new Configuration();
sourceConfig.set(DynamodbStreamsSourceConfigConstants.STREAM_INITIAL_POSITION, DynamodbStreamsSourceConfigConstants.InitialPosition.TRIM_HORIZON); // This is optional, by default connector will read from LATEST
// Create a new DynamoDbStreamsSource to read from the specified DynamoDB Stream.
DynamoDbStreamsSource<String> dynamoDbStreamsSource =
DynamoDbStreamsSource.<String>builder()
.setStreamArn("arn:aws:dynamodb:us-east-1:1231231230:table/test/stream/2024-04-11T07:14:19.380")
.setSourceConfig(sourceConfig)
// User must implement their own deserialization schema to translate change data capture events into custom data types
.setDeserializationSchema(dynamodbDeserializationSchema)
.build();
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// Specify watermarking strategy and the name of the DynamoDB Streams Source operator.
// Specify return type using TypeInformation.
// Specify UID of operator in line with Flink best practice.
DataStream<String> cdcEventsWithEventTimeWatermarks = env.fromSource(dynamoDbStreamsSource, WatermarkStrategy.<String>forMonotonousTimestamps().withIdleness(Duration.ofSeconds(1)), "DynamoDB Streams source")
.returns(TypeInformation.of(String.class))
.uid("custom-uid");
Scala
// Configure the DynamodbStreamsSource
val sourceConfig = new Configuration()
sourceConfig.set(DynamodbStreamsSourceConfigConstants.STREAM_INITIAL_POSITION, DynamodbStreamsSourceConfigConstants.InitialPosition.TRIM_HORIZON) // This is optional, by default connector will read from LATEST
// Create a new DynamoDbStreamsSource to read from the specified DynamoDB Stream.
val dynamoDbStreamsSource = DynamoDbStreamsSource.builder[String]()
.setStreamArn("arn:aws:dynamodb:us-east-1:1231231230:table/test/stream/2024-04-11T07:14:19.380")
.setSourceConfig(sourceConfig)
// User must implement their own deserialization schema to translate change data capture events into custom data types
.setDeserializationSchema(dynamodbDeserializationSchema)
.build()
val env = StreamExecutionEnvironment.getExecutionEnvironment()
// Specify watermarking strategy and the name of the DynamoDB Streams Source operator.
// Specify return type using TypeInformation.
// Specify UID of operator in line with Flink best practice.
val cdcEventsWithEventTimeWatermarks = env.fromSource(dynamoDbStreamsSource, WatermarkStrategy.<String>forMonotonousTimestamps().withIdleness(Duration.ofSeconds(1)), "DynamoDB Streams source")
.uid("custom-uid")
위는 DynamoDbStreamsSource를 사용하는 간단한 예시입니다.
- 읽고 있는 DynamoDB Stream은 스트림 ARN으로 지정됩니다.
Source의 구성은 Flink의Configuration클래스 인스턴스로 제공됩니다. 구성 키는AWSConfigOptions(AWS 특화 구성)와DynamodbStreamsSourceConfigConstants(DynamoDB Streams Source 구성)에서 가져올 수 있습니다.- 예시는 시작 위치를
TRIM_HORIZON으로 지정합니다(자세한 내용은 Configuring Starting Position 참고). - 역직렬화 포맷은
SimpleStringSchema입니다(자세한 내용은 Deserialization Schema 참고). - 서브태스크 간 샤드 분포는
UniformShardAssigner로 제어됩니다(자세한 내용은 Shard Assignment Strategy 참고). - 예시는 증가하는(Increasing)
WatermarkStrategy도 지정하며, 이는 각 레코드가approximateCreationDateTime으로 지정된 이벤트 시간으로 태그됨을 의미합니다. 단조 증가하는 워터마크가 생성되며, 1초 동안 레코드가 내보내지지 않으면 서브태스크가 유휴으로 간주됩니다.
시작 위치 구성 (Configuring Starting Position)
DynamodbStreamsSource의 시작 위치를 지정하려면 사용자는 구성에서 DynamodbStreamsSourceConfigConstants.STREAM_INITIAL_POSITION을 설정할 수 있습니다.
LATEST: 가장 최근 레코드부터 스트림의 모든 샤드를 읽습니다.TRIM_HORIZON: 가능한 가장 이른 레코드부터 스트림의 모든 샤드를 읽습니다(DynamoDB가 24시간 후 데이터를 트리밍합니다).
역직렬화 스키마 (Deserialization Schema)
DynamoDbStreamsSource는 사용자가 자신만의 역직렬화 스키마를 구현해 DynamoDB 변경 데이터 캡처 이벤트를 커스텀 이벤트 타입으로 변환할 수 있게 하는 DynamoDbStreamsDeserializationSchema<T> 인터페이스를 제공합니다.
DynamoDbStreamsDeserializationSchema<T>#deserialize 메서드는 DynamoDB 모델의 Record 인스턴스를 받습니다. Record는 DynamoDB Stream의 구성에 따라 다른 내용을 포함할 수 있습니다. 자세한 내용은 AWS 문서를 참고하세요.
이벤트 순서 (Event Ordering)
이벤트는 DynamoDB Streams에 기록되며, 같은 기본 키 내에서 순서를 유지합니다. 이는 같은 기본 키 내의 이벤트가 같은 샤드 계보(lineage)에 기록되도록 보장해 수행됩니다. 샤드 스플릿이 있을 때(부모 샤드 하나가 두 자식 샤드로 분할), 부모 샤드를 자식 샤드 읽기를 시작하기 전에 완전히 읽는 한 순서가 유지됩니다.
DynamoDbStreamsSource는 부모-자식 샤드 순서를 존중하는 방식으로 샤드가 할당되도록 보장합니다. 즉, 부모 샤드가 완전히 읽힌 경우에만 샤드가 샤드 할당기에 전달됩니다. 이는 같은 DynamoDB 기본 키 내에서 변경 데이터 캡처 스트림의 이벤트가 순서대로 읽히도록 보장하는 데 도움이 됩니다.
샤드 할당 전략 (Shard Assignment Strategy)
UniformShardAssigner는 DynamoDB Stream의 샤드를 소스 연산자의 병렬 서브태스크에 고르게 할당합니다. DynamoDB Stream 샤드는 일시적이며 필요에 따라 자동으로 생성·삭제됩니다. UniformShardAssigner는 현재 할당된 샤드 수가 가장 적은 서브태스크에 새 샤드를 할당합니다.
사용자는 DynamoDbStreamsShardAssigner 인터페이스를 구현해 자신만의 샤드 할당 전략을 구현할 수도 있습니다.
구성 (Configuration)
재시도 전략 (Retry Strategy)
DynamoDbStreamsSource는 AWS v2 SDK for Java를 사용해 Amazon DynamoDB와 상호작용합니다.
AWS SDK 클라이언트가 사용하는 재시도 전략은 다음 구성 옵션으로 튜닝할 수 있습니다:
DYNAMODB_STREAMS_RETRY_COUNT: 재시도 가능한 오류에 대한 API 재시도의 최대 횟수. 초과하면 Flink 작업을 재시작합니다.DYNAMODB_STREAMS_EXPONENTIAL_BACKOFF_MIN_DELAY: 지수 백오프 계산에 사용되는 기본 지연.DYNAMODB_STREAMS_EXPONENTIAL_BACKOFF_MAX_DELAY: 지수 백오프의 최대 지연.
샤드 발견 (Shard Discovery)
DynamoDbStreamsSource는 주기적으로 DynamoDB Stream에서 새로 생성된 샤드를 발견합니다. 이는 샤드 스플릿이나 샤드 회전에서 올 수 있습니다. 기본적으로 매 60초마다 샤드를 발견하도록 설정됩니다. 그러나 사용자는 SHARD_DISCOVERY_INTERVAL을 구성해 더 작은 값으로 커스터마이즈할 수 있습니다.
샤드 발견에는 DynamoDB가 반환하는 샤드 그래프가 불일치할 수 있는 문제가 있습니다. 이 경우 DynamoDbStreamsSource가 자동으로 불일치를 감지하고 샤드 발견 과정을 재시도합니다. 최대 재시도 횟수는 DESCRIBE_STREAM_INCONSISTENCY_RESOLUTION_RETRY_COUNT로 구성할 수 있습니다.
Amazon DynamoDB 싱크 (Amazon DynamoDB Sink)
DynamoDB 싱크는 AWS v2 SDK for Java를 사용해 Amazon DynamoDB에 씁니다. 테이블을 설정하려면 Amazon DynamoDB Developer Guide의 지침을 따르세요.
Java
Properties sinkProperties = new Properties();
// Required
sinkProperties.put(AWSConfigConstants.AWS_REGION, "eu-west-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");
ElementConverter<InputType, DynamoDbWriteRequest> elementConverter = new CustomElementConverter();
DynamoDbSink<String> dynamoDbSink =
DynamoDbSink.<InputType>builder()
.setDynamoDbProperties(sinkProperties) // Required
.setTableName("my-dynamodb-table") // Required
.setElementConverter(elementConverter) // Required
.setOverwriteByPartitionKeys(singletonList("key")) // Optional
.setFailOnError(false) // Optional
.setMaxBatchSize(25) // Optional
.setMaxInFlightRequests(50) // Optional
.setMaxBufferedRequests(10_000) // Optional
.setMaxTimeInBufferMS(5000) // Optional
.build();
flinkStream.sinkTo(dynamoDbSink);
Scala
val sinkProperties = new Properties()
// Required
sinkProperties.put(AWSConfigConstants.AWS_REGION, "eu-west-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")
val elementConverter = new CustomElementConverter();
val dynamoDbSink =
DynamoDbSink.<InputType>builder()
.setDynamoDbProperties(sinkProperties) // Required
.setTableName("my-dynamodb-table") // Required
.setElementConverter(elementConverter) // Required
.setOverwriteByPartitionKeys(singletonList("key")) // Optional
.setFailOnError(false) // Optional
.setMaxBatchSize(25) // Optional
.setMaxInFlightRequests(50) // Optional
.setMaxBufferedRequests(10_000) // Optional
.setMaxTimeInBufferMS(5000) // Optional
.build()
flinkStream.sinkTo(dynamoDbSink)
구성 (Configurations)
Flink의 DynamoDB 싱크는 정적 빌더 DynamoDBSink.<InputType>builder()로 생성됩니다.
setDynamoDbProperties(Properties sinkProperties)- 필수.
- DynamoDB 클라이언트에 자격 증명, 지역 및 기타 파라미터를 제공합니다.
setTableName(String tableName)- 필수.
- 싱크할 테이블의 이름.
setElementConverter(ElementConverter<InputType, DynamoDbWriteRequest> elementConverter)- 필수.
InputType타입의 일반 레코드를DynamoDbWriteRequest로 변환합니다.
setOverwriteByPartitionKeys(List partitionKeys)- 선택. 기본값: [].
- DynamoDB에 푸시되는 각 배치 내에서 쓰기 요청을 중복 제거하는 데 사용됩니다.
setFailOnError(boolean failOnError)- 선택. 기본값:
false. - 레코드 쓰기 실패 요청이 싱크에서 치명적인 예외로 처리되는지 여부.
- 선택. 기본값:
setMaxBatchSize(int maxBatchSize)- 선택. 기본값:
25. - 쓸 배치의 최대 크기.
- 선택. 기본값:
setMaxInFlightRequests(int maxInFlightRequests)- 선택. 기본값:
50. - 싱크가 백프레셔를 적용하기 전에 허용되는 in-flight 요청의 최대 수.
- 선택. 기본값:
setMaxBufferedRequests(int maxBufferedRequests)- 선택. 기본값:
10_000. - 백프레셔가 적용되기 전에 싱크에 버퍼링될 수 있는 레코드의 최대 수.
- 선택. 기본값:
setMaxBatchSizeInBytes(int maxBatchSizeInBytes)- N/A.
- 이 구성은 지원되지 않습니다. FLINK-29854 참고.
setMaxTimeInBufferMS(int maxTimeInBufferMS)- 선택. 기본값:
5000. - 플러시되기 전에 레코드가 싱크에 머무를 수 있는 최대 시간.
- 선택. 기본값:
setMaxRecordSizeInBytes(int maxRecordSizeInBytes)- N/A.
- 이 구성은 지원되지 않습니다. FLINK-29854 참고.
build()- DynamoDB 싱크를 구성하고 반환합니다.
요소 변환기 (Element Converter)
요소 변환기는 DataStream의 레코드를 싱크가 대상 DynamoDB 테이블에 쓸 DynamoDbWriteRequest로 변환하는 데 사용됩니다. DynamoDB 싱크는 사용자가 커스텀 요소 변환기를 제공하거나, 요소 클래스에서 항목 스키마를 추출하는 제공된 DefaultDynamoDbElementConverter를 사용할 수 있게 해줍니다. 이는 요소 클래스가 복합 타입(즉, Pojo, Tuple 또는 Row)이어야 합니다. 요소의 TypeInformation이 있으면 new DynamoDbTypeInformedElementConverter(TypeInformation.of(MyPojo.class))처럼 DynamoDbTypeInformedElementConverter를 사용해 스키마를 즉시 구성합니다.
또는 @DynamoDbBean 객체로 작업할 때 DynamoDbBeanElementConverter를 사용할 수 있습니다. 지원되는 어노테이션에 대한 자세한 내용은 여기를 참고하세요.
커스텀 ElementConverter를 사용하는 샘플 애플리케이션은 여기, DynamoDbBeanElementConverter를 사용하는 샘플 애플리케이션은 여기에서 찾을 수 있습니다.
커스텀 DynamoDB 엔드포인트 사용 (Using Custom DynamoDB Endpoints)
때로는 Flink가 DynamoDB VPC 엔드포인트나 Localstack 같은 비-AWS DynamoDB 엔드포인트에 대해 컨슈머 또는 프로듀서로 동작하게 하는 것이 바람직할 수 있습니다. 이는 Flink 애플리케이션의 기능 테스트를 수행할 때 특히 유용합니다. Flink 구성에 설정된 AWS 지역에서 보통 추론되는 AWS 엔드포인트는 구성 속성으로 재정의되어야 합니다.
AWS 엔드포인트를 재정의하려면 AWSConfigConstants.AWS_ENDPOINT와 AWSConfigConstants.AWS_REGION 속성을 설정하세요. 지역은 엔드포인트 URL에 서명하는 데 사용됩니다.
Java
Properties producerConfig = new Properties();
producerConfig.put(AWSConfigConstants.AWS_REGION, "eu-west-1");
producerConfig.put(AWSConfigConstants.AWS_ACCESS_KEY_ID, "aws_access_key_id");
producerConfig.put(AWSConfigConstants.AWS_SECRET_ACCESS_KEY, "aws_secret_access_key");
producerConfig.put(AWSConfigConstants.AWS_ENDPOINT, "http://localhost:4566");
Scala
val producerConfig = new Properties()
producerConfig.put(AWSConfigConstants.AWS_REGION, "eu-west-1")
producerConfig.put(AWSConfigConstants.AWS_ACCESS_KEY_ID, "aws_access_key_id")
producerConfig.put(AWSConfigConstants.AWS_SECRET_ACCESS_KEY, "aws_secret_access_key")
producerConfig.put(AWSConfigConstants.AWS_ENDPOINT, "http://localhost:4566")