Kinesis Sink 커넥터
Kinesis Sink 커넥터
Kinesis sink 커넥터는 Pulsar에서 데이터를 가져와 Amazon Kinesis에 저장하는 커넥터예요. Pulsar의 데이터를 AWS Kinesis 스트림으로 전달할 때 사용해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
출처: 문서
본문
구성 (Configuration)
Kinesis sink 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (Property)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| messageFormat | MessageFormat | true | ONLY_RAW_PAYLOAD | Kinesis sink가 Pulsar 메시지를 변환해 Kinesis 스트림에 발행하는 메시지 형식이에요. 사용 가능한 옵션은 다음과 같아요. ONLY_RAW_PAYLOAD: Kinesis sink가 Pulsar 메시지 페이로드를 구성된 Kinesis 스트림에 메시지로 직접 발행해요. FULL_MESSAGE_IN_JSON: Kinesis sink가 Pulsar 메시지 페이로드, 프로퍼티, encryptionCtx로 JSON 페이로드를 만들어 구성된 Kinesis 스트림에 JSON 페이로드를 발행해요. FULL_MESSAGE_IN_FB: Kinesis sink가 Pulsar 메시지 페이로드, 프로퍼티, encryptionCtx로 flatbuffer 직렬화 페이로드를 만들어 구성된 Kinesis 스트림에 flatbuffer 페이로드를 발행해요. FULL_MESSAGE_IN_JSON_EXPAND_VALUE: Kinesis sink가 레코드 토픽 이름, 키, 페이로드, 프로퍼티, 이벤트 시간을 담은 JSON 구조를 보내요. 레코드 스키마가 값을 JSON으로 변환하는 데 사용돼요. |
| jsonIncludeNonNulls | boolean | false | true | 메시지 형식이 FULL_MESSAGE_IN_JSON_EXPAND_VALUE일 때 null이 아닌 값을 가진 프로퍼티만 포함돼요. |
| jsonFlatten | boolean | false | false | true로 설정하고 메시지 형식이 FULL_MESSAGE_IN_JSON_EXPAND_VALUE이면 출력 JSON이 평탄화(flatten)돼요. |
| retainOrdering | boolean | false | false | Pulsar 커넥터가 Pulsar에서 Kinesis로 메시지를 이동할 때 순서를 유지할지 여부예요. |
| awsEndpoint | String | false | " " (빈 문자열) | Kinesis 엔드포인트 URL이에요. 여기에서 찾을 수 있어요. |
| awsRegion | String | false | " " (빈 문자열) | AWS 리전이에요. 예: us-west-1, us-west-2 |
| awsKinesisStreamName | String | true | " " (빈 문자열) | Kinesis 스트림 이름이에요. |
| awsCredentialPluginName | String | false | " " (빈 문자열) | AwsCredentialProviderPlugin 구현의 완전한 클래스 이름이에요. Kinesis sink가 사용하는 AWSCredentialsProvider를 만드는 팩토리 클래스예요. 비어 있으면, Kinesis sink는 awsCredentialPluginParam의 자격 증명 json-map을 받는 기본 AWSCredentialsProvider를 만들어요. |
| awsCredentialPluginParam | String | false | " " (빈 문자열) | awsCredentialsProviderPlugin을 초기화하는 JSON 파라미터예요. |
내장 플러그인 (Built-in plugins)
다음은 내장 AwsCredentialProviderPlugin 플러그인들이에요.
org.apache.pulsar.io.aws.AwsDefaultProviderChainPlugin
이 플러그인은 구성이 필요 없으며 기본 AWS 프로바이더 체인을 사용해요. 더 자세한 내용은 AWS 문서를 참고하세요.
org.apache.pulsar.io.aws.STSAssumeRoleProviderPlugin
이 플러그인은 KCL 실행 시 맡을(assume) 역할을 설명하는 구성(awsCredentialPluginParam 통해)을 받아요. 이 구성은 다음과 같은 작은 JSON 문서의 형태예요.
{"roleArn": "arn...", "roleSessionName": "name"}
예제 (Example)
Kinesis sink 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
- JSON
{
"configs": {
"awsEndpoint": "some.endpoint.aws",
"awsRegion": "us-east-1",
"awsKinesisStreamName": "my-stream",
"awsCredentialPluginParam": "{\"accessKey\":\"myKey\",\"secretKey\":\"my-Secret\"}",
"messageFormat": "ONLY_RAW_PAYLOAD",
"retainOrdering": "true"
}
}
- YAML
configs:
awsEndpoint: "some.endpoint.aws"
awsRegion: "us-east-1"
awsKinesisStreamName: "my-stream"
awsCredentialPluginParam: "{\"accessKey\":\"myKey\",\"secretKey\":\"my-Secret\"}"
messageFormat: "ONLY_RAW_PAYLOAD"
retainOrdering: "true"