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를 구성하고 반환해요.