Apache Pulsar에서 수집
Apache Pulsar에서 수집 (Ingest from Apache Pulsar)
Apache Pulsar 토픽의 레코드 스트림을 Pinot 테이블로 수집하는 방법을 보여 주는 가이드예요.
본문
Pinot은 pinot-pulsar 플러그인을 통해 Apache Pulsar에서 데이터 소비를 지원해요. 클래스패스에 Pulsar 관련 라이브러리가 있도록 이 플러그인을 활성화해야 해요.
Pinot 설정 시 다음 설정으로 Pulsar 플러그인을 활성화해요: -Dplugins.include=pinot-pulsar
{% hint style="info" %}
pinot-pulsar 플러그인은 Pinot 0.11.0부터 공식 바이너리 배포판에 포함됩니다. 이전 버전을 실행 중이라면 Apache Pinot 외부 저장소에서 플러그인을 다운로드해 plugins 디렉터리에 추가할 수 있습니다.
{% endhint %}
Pulsar 테이블 설정
다음은 샘플 Pulsar 스트림 설정이에요. 이 샘플의 streamConfigs 섹션을 사용해 해당 테이블에 맞게 변경할 수 있어요.
{
"tableName": "pulsarTable",
"tableType": "REALTIME",
"segmentsConfig": {
"timeColumnName": "timestamp",
"replicasPerPartition": "1"
},
"tenants": {},
"tableIndexConfig": {
"loadMode": "MMAP",
"streamConfigs": {
"streamType": "pulsar",
"stream.pulsar.topic.name": "<your pulsar topic name>",
"stream.pulsar.bootstrap.servers": "pulsar://localhost:6650,pulsar://localhost:6651",
"stream.pulsar.consumer.prop.auto.offset.reset" : "smallest",
"stream.pulsar.fetch.timeout.millis": "30000",
"stream.pulsar.decoder.class.name": "org.apache.pinot.plugin.inputformat.json.JSONMessageDecoder",
"stream.pulsar.consumer.factory.class.name": "org.apache.pinot.plugin.stream.pulsar.PulsarConsumerFactory",
"realtime.segment.flush.threshold.rows": "1000000",
"realtime.segment.flush.threshold.time": "6h"
}
},
"metadata": {
"customConfigs": {}
}
}
Pulsar 설정 옵션
테이블에 대해 다음 Pulsar 특정 설정을 변경할 수 있어요:
| 속성 (Property) | 설명 (Description) |
|---|---|
streamType |
"pulsar"로 설정해야 함 |
stream.pulsar.topic.name |
Pulsar 토픽 이름 |
stream.pulsar.bootstrap.servers |
Apache Pulsar용 쉼표로 구분된 브로커 목록 |
stream.pulsar.metadata.populate |
메타데이터 채우려면 true로 설정 |
stream.pulsar.metadata.fields |
채울 메타데이터 필드의 쉼표로 구분된 목록으로 설정 |
인증 (Authentication)
Pinot-Pulsar 커넥터는 보안 토큰을 사용한 인증을 지원해요. 토큰 생성은 Pulsar 문서의 지침을 따르세요. 생성 후 각 요청에 인증 토큰을 추가하려면 streamConfigs에 다음 속성을 추가해요:
"stream.pulsar.authenticationToken":"your-auth-token"
OAuth2 인증
Pinot-Pulsar 커넥터는 예를 들어 StreamNative Pulsar 클러스터에 연결할 때 OAuth2를 사용한 인증을 지원해요. 자세한 내용은 Pulsar 클라이언트에서 OAuth2 인증 구성을 참고하세요. 구성 후 streamConfigs에 다음 속성을 추가할 수 있어요:
"stream.pulsar.issuerUrl": "https://auth.streamnative.cloud"
"stream.pulsar.credsFilePath": "file:///path/to/private_creds_file
"stream.pulsar.audience": "urn:sn:pulsar:test:test-cluster"
TLS 지원
Pinot-pulsar 커넥터는 암호화 연결을 위한 TLS도 지원해요. 공식 pulsar 문서를 따라 pulsar 클러스터에서 TLS를 활성화할 수 있어요. 완료 후 이전 단계에서 생성한 신뢰 인증서 파일 위치를 제공해 pulsar 커넥터에서 TLS를 활성화할 수 있어요.
"stream.pulsar.tlsTrustCertsFilePath": "/path/to/ca.cert.pem"
또한 브로커 url을 pulsar://localhost:6650에서 pulsar+ssl://localhost:6650으로 변경해 보안 연결이 사용되도록 하세요.
다른 테이블·스트림 설정은 테이블 설정 레퍼런스 (Table configuration Reference) 참고하세요.
지원되는 Pulsar 버전
Pinot은 현재 Pulsar 클라이언트 버전 4.0.x에 의존해요. Pulsar 브로커가 이 클라이언트 버전과 호환되는지 확인하세요.
레코드 헤더를 Pinot 테이블 컬럼으로 추출
Pinot의 Pulsar 커넥터는 레코드 헤더와 메타데이터를 Pinot 테이블 컬럼으로 자동 추출하는 것을 지원해요. Pulsar는 레코드당 많은 메타데이터를 지원해요. 메타데이터 필드의 의미는 공식 Pulsar 문서를 참고하세요.
다음 표는 레코드 헤더/메타데이터에서 Pinot 테이블 컬럼 이름으로의 매핑이에요:
| Pulsar 메시지 (Message) | Pinot 테이블 컬럼 (Column) | 주석 (Comments) | 기본 제공 (Available By Default) |
|---|---|---|---|
| key : String | __key : String |
예 | |
| properties : Map<String, String> | 각 헤더 키는 별도 컬럼으로: __header$HeaderKeyName : String |
예 | |
| publishTime : Long | __metadata$publishTime : String |
프로듀서가 결정한 publish time | 예 |
| brokerPublishTime: Optional | __metadata$brokerPublishTime : String |
브로커가 결정한 publish time | 예 |
| eventTime : Long | __metadata$eventTime : String |
예 | |
| messageId : MessageId -> String | __metadata$messageId : String |
MessagId 필드의 문자열 표현. 형식은 ledgerId:entryId:partitionIndex | |
| messageId : MessageId -> bytes | __metadata$messageBytes : String |
MessageId.toByteArray() 호출에서 반환된 바이트의 Base64 인코딩 버전 | |
| producerName : String | __metadata$producerName : String |
||
| schemaVersion : byte[] | __metadata$schemaVersion : String |
Base64 인코딩 값 | |
| sequenceId : Long | __metadata$sequenceId : String |
||
| orderingKey : byte[] | __metadata$orderingKey : String |
Base64 인코딩 값 | |
| size : Integer | __metadata$size : String |
||
| topicName : String | __metadata$topicName : String |
||
| index : String | __metadata$index : String |
||
| redeliveryCount : Integer | __metadata$redeliveryCount : String |
Pulsar 테이블에서 메타데이터 추출을 활성화하려면 스트림 설정 metadata.populate를 true로 설정해요. eventTime, publishTime, brokerPublishTime, key 필드는 기본적으로 채워져요. Pulsar Message에서 추가 필드를 추출하려면 metadataFields 설정을 채울 필드의 쉼표로 구분된 목록으로 채워요. 필드는 Pulsar Message의 필드 이름으로 참조돼요. 예를 들어 설정:
"streamConfigs": {
...
"stream.pulsar.metadata.populate": "true",
"stream.pulsar.metadata.fields": "messageId,messageIdBytes,eventTime,topicName",
...
}
그러면 __metadata$messageId, __metadata$messageBytes, __metadata$eventTime, __metadata$topicName 필드를 Pinot 스키마의 컬럼에 매핑할 수 있게 돼요.
이에 더해, 이 컬럼들 중 어떤 것이든 테이블에서 사용하려면 테이블의 스키마에 명시적으로 나열해야 해요.
예를 들어 Pinot 테이블에 offset과 key만 dimension 컬럼으로 추가하고 싶다면 스키마에 다음과 같이 나열할 수 있어요:
"dimensionFieldSpecs": [
{
"name": "__key",
"dataType": "STRING"
},
{
"name": "__metadata$messageId",
"dataType": "STRING"
},
...
],
스키마가 업데이트되면 이 컬럼들은 다른 pinot 컬럼과 유사해져요. 여기에 수집 변환을 적용하거나 인덱스를 정의할 수 있어요.
{% hint style="info" %} 기존 테이블의 스키마를 업데이트할 때는 스키마 진화 가이드라인을 따르는 것을 잊지 마세요! {% endhint %}