Kinesis Source 커넥터
Kinesis Source 커넥터
Kinesis source 커넥터는 Amazon Kinesis에서 데이터를 가져와 Pulsar에 저장하는 커넥터예요. 이 커넥터는 Kinesis Consumer Library(KCL)로 실제 메시지 소비를 수행해요. KCL은 컨슈머 상태 추적에 DynamoDB를 사용해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
참고: 현재 Kinesis source 커넥터는 raw 메시지만 지원해요. KMS로 암호화된 메시지를 사용하면, 암호화된 메시지가 다운스트림으로 전송돼요. 이 커넥터는 향후 릴리스에서 메시지 복호화를 지원할 예정이에요.
출처: 문서
본문
구성 (Configuration)
Kinesis source 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (Property)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| initialPositionInStream | InitialPositionInStream | false | LATEST | 커넥터가 시작하는 위치예요. 사용 가능한 옵션은 다음과 같아요. AT_TIMESTAMP: 지정한 타임스탬프 또는 이후의 레코드에서 시작해요. LATEST: 가장 최근 데이터 레코드 이후에서 시작해요. TRIM_HORIZON: 사용 가능한 가장 오래된 데이터 레코드에서 시작해요. |
| startAtTime | Date | false | " " (빈 문자열) | AT_TIMESTAMP로 설정하면 소비를 시작할 시점을 지정해요. |
| applicationName | String | false | Pulsar IO connector | Amazon Kinesis 애플리케이션의 이름이에요. 기본적으로 애플리케이션 이름은 AWS 요청에 사용되는 user agent 문자열에 포함돼요. 이는 예를 들어 서로 다른 커넥터 인스턴스가 보낸 요청을 구분하는 등 문제 해결에 도움이 돼요. |
| checkpointInterval | long | false | 60000 | Kinesis 스트림 checkpoint의 빈도(밀리초)예요. |
| backoffTime | long | false | 3000 | AWS Kinesis로부터 throttling 예외를 만났을 때 요청 사이에 지연할 시간(밀리초)이에요. |
| numRetries | int | false | 3 | checkpoint 설정을 시도할 때 예외를 만나면 재시도하는 횟수예요. |
| receiveQueueSize | int | false | 1000 | 커넥터 내부에 버퍼링할 수 있는 최대 AWS 레코드 수예요. receiveQueueSize에 도달하면, 큐의 일부 메시지가 성공적으로 소비될 때까지 커넥터는 Kinesis에서 메시지를 소비하지 않아요. |
| dynamoEndpoint | String | false | " " (빈 문자열) | Dynamo 엔드포인트 URL이에요. 여기에서 찾을 수 있어요. |
| cloudwatchEndpoint | String | false | " " (빈 문자열) | Cloudwatch 엔드포인트 URL이에요. 여기에서 찾을 수 있어요. |
| useEnhancedFanOut | boolean | false | true | true로 설정하면 Kinesis enhanced fan-out을 사용해요. false로 설정하면 폴링(polling)을 사용해요. |
| awsEndpoint | String | false | " " (빈 문자열) | Kinesis 엔드포인트 URL이에요. 여기에서 찾을 수 있어요. |
| awsRegion | String | false | " " (빈 문자열) | AWS 리전이에요. 예: us-west-1, us-west-2 |
| awsKinesisStreamName | String | true | " " (빈 문자열) | Kinesis 스트림 이름이에요. |
| awsCredentialPluginName | String | false | " " (빈 문자열) | AwsCredentialProviderPlugin 구현의 완전한 클래스 이름이에요. awsCredentialProviderPlugin에는 다음과 같은 내장 플러그인이 있어요.org.apache.pulsar.io.kinesis.AwsDefaultProviderChainPlugin: 기본 AWS 프로바이더 체인을 사용하는 플러그인이에요. 더 자세한 내용은 default credential provider chain 사용하기를 참고하세요.org.apache.pulsar.io.kinesis.STSAssumeRoleProviderPlugin: KCL 실행 시 맡을(assume) 역할을 설명하는 구성을 awsCredentialPluginParam을 통해 받는 플러그인이에요. JSON 구성 예시: {"roleArn": "arn...", "roleSessionName": "name"}. awsCredentialPluginName은 Kinesis sink가 사용하는 AWSCredentialsProvider를 만드는 팩토리 클래스예요. awsCredentialPluginName이 비어 있으면, Kinesis sink는 awsCredentialPluginParam의 자격 증명 json-map을 받는 기본 AWSCredentialsProvider를 만들어요. |
| awsCredentialPluginParam | String | false | " " (빈 문자열) | awsCredentialsProviderPlugin을 초기화하는 JSON 파라미터예요. |
예제 (Example)
Kinesis source 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
- JSON
{
"configs": {
"awsEndpoint": "https://some.endpoint.aws",
"awsRegion": "us-east-1",
"awsKinesisStreamName": "my-stream",
"awsCredentialPluginParam": "{\"accessKey\":\"myKey\",\"secretKey\":\"my-Secret\"}",
"applicationName": "My test application",
"checkpointInterval": "30000",
"backoffTime": "4000",
"numRetries": "3",
"receiveQueueSize": 2000,
"initialPositionInStream": "TRIM_HORIZON",
"startAtTime": "2019-03-05T19:28:58.000Z"
}
}
- YAML
configs:
awsEndpoint: "https://some.endpoint.aws"
awsRegion: "us-east-1"
awsKinesisStreamName: "my-stream"
awsCredentialPluginParam: "{\"accessKey\":\"myKey\",\"secretKey\":\"my-Secret\"}"
applicationName: "My test application"
checkpointInterval: 30000
backoffTime: 4000
numRetries: 3
receiveQueueSize: 2000
initialPositionInStream: "TRIM_HORIZON"
startAtTime: "2019-03-05T19:28:58.000Z"