스트림 수집 플러그인
스트림 수집 플러그인 (Stream Ingestion Plugin)
Pinot용 자체 스트림 수집 플러그인을 작성하는 방법을 다루는 문서예요. Pravega, Kinesis 같은 자체 스트리밍 플랫폼 지원을 추가하기 위해 커스텀 스트림 수집 플러그인을 작성할 수 있어요.
출처: 문서
본문
이 페이지는 Pinot용 자체 스트림 수집 플러그인을 작성하는 방법을 설명해요.
Pravega, Kinesis 같은 자체 스트리밍 플랫폼 지원을 추가하기 위해 커스텀 스트림 수집 플러그인을 작성할 수 있어요.
Pinot는 각 스트림 파티션을 독립적으로 소비하며 파티션별 오프셋을 추적해요. 통합하려는 스트림은 따라서 파티션별 오프셋을 노출하고 주어진 오프셋부터 소비하는 것을 지원해야 해요.
스트림 지원 요구 사항 (Requirements to support a stream)
파티션 수준에서 행을 소비하는 동안 스트림은 다음 속성을 지원해야 해요:
- 스트림은 현재 파티션 수를 얻는 메커니즘을 제공해야 해요.
- 파티션의 각 이벤트는 64비트 이하 길이의 고유 오프셋을 가져야 해요.
- 파티션을 32비트 이하 길이의 숫자로 참조해야 해요.
- 스트림은 주어진 스트림 파티션의 오프셋을 얻는 다음 메커니즘을 제공해야 해요:
- 파티션에서 사용 가능한 가장 오래된 이벤트의 오프셋 얻기 (이벤트가 주기적으로 만료된다고 가정)
- 파티션에 게시된 가장 최근 이벤트의 오프셋 얻기
- (선택 사항) 지정된 시각에 게시된 이벤트의 오프셋 얻기
- 스트림은 지정된 오프셋에서 시작해 파티션의 이벤트 집합을 소비하는 메커니즘을 제공해야 해요.
- Pinot는 들어오는 이벤트의 오프셋이 단조 증가한다고 가정해요. 즉 Pinot가 오프셋
o1에서 이벤트를 소비하면 다음 이벤트의 오프셋o2는o2 > o1이어야 해요.
추가로, 시간이 지나도 파티션 수가 줄어들지 않아야 한다는 운영 요구 사항이 있어요.
스트림 플러그인 구현 (Stream plug-in implementation)
새 스트림 유형(예: Foo)을 추가하려면 다음 클래스를 구현해요:
- FooConsumerFactory는 StreamConsumerFactory를 확장
- FooPartitionGroupConsumer는 PartitionGroupConsumer를 구현
- FooMetadataProvider는 StreamMetadataProvider를 구현
- FooMessageDecoder는 StreamMessageDecoder를 구현
스트림 구현의 속성은 테이블 구성의 streamConfigs 섹션 안에 설정해요.
streamType 속성을 사용해 스트림 유형을 정의해요. 예를 들어 스트림 foo 구현의 경우 "streamType" : "foo" 속성을 설정해요.
스트림의 나머지 구성 속성은 "stream.foo" 접두사로 설정해야 해요. 다음에 대해 동일한 접미사를 사용해야 해요 (아래 예시 참고):
- topic
- 스트림 컨슈머 팩토리
- offset
- 디코더 클래스 이름
- 디코더 속성
- 연결 타임아웃
- fetch 타임아웃
모든 값은 문자열이어야 해요. 예를 들어:
"streamType" : "foo",
"stream.foo.topic.name" : "SomeTopic",
"stream.foo.consumer.factory.class.name": "fully.qualified.pkg.ConsumerFactoryClassName",
"stream.foo.consumer.prop.auto.offset.reset": "largest",
"stream.foo.decoder.class.name" : "fully.qualified.pkg.DecoderClassName",
"stream.foo.decoder.prop.a.decoder.property" : "decoderPropValue",
"stream.foo.connection.timeout.millis" : "10000", // default 30_000
"stream.foo.fetch.timeout.millis" : "10000" // default 5_000
스트림에 특화된 추가 속성을 가질 수도 있어요. 예를 들어:
"stream.foo.some.buffer.size" : "24g"
이 속성들에 더해, 소비 세그먼트의 임계값을 정의할 수 있어요:
- 행 임계값
- 시간 임계값
임계값 속성은 다음과 같아요:
"realtime.segment.flush.threshold.rows" : "100000"
"realtime.segment.flush.threshold.time" : "6h"
이 구현의 예시는 kafka 스트림 구현인 KafkaConsumerFactory에서 볼 수 있어요.