Amazon Kinesis Data Firehose Sink
Amazon Kinesis Data Firehose Sink
Firehose sink는 Amazon Kinesis Data Firehose에 씁니다.
Amazon Kinesis Data Firehose Developer Guide의 지침에 따라 Kinesis Data Firehose delivery stream을 설정하세요.
커넥터를 사용하려면 프로젝트에 다음 Maven 의존성을 추가하세요.
Flink 2.3 버전에 사용 가능한 커넥터는 아직 없습니다.
PyFlink 작업에서 사용하려면 다음 의존성이 필요합니다.
| Version | PyFlink JAR |
|---|---|
| flink-connector-aws-kinesis-firehose | Flink 2.3 버전에 사용 가능한 SQL jar는 아직 없습니다. |
PyFlink에서 JAR을 사용하는 방법에 대한 자세한 내용은 Python dependency management를 참조하세요.
KinesisFirehoseSink는 AWS v2 SDK for Java를 사용해 Flink 스트림의 데이터를 Firehose delivery stream에 씁니다.
출처: 문서
본문
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");
KinesisFirehoseSink<String> kdfSink =
KinesisFirehoseSink.<String>builder()
.setFirehoseClientProperties(sinkProperties) // Required
.setSerializationSchema(new SimpleStringSchema()) // Required
.setDeliveryStreamName("your-stream-name") // Required
.setFailOnError(false) // Optional
.setMaxBatchSize(500) // Optional
.setMaxInFlightRequests(50) // Optional
.setMaxBufferedRequests(10_000) // Optional
.setMaxBatchSizeInBytes(4 * 1024 * 1024) // Optional
.setMaxTimeInBufferMS(5000) // Optional
.setMaxRecordSizeInBytes(1000 * 1024) // Optional
.build();
flinkStream.sinkTo(kdfSink);
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 kdfSink =
KinesisFirehoseSink.<String>builder()
.setFirehoseClientProperties(sinkProperties) // Required
.setSerializationSchema(new SimpleStringSchema()) // Required
.setDeliveryStreamName("your-stream-name") // Required
.setFailOnError(false) // Optional
.setMaxBatchSize(500) // Optional
.setMaxInFlightRequests(50) // Optional
.setMaxBufferedRequests(10_000) // Optional
.setMaxBatchSizeInBytes(4 * 1024 * 1024) // Optional
.setMaxTimeInBufferMS(5000) // Optional
.setMaxRecordSizeInBytes(1000 * 1024) // Optional
.build()
flinkStream.sinkTo(kdfSink)
Python
sink_properties = {
# Required
'aws.region': 'eu-west-1',
# Optional, provide via alternative routes e.g. environment variables
'aws.credentials.provider.basic.accesskeyid': 'aws_access_key_id',
'aws.credentials.provider.basic.secretkey': 'aws_secret_access_key'
}
kdf_sink = KinesisFirehoseSink.builder() \
.set_firehose_client_properties(sink_properties) \ # Required
.set_serialization_schema(SimpleStringSchema()) \ # Required
.set_delivery_stream_name('your-stream-name') \ # Required
.set_fail_on_error(False) \ # Optional
.set_max_batch_size(500) \ # Optional
.set_max_in_flight_requests(50) \ # Optional
.set_max_buffered_requests(10000) \ # Optional
.set_max_batch_size_in_bytes(5 * 1024 * 1024) \ # Optional
.set_max_time_in_buffer_ms(5000) \ # Optional
.set_max_record_size_in_bytes(1 * 1024 * 1024) \ # Optional
.build()
구성
Flink의 Firehose sink는 정적 빌더 KinesisFirehoseSink.<InputType>builder()로 생성합니다.
- setFirehoseClientProperties(Properties sinkProperties)
- 필수.
- Firehose 클라이언트에 자격 증명, 리전 및 기타 매개변수를 제공합니다.
- setSerializationSchema(SerializationSchema serializationSchema)
- 필수.
- Sink에 직렬화 스키마를 제공합니다. 이 스키마는 Firehose로 보내기 전에 요소를 직렬화하는 데 사용됩니다.
- setDeliveryStreamName(String deliveryStreamName)
- 필수.
- 싱크할 delivery stream의 이름입니다.
- setFailOnError(boolean failOnError)
- 선택. 기본값:
false. - Firehose에 레코드 쓰기 실패 요청을 sink에서 치명적 예외로 처리할지 여부입니다.
- 선택. 기본값:
- setMaxBatchSize(int maxBatchSize)
- 선택. 기본값:
500. - Firehose에 쓸 배치의 최대 크기입니다.
- 선택. 기본값:
- setMaxInFlightRequests(int maxInFlightRequests)
- 선택. 기본값:
50. - sink가 백프레셔를 적용하기 전에 허용되는 진행 중(in flight) 요청의 최대 수입니다.
- 선택. 기본값:
- setMaxBufferedRequests(int maxBufferedRequests)
- 선택. 기본값:
10_000. - 백프레셔가 적용되기 전에 sink에 버퍼링될 수 있는 레코드의 최대 수입니다.
- 선택. 기본값:
- setMaxBatchSizeInBytes(int maxBatchSizeInBytes)
- 선택. 기본값:
4 * 1024 * 1024. - 배치가 될 수 있는 최대 크기(바이트)입니다. 전송되는 모든 배치는 이 크기보다 작거나 같습니다.
- 선택. 기본값:
- setMaxTimeInBufferMS(int maxTimeInBufferMS)
- 선택. 기본값:
5000. - 레코드가 플러시되기 전에 sink에 머무를 수 있는 최대 시간입니다.
- 선택. 기본값:
- setMaxRecordSizeInBytes(int maxRecordSizeInBytes)
- 선택. 기본값:
1000 * 1024. - sink가 수락할 최대 레코드 크기입니다. 이보다 큰 레코드는 자동으로 거부됩니다.
- 선택. 기본값:
- build()
- Firehose sink를 구성하고 반환합니다.
사용자 지정 Firehose 엔드포인트 사용
때로 Flink가 Firehose VPC 엔드포인트나 Localstack 같은 비-AWS Firehose 엔드포인트에 대해 컨슈머 또는 프로듀서로 동작하게 하는 것이 바람직합니다. 특히 Flink 애플리케이션의 기능 테스트를 수행할 때 유용합니다. 일반적으로 Flink 구성에 설정된 AWS 리전으로 추론되는 AWS 엔드포인트는 구성 속성을 통해 재정의해야 합니다.
AWS 엔드포인트를 재정의하려면 AWSConfigConstants.AWS_ENDPOINT와 AWSConfigConstants.AWS_REGION 속성을 설정하세요. 리전은 엔드포인트 URL에 서명하는 데 사용됩니다.
Java
Properties producerConfig = new Properties();
producerConfig.put(AWSConfigConstants.AWS_REGION, "us-east-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, "us-east-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")
Python
producer_config = {
'aws.region': 'us-east-1',
'aws.credentials.provider.basic.accesskeyid': 'aws_access_key_id',
'aws.credentials.provider.basic.secretkey': 'aws_secret_access_key',
'aws.endpoint': 'http://localhost:4566'
}