Amazon SQS Sink

Amazon SQS Sink

SQS sink는 AWS v2 SDK for Java를 사용해 Amazon SQS에 데이터를 기록해요. Amazon SQS Developer Guide의 지침을 따라 SQS 메시지 큐를 설정하세요.

출처: Amazon SQS Sink

본문

커넥터를 사용하려면 프로젝트에 다음 Maven 의존성을 추가하세요:

Flink 버전 2.3용 커넥터는 아직 제공되지 않아요.

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");

// Optional, use following if you want to provide access via AssumeRole, Please make sure given IAM role has "sqs:SendMessage" permission
sinkProperties.setProperty(AWSConfigConstants.AWS_CREDENTIALS_PROVIDER, "ASSUME_ROLE");
sinkProperties.setProperty(AWSConfigConstants.AWS_ROLE_ARN, "replace-this-with-IAMRole-arn");
sinkProperties.setProperty(AWSConfigConstants.AWS_ROLE_SESSION_NAME, "any-session-name-string");

SqsSink<String> sqsSink =
        SqsSink.<String>builder()
                .setSerializationSchema(new SimpleStringSchema())                // Required
                .setSqsUrl("https://sqs.us-east-1.amazonaws.com/xxxx/test-sqs")  // Required
                .setSqsClientProperties(sinkProperties)                          // Required
                .setFailOnError(false)                                           // Optional
                .setMaxBatchSize(10)                                             // Optional
                .setMaxInFlightRequests(50)                                      // Optional
                .setMaxBufferedRequests(1_000)                                   // Optional
                .setMaxBatchSizeInBytes(256 * 1024)                         // Optional
                .setMaxTimeInBufferMS(5000)                                      // Optional
                .setMaxRecordSizeInBytes(256 * 1024)                            // Optional
                .build();

flinkStream.sinkTo(sqsSink)

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 SqsSink<String> sqsSink =
                SqsSink.<String>builder()
                        .setSerializationSchema(new SimpleStringSchema())                // Required
                        .setSqsUrl("https://sqs.us-east-1.amazonaws.com/xxxx/test-sqs")  // Required
                        .setSqsClientProperties(sinkProperties)                          // Required
                        .setFailOnError(false)                                           // Optional
                        .setMaxBatchSize(10)                                             // Optional
                        .setMaxInFlightRequests(50)                                      // Optional
                        .setMaxBufferedRequests(1_000)                                   // Optional
                        .setMaxBatchSizeInBytes(256 * 1024)                              // Optional
                        .setMaxTimeInBufferMS(5000)                                      // Optional
                        .setMaxRecordSizeInBytes(256 * 1024)                             // Optional
                        .build();
                        
                       
flinkStream.sinkTo(sqsSink)

Configurations

Flink의 SQS sink는 정적 빌더 SqsSink.<String>builder()로 만들어져요.

  • setSqsClientProperties(Properties sinkProperties) — 필수. SQS 클라이언트에 자격 증명, 리전, 기타 파라미터를 제공해요.
  • setSerializationSchema(SerializationSchema serializationSchema) — 필수. Sink에 직렬화 스키마를 제공해요. 이 스키마는 SQS로 보내기 전에 요소를 직렬화하는 데 사용돼요.
  • setSqsUrl(String sqsUrl) — 필수. 싱크할 SQS의 URL.
  • setFailOnError(boolean failOnError) — 선택. 기본값: false. SQS에 레코드를 기록하려는 실패한 요청이 Flink 잡을 재시작하게 하는 sink의 치명적 예외로 취급되는지 여부.
  • setMaxBatchSize(int maxBatchSize) — 선택. 기본값: 10. SQS에 기록할 배치의 최대 크기.
  • setMaxInFlightRequests(int maxInFlightRequests) — 선택. 기본값: 50. sink가 백프레셔를 적용하기 전에 허용되는 in-flight 요청의 최대 수.
  • setMaxBufferedRequests(int maxBufferedRequests) — 선택. 기본값: 5_000. 백프레셔가 적용되기 전에 sink에 버퍼링될 수 있는 레코드의 최대 수.
  • setMaxBatchSizeInBytes(int maxBatchSizeInBytes) — 선택. 기본값: 256 * 1024. 배치가 될 수 있는 최대 크기(바이트). 전송되는 모든 배치는 이 크기보다 작거나 같아요.
  • setMaxTimeInBufferMS(int maxTimeInBufferMS) — 선택. 기본값: 5000. 레코드가 플러시되기 전에 sink에 머무를 수 있는 최대 시간.
  • setMaxRecordSizeInBytes(int maxRecordSizeInBytes) — 선택. 기본값: 256 * 1024. sink가 수용할 최대 레코드 크기. 이보다 큰 레코드는 자동으로 거부돼요.
  • build() — SQS sink를 구성하고 반환해요.

더 알아보기 (Learn more)