풀 기반 수집 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)는 지원되지 않아요.
  • 풀 기반 수집 관리 (Pull-based ingestion management)

더 알아보기 (Learn more)