풀 기반 수집 API
풀 기반 수집 API (Pull-based Ingestion API)
3.0에서 도입되었어요. 풀 기반 수집(pull-based ingestion)은 OpenSearch가 Apache Kafka나 Amazon Kinesis 같은 스트리밍 소스에서 데이터를 수집할 수 있게 해줘요. 클라이언트가 REST API를 통해 OpenSearch에 데이터를 적극적으로 밀어 넣는 전통적인 수집 방식과 달리, 풀 기반 수집은 OpenSearch가 스트리밍 소스에서 직접 데이터를 가져와 데이터 흐름을 제어해요. 이 방식은 네이티브 백프레셔(backpressure) 처리를 제공해 트래픽 급증 시 서버 과부하를 막는 데 도움이 돼요. 풀 기반 수집은 최소 한 번(at-least-once) 수집 의미론을 보장하고, 외부 버전 관리(external versioning)를 사용해 데이터 일관성을 유지해요.
출처: 문서
본문
사전 요구 사항 (Prerequisites)
풀 기반 수집을 사용하기 전에 다음 사전 요구 사항이 충족되었는지 확인하세요.
bin/opensearch-plugin install <plugin-name>명령으로 스트리밍 소스용 ingestion plugin을 설치해요. 자세한 내용은 Additional plugins를 참고하세요. OpenSearch가 지원하는 ingestion plugin은 다음과 같아요.
- ingestion-kafka
- ingestion-kinesis (실험적)- 인덱스 생성 중에 풀 기반 수집을 구성해요. 기존의 푸시 기반(push-based) 인덱스를 풀 기반 인덱스로 전환할 수는 없어요.
풀 기반 수집용 인덱스 만들기 (Creating an index for pull-based ingestion)
스트리밍 소스에서 데이터를 수집하려면 먼저 풀 기반 수집 설정으로 인덱스를 만들어요. 다음 요청은 세그먼트 복제(segment replication) 모드로 Kafka 토픽에서 데이터를 가져오는 인덱스를 만들어요. 다른 사용 가능한 모드는 Ingestion modes를 참고하세요.
PUT /my-index
{
"settings": {
"ingestion_source": {
"type": "kafka",
"pointer.init.reset": "earliest",
"param": {
"topic": "test",
"bootstrap_servers": "localhost:49353"
}
},
"index.number_of_shards": 1,
"index.number_of_replicas": 1,
"index": {
"replication.type": "SEGMENT"
}
},
"mappings": {
"properties": {
"name": {
"type": "text"
},
"age": {
"type": "integer"
}
}
}
}
수집 소스 설정 (Ingestion source settings)
ingestion_source 설정은 OpenSearch가 스트리밍 소스에서 데이터를 가져오는 방식을 제어해요. 폴(poll)은 OpenSearch가 스트리밍 소스에 데이터 배치를 적극적으로 요청하는 작업이에요. 다음 표는 ingestion_source가 지원하는 모든 설정을 보여줘요.
동적(dynamic) 설정은 수집 과정을 재시작하지 않고 Update Settings API로 업데이트할 수 있어요. 정적(static) 설정은 인덱스 생성 후에는 변경할 수 없어요. 정적·동적 설정에 대한 자세한 내용은 Configuring OpenSearch를 참고하세요.
| 설정 | 동적 | 설명 |
|---|---|---|
| type | 아니요 | 스트리밍 소스 유형이에요. 필수예요. 유효한 값은 kafka 또는 kinesis예요. |
| pointer.init.reset | 아니요 | 읽기를 시작할 스트림 위치를 결정해요. 선택 사항이에요. 유효한 값은 earliest, latest, reset_by_offset, reset_by_timestamp, 또는 none이에요. Stream position을 참고하세요. |
| pointer.init.reset.value | 아니요 | reset_by_offset 또는 reset_by_timestamp에서만 필요해요. 오프셋 값이나 밀리초 단위의 타임스탬프를 지정해요. Stream position을 참고하세요. |
| error_strategy | 예 | 실패한 메시지를 처리하는 방법이에요. 선택 사항이에요. 유효한 값은 DROP(실패한 메시지를 건너뛰고 수집을 계속)과 BLOCK(메시지가 실패하면 수집을 중지)이에요. 기본값은 DROP이에요. |
| poll.max_batch_size | 예 | 각 폴 작업에서 가져올 최대 레코드 수예요. 선택 사항이에요. |
| poll.timeout | 예 | 각 폴 작업에서 데이터를 기다리는 최대 시간이에요. 선택 사항이에요. |
| num_processor_threads | 아니요 | 수집된 데이터를 처리하는 스레드 수예요. 선택 사항이에요. 기본값은 1이에요. |
| internal_queue_size | 아니요 | 고급 튜닝을 위한 내부 블로킹 큐의 크기예요. 유효한 값은 1부터 100,000까지(포함)예요. 선택 사항이에요. 기본값은 100이에요. |
| all_active | 아니요 | all-active 수집 모드를 활성화할지 여부예요. 세그먼트 복제 모드를 사용하는 인덱스에서는 활성화할 수 없어요. 기본값은 false예요. Ingestion modes를 참고하세요. |
| pointer_based_lag_update_interval | 아니요 | 포인터 기반 지연(lag)이 계산되는 간격이에요. 시간 단위를 받아요. 기본값은 10s예요. 이 값을 0으로 설정하면 포인터 기반 지연 계산을 비활성화해요. |
| warmup.timeout | 예 | 노드 재시작이나 샤드 재배치 후 워밍업(warmup) 단계에서 샤드가 스트리밍 소스를 따라잡기를 기다리는 최대 시간이에요. 워밍업이 완료되거나 제한 시간이 끝날 때까지 샤드는 쿼리를 제공하지 않아요. 시간 단위를 받아요. 선택 사항이에요. 기본값은 -1(비활성화)이에요. |
| warmup.lag_threshold | 예 | 워밍업 완료를 위한 허용 가능한 포인터 기반 지연 임계값이에요. 지연이 이 값 이하가 되면 워밍업이 완료돼요. 값 0은 샤드가 소스와 동기화되었음을 의미해요. 선택 사항이에요. 기본값은 100이에요. |
| mapper_type | 아니요 | 입력 메시지 형식에 대한 매퍼(mapper)를 정의해요. 유효한 값은 default와 raw_payload예요. Message format을 참고하세요. |
| param | 예 | 소스별 구성 파라미터예요. 필수예요. - ingest-kafka plugin은 다음을 요구해요. - topic : 소비할 Kafka 토픽 - bootstrap_servers : Kafka 서버 주소 선택적으로 fetch.min.bytes 같은 추가 표준 Kafka 소비자 파라미터를 제공할 수 있어요. 이 파라미터들은 Kafka 소비자에게 직접 전달돼요. - ingest-kinesis plugin은 다음을 요구해요. - stream : Kinesis 스트림 이름 - region : AWS 리전 - access_key : AWS 액세스 키 - secret_key : AWS 시크릿 키 선택적으로 endpoint_override를 제공할 수 있어요. |
기타 설정 (Other settings)
풀 기반 수집은 다음 OpenSearch 설정을 지원해요.
| 설정 | 동적 | 설명 |
|---|---|---|
| index.periodic_flush_interval | 예 | OpenSearch가 flush 작업을 트리거하는 간격이에요. 풀 기반 수집 인덱스의 기본값은 10m이에요. Index settings를 참고하세요. |
수집 모드 (Ingestion modes)
풀 기반 수집은 다음 모드를 지원해요.
세그먼트 복제 모드 (Segment replication mode)
세그먼트 복제 모드에서 프라이머리 샤드는 스트리밍 소스에서 이벤트를 수집하고 문서를 인덱싱해요. 풀 기반 인덱스는 세그먼트 파일을 프라이머리에서 복제본 샤드로 복사하기 위해 세그먼트 복제를 사용하도록 구성돼요.
이 모드는 원격 백업 스토리지(remote-backed storage)와 함께 사용하는 것을 권장해요.
All-active 모드 (All-active mode)
all-active 모드를 활성화하면 프라이머리와 복제본 샤드가 모두 스트리밍 소스에서 이벤트를 독립적으로 수집하고 인덱싱해요.
샤드 사이에 복제나 조정은 없어요. 다만 부트스트래핑 중 로컬 복사본을 사용할 수 없으면 복제본 샤드가 프라이머리 샤드에서 세그먼트 파일을 가져올 수 있어요. 이 모드는 현재 세그먼트 복제와 함께 지원되지 않아요.
스트림 위치 (Stream position)
인덱스를 만들 때 ingestion_source 파라미터의 pointer.init.reset과 pointer.init.reset.value 설정을 구성해서 OpenSearch가 스트림에서 읽기를 시작할 위치를 지정할 수 있어요. 기존 인덱스의 경우 OpenSearch는 마지막으로 커밋된 위치부터 읽기를 재개해요.
다음 표는 유효한 pointer.init.reset 값과 그에 해당하는 pointer.init.reset.value 값을 제공해요.
| pointer.init.reset | 시작 수집 지점 | pointer.init.reset.value |
|---|---|---|
| earliest | 스트림의 시작 | 없음 |
| latest | 스트림의 현재 끝 | 없음 |
| reset_by_offset | 스트림의 특정 오프셋 | 양의 정수 오프셋. 필수예요. |
| reset_by_timestamp | 특정 타임스탬프 | 밀리초 단위의 Unix 타임스탬프. 필수예요. Kafka 스트림의 경우 주어진 타임스탬프에 대한 메시지가 없으면 Kafka의 auto.offset.reset 정책을 기본값으로 사용해요. |
| none | 기존 인덱스의 마지막 커밋 위치 | 없음 |
스트림 파티셔닝 (Stream partitioning)
파티셔닝된 스트림(예: Kafka 토픽이나 Kinesis 샤드)을 사용할 때는 스트림 파티션과 OpenSearch 샤드 사이의 다음 관계에 유의하세요.
- OpenSearch 샤드는 스트림 파티션에 일대일로 매핑돼요.
- 인덱스 샤드 수는 스트림 파티션 수보다 크거나 같아야 해요.
- 파티션 수를 넘는 추가 샤드는 비어 있는 상태로 남아요.
- 성공적인 업데이트를 위해 문서는 같은 파티션으로 보내야 해요.
풀 기반 수집을 사용하면 인덱스에 대한 전통적인 REST API 기반 수집은 비활성화돼요.
오류 정책 업데이트 (Updating the error policy)
Update Settings API로 index.ingestion_source.error_strategy를 DROP 또는 BLOCK으로 설정해서 오류 정책을 동적으로 업데이트할 수 있어요.
다음 예제는 오류 정책을 업데이트하는 방법을 보여줘요.
메시지 형식 (Message format)
OpenSearch가 올바르게 처리하려면 스트리밍 소스의 메시지는 다음 형식이어야 해요.
{"_id":"1", "_version":"1", "_source":{"name": "alice", "age": 30}, "_op_type": "index"}
{"_id":"2", "_version":"2", "_source":{"name": "alice", "age": 30}, "_op_type": "delete"}
스트리밍 소스(Kafka 메시지 또는 Kinesis 레코드)의 각 데이터 단위는 OpenSearch 문서를 만들거나 수정하는 방법을 지정하는 다음 필드를 포함해야 해요. 이것은 풀 기반 수집이 지원하는 기본 형식이에요.
| 필드 | 데이터 타입 | 필수 | 설명 |
|---|---|---|---|
| _id | String | 아니요 | 문서의 고유 식별자예요. 제공하지 않으면 OpenSearch가 ID를 자동 생성해요. 문서 업데이트나 삭제에는 필수예요. |
| _version | Long | 아니요 | 문서 버전 번호로, 외부에서 유지 관리해야 해요. 제공하면 OpenSearch는 현재 문서 버전보다 이전 버전의 메시지를 버려요. 제공하지 않으면 버전 확인이 일어나지 않아요. |
| _op_type | String | 아니요 | 수행할 작업이에요. 유효한 값은 다음과 같아요. - index : 새 문서를 만들거나 기존 문서를 업데이트해요. - create : append 모드로 새 문서를 만들어요. 기존 문서는 업데이트하지 않아요. - delete : 문서를 소프트 삭제해요. |
| _source | Object | 예 | 문서 데이터를 담고 있는 메시지 페이로드예요. |
풀 기반 수집은 최소 한 번 수집 의미론을 제공하므로 중복을 막기 위해 문서 _id 필드를 사용하는 것을 권장해요. 프로듀서가 이벤트 순서를 보장할 수 없다면 데이터 일관성을 위해 _version 필드도 설정하세요.
또는 풀 기반 수집은 변환 없이 append-only 모드로 원시 페이로드를 인덱싱하는 것을 지원해요. 이 동작을 활성화하려면 index.ingestion_source.mapper_type을 raw_payload로 설정하세요. 이 모드에서는 동적 매핑을 지원하지 않으므로 인덱스 매핑이 메시지 구조와 일치해야 해요. raw_payload를 사용할 때는 다음 예제처럼 들어오는 데이터 스트림에 나타나는 그대로 정확한 원시 JSON 객체를 제공해야 해요.
{"name": "alice", "age": 30}
{"name": "bob", "age": 30}
풀 기반 수집 메트릭 (Pull-based ingestion metrics)
풀 기반 수집은 수집 과정을 모니터링하는 데 사용할 수 있는 메트릭을 제공해요. 현재 polling_ingest_stats 메트릭이 지원되며 샤드 수준에서 사용할 수 있어요.
다음 표는 사용 가능한 polling_ingest_stats 메트릭을 보여줘요.
| 메트릭 | 설명 |
|---|---|
| message_processor_stats.total_processed_count | 메시지 프로세서가 처리한 총 메시지 수예요. |
| message_processor_stats.total_invalid_message_count | 만난 잘못된(invalid) 메시지 수예요. |
| message_processor_stats.total_version_conflicts_count | 버전 충돌로 인해 더 오래된 버전의 메시지가 버려질 횟수예요. |
| message_processor_stats.total_failed_count | 처리 중 오류가 발생한 실패 메시지의 총 수예요. |
| message_processor_stats.total_failures_dropped_count | 재시도를 모두 소진한 후 버려진 실패 메시지의 총 수예요. 메시지는 DROP 오류 정책을 사용할 때만 버려진다는 점에 유의하세요. |
| message_processor_stats.total_processor_thread_interrupt_count | 프로세서 스레드에서 발생한 스레드 인터럽트 수를 나타내요. |
| consumer_stats.total_polled_count | 스트림 소비자에서 폴링된 총 메시지 수예요. |
| consumer_stats.total_consumer_error_count | 치명적인 소비자 읽기 오류의 총 수예요. |
| consumer_stats.total_poller_message_failure_count | 폴러에서 실패한 메시지의 총 수예요. |
| consumer_stats.total_poller_message_dropped_count | 폴러에서 버려진 실패 메시지의 총 수예요. |
| consumer_stats.lag_in_millis | 밀리초 단위의 지연으로, 마지막으로 처리된 메시지 타임스탬프 이후 경과된 시간으로 계산돼요. |
| consumer_stats.pointer_based_lag | Apache Kafka 오프셋 기반 지연으로, 사용 가능한 최신 오프셋과 현재 메시지 오프셋의 차이로 계산돼요. 이 메트릭은 Apache Kafka를 스트리밍 소스로 사용할 때만 적용돼요. |
샤드 수준의 풀 기반 수집 메트릭을 가져오려면 Nodes Stats API를 사용해요.
제한 사항 (Limitations)
풀 기반 수집을 사용할 때 다음 제한 사항이 적용돼요.
- ingest pipeline은 풀 기반 수집과 호환되지 않아요.
- 동적 매핑은 지원되지 않아요.
- 인덱스 롤오버(rollover)는 지원되지 않아요.
- 작업 리스너(operation listener)는 지원되지 않아요.
관련 문서 (Related documentation)
- 풀 기반 수집 관리 (Pull-based ingestion management)