다양한 서비스가 결합된 대규모 플랫폼에서는 사용자들이 원하는 정보를 빠르고 정확하게 찾을 수 있도록 지원하는 검색 기능이 매우 중요합니다. 저희 팀은 여러 서비스들이 공통으로 사용할 수 있는 검색 플랫폼을 만들고, 다양한 서비스의 검색 기능을 안정적으로 지원하는 역할을 담당하고 있습니다.
이를 위해서는 가장 먼저 각 서비스의 데이터를 검색 엔진으로 가져와 색인(Indexing)하는 과정이 필요한데요, 이 글에서는 기존의 색인 파이프라인(Indexing Pipeline)을 운영하면서 마주했던 문제점들과, 이를 해결하기 위해 어떻게 파이프라인을 개선해 나갔는지 그 여정을 공유하고자 합니다.
저희가 운영하던 색인 파이프라인은 크게 두 가지 방식으로 동작했습니다. 첫째는 데이터 변경 사항을 실시간으로 반영하기 위한 온라인 색인(Online Indexing)이고, 둘째는 데이터 보정이나 사전 데이터 변경 등을 위해 하루에 한 번 전체 데이터를 다시 색인하는 오프라인 색인(Offline Indexing)입니다. 이 두 가지 파이프라인을 운영하면서 다음과 같은 문제점들에 직면하게 되었습니다.
위에서 언급한 4가지 문제를 해결하기 위해, 우리는 다음과 같은 명확한 목표를 세웠습니다.
생산성 높은 시스템을 만들기 위해서는, 반복적인 작업을 자동화하고 개발자가 핵심 비즈니스 로직에만 집중할 수 있도록 해야 합니다. 이를 위해 우리는 다음과 같은 원칙을 세우고, 설정 기반의 인터페이스를 설계했습니다.
원칙:
YAML 기반 설정 인터페이스:
프로젝트 내부 설정 형식 예시 — 아래 필드는 설명용 내부 스키마이며, Kafka나 Airflow가 그대로 인식하는 표준 설정이 아닙니다. 런타임에서 필드·타입·조건식과 참조하는 변환 함수를 검증해야 합니다.
# 예시: 데이터 소스 및 색인 관련 설정을 정의하는 YAML 파일
version: 1
name: my_service_datastore
storage: datastore
sources:
- table: my_service_table_v1
primary_key: id
schema: dbschema/my_service_table.yaml # 데이터 스키마 파일 경로
scan:
type: time
target_field: updated_at # 증분 색인을 위한 시간 기준 필드
ttl_interval_ms: 86400000 # 데이터 보존 기간 (예: 1일)
subscribe_message:
- topic: event.my_service.data_changed # 실시간 변경 이벤트를 수신할 Kafka 토픽
transform: my_custom_transform # 해당 이벤트를 처리할 변환 로직 이름
auto_indexing:
- option: simple
active_conditions: # 검색 엔진에 색인될 데이터의 조건
- deleted_at is null
# ... (배포, 리소스, 변환 파라미터 등 기타 설정)
자동 생성 모델의 구조 예시
// Code generated by indexer/message. DO NOT EDIT.
// source: message/schema/my_service_datastore.yaml
package my_service_datastore
import (
"time"
)
type MyServiceTableV1 struct {
ID int64 `db:"id" json:"id"`
Title *string `db:"title" json:"title"`
Content *string `db:"content" json:"content"`
UpdatedAt *time.Time `db:"updated_at" json:"updated_at"`
// ...
}
이러한 설정 기반 자동화 덕분에, 새로운 서비스를 검색 플랫폼에 연동하는 작업이 매우 간소화되었고, Airflow와 같은 워크플로우 관리 도구에 오프라인 색인 파이프라인(DAG)이 자동으로 생성되도록 연동하여 생산성을 극대화할 수 있었습니다.
운영 DB(Online DB)는 실시간으로 상태가 변하고, 직접적인 접근은 장애 포인트를 증가시킵니다. 또한, 매일 풀 스캔을 수행하는 것은 비용과 부하 측면에서 매우 비효율적입니다. 우리는 이 문제를 해결하기 위해 오프라인 스토리지(Offline Storage)를 중간 저장소로 활용하는 아키텍처를 도입했습니다.
updated_at만으로 판단하면 시계 차이·동일 시각 갱신을 구분하지 못할 수 있습니다.이러한 구조를 통해 배치 색인기의 운영 DB 전체 조회 의존성을 줄이고, 안정적이고 정제된 오프라인 데이터를 기반으로 색인 작업을 수행함으로써 파이프라인 전체의 안정성을 크게 향상시킬 수 있었습니다.
이 구조에도 최신 상태를 다시 조회하는 수집기는 운영 DB에 의존합니다. 요청을 모아 읽고 조회 동시성을 제한하는 등 부하 관리가 필요합니다. 삭제된 행은 재조회 결과만으로 삭제 시점과 버전을 알 수 없으므로 삭제 이벤트나 tombstone을 보존해야 합니다.
초기 스냅샷과 실시간 변경을 연결할 때는 스냅샷 기준점과 스트림 오프셋을 함께 기록합니다. 스냅샷 이후 변경을 재생하고, 새 인덱스의 건수·표본·삭제 반영을 확인한 뒤 읽기 대상을 전환합니다. 단순히 배치가 끝난 시각부터 이벤트를 읽으면 수집 중 발생한 변경을 놓칠 수 있습니다.
하루에 한 번, 수억 건의 데이터를 오프라인 스토리지에서 가져와 풀 색인하는 과정은 여전히 비용과 시간 측면에서 부담이었습니다. BigQuery의 읽기 비용은 조회 컬럼과 파티션·클러스터 프루닝, 요금 모델 등에 영향을 받습니다. 파티션 컬럼을 활용하지 못하는 필터는 많은 데이터를 읽을 수 있지만, 파일 기반이라는 이유만으로 모든 범위 조회가 전체 스캔을 하는 것은 아닙니다. BigQuery 파티션 조회
작업을 나누면 병렬성은 높일 수 있지만, 데이터를 나누어 복사하는 전처리 자체의 스캔·저장 비용도 계산해야 합니다. BigQuery의 네이티브 시간·정수 범위 파티셔닝과 임시 테이블을 만드는 작업 분할은 구분합니다. 해시 조건으로 여러 번 읽는다면 각 작업이 같은 원본을 반복 스캔하는지도 확인해야 합니다.
초기 기록의 처리 시간 1~2시간과 수만 원 수준 비용은 당시 작업에 대한 설명입니다. 데이터 크기, 가격·요금 모델, 전처리 포함 여부가 정리되지 않았으므로 일반적인 비용 절감 효과로 확장하지 않습니다. 실행 계획과 처리 바이트, 전처리부터 색인 완료까지의 시간을 함께 비교해야 합니다.
실시간 색인을 위해 사용하는 Kafka에는 특정 서비스로부터 초당 수천 건 이상의 이벤트가 쏟아져 들어오는 경우가 있습니다. 만약 이러한 이벤트를 하나씩 개별적으로 처리한다면, Kafka Lag가 발생하거나 DB 및 검색 엔진에 급격한 부하가 발생하여 다른 서비스에까지 영향을 미칠 수 있습니다.
다음은 특정 Go 라이브러리의 실행 코드가 아닌 처리 순서를 나타낸 의사 코드입니다.
이벤트를 제한된 크기의 버퍼에 추가
시간 또는 건수·바이트 상한에 도달하면 배치 생성
문서 ID별 버전과 삭제 정보를 확인하며 변환
검색 엔진 Bulk 요청 전송
응답의 각 항목을 검사
일시 실패: 해당 항목만 제한적으로 재시도
영구 실패: 실패 저장소에 기록하고 운영 정책에 따라 처리
완료가 확인된 연속 구간까지 파티션별 오프셋 커밋
Bulk 요청의 HTTP 응답이 성공이어도 개별 항목은 실패할 수 있습니다. 항목별 결과를 검사하지 않고 오프셋을 커밋하면 재처리 기회를 잃습니다. 반대로 색인 성공 직후 커밋 전에 장애가 나면 중복 처리될 수 있으므로 문서 ID와 버전 기준으로 멱등성을 확보해야 합니다. Elasticsearch Bulk API
배치 처리는 호출 오버헤드를 줄이지만 고가용성 자체를 보장하지는 않습니다. 버퍼 상한, 소비 지연, 재시도 대기열, 리밸런싱 복구와 검색 엔진 장애를 함께 검증해야 합니다. 과거 기록에 등장한 Pod 수와 처리량 역시 메시지 크기·검색 부하·내구성 조건이 없으면 용량 산정 기준으로 사용할 수 없습니다.
지금까지 검색 플랫폼의 색인 파이프라인이 가진 문제점들을 해결하고, 안정성과 생산성을 높이기 위해 고민하고 개선했던 과정들을 공유했습니다.
서비스가 성장할수록 더 많은 데이터를 더 효율적으로 처리할 수 있는 구조는 필수적입니다. 저희는 여기서 멈추지 않고, Transform 로직 자동 생성, 벡터 임베딩 및 LLM을 활용한 시맨틱 검색, 모델 서빙(Inference) 파이프라인 통합, 지속적인 테스트 및 모니터링 강화 등 더욱 편리하고 신뢰할 수 있는 검색 환경을 제공하기 위해 다양한 시도를 계속하고 있습니다.
이 글이 대규모 데이터 처리와 검색 시스템 아키텍처에 대해 고민하는 분들께 작은 영감이나마 드릴 수 있었기를 바랍니다. 긴 글 읽어주셔서 감사합니다.