AWS DynamoDB Source 커넥터
AWS DynamoDB Source 커넥터
DynamoDB source 커넥터는 DynamoDB 테이블 스트림에서 데이터를 가져와 Pulsar에 저장하는 커넥터예요. 이 커넥터는 DynamoDB Streams Kinesis Adapter를 사용하는데, 실제 메시지 소비는 Kinesis Consumer Library(KCL)가 담당해요. KCL은 컨슈머 상태 추적에 DynamoDB를 사용하고, 지표 로깅을 위해 cloudwatch 접근이 필요해요.
참고: 모든 Pulsar 커넥터는 download page에서 내려받을 수 있어요.
출처: 문서
본문
구성 (Configuration)
DynamoDB source 커넥터의 구성에는 다음과 같은 프로퍼티가 있어요.
프로퍼티 (Property)
| 이름 | 타입 | 필수 | 기본값 | 설명 |
|---|---|---|---|---|
| initialPositionInStream | InitialPositionInStream | false | LATEST | 커넥터가 시작하는 위치예요. 사용 가능한 옵션은 다음과 같아요. AT_TIMESTAMP: 지정한 타임스탬프 또는 이후의 레코드에서 시작해요. LATEST: 가장 최근 데이터 레코드 이후에서 시작해요. TRIM_HORIZON: 사용 가능한 가장 오래된 데이터 레코드에서 시작해요. |
| startAtTime | Date | false | " " (빈 문자열) | AT_TIMESTAMP로 설정하면 소비를 시작할 시점을 지정해요. |
| applicationName | String | false | Pulsar IO connector | KCL 애플리케이션의 이름이에요. 상태 추적에 사용되는 dynamo 테이블의 테이블 이름을 정의하는 데 쓰이므로 고유해야 해요. 기본적으로 애플리케이션 이름은 AWS 요청에 사용되는 user agent 문자열에 포함돼요. 이는 예를 들어 서로 다른 커넥터 인스턴스가 보낸 요청을 구분하는 등 문제 해결에 도움이 돼요. |
| checkpointInterval | long | false | 60000 | KCL 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이에요. 여기에서 찾을 수 있어요. |
| awsEndpoint | String | false | " " (빈 문자열) | DynamoDB Streams 엔드포인트 URL이에요. 여기에서 찾을 수 있어요. |
| awsRegion | String | false | " " (빈 문자열) | AWS 리전이에요. 예: us-west-1, us-west-2 |
| awsDynamodbStreamArn | String | true | " " (빈 문자열) | DynamoDB 스트림 ARN이에요. |
| 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)
DynamoDB source 커넥터를 사용하기 전에 다음 방법 중 하나로 구성 파일을 만들어야 해요.
- JSON
{
"configs": {
"awsEndpoint": "https://some.endpoint.aws",
"awsRegion": "us-east-1",
"awsDynamodbStreamArn": "arn:aws:dynamodb:us-west-2:111122223333:table/TestTable/stream/2015-05-11T21:21:33.291",
"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"
awsDynamodbStreamArn: "arn:aws:dynamodb:us-west-2:111122223333:table/TestTable/stream/2015-05-11T21:21:33.291"
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"