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"

더 알아보기 (Learn more)