본문 바로가기
WIKI 기술 지식 베이스

스트리밍 서비스 (Streaming Service)

원문 보기 위키 갱신

Streaming Service는 Milvus 내부 스트리밍 시스템 모듈의 개념으로, WAL(Write-Ahead Log)을 중심으로 다양한 스트리밍 관련 기능을 지원해요. 여기에는 스트리밍 데이터 인제스트/구독, 클러스터 상태의 장애 복구, 스트리밍 데이터의 과거 데이터로의 변환, 증가 데이터 쿼리 등이 포함돼요. 아키텍처적으로 Streaming Service는 세 가지 주요 구성 요소로 이루어져 있어요.

Streaming Distributed Arc

  • Streaming Coordinator: 코디네이터 노드의 논리적 구성 요소예요. Etcd를 사용한 서비스 디스커버리로 사용 가능한 스트리밍 노드를 찾고, WAL을 해당 스트리밍 노드에 바인딩하는 역할을 해요. 또한 WAL 분산 토폴로지를 노출하는 서비스를 등록해서 스트리밍 클라이언트가 주어진 WAL에 적절한 스트리밍 노드를 알 수 있게 해 줘요.
  • Streaming Node Cluster: 모든 스트리밍 처리 작업(wal appending, state recovering, growing data querying 등)을 담당하는 스트리밍 워커 노드의 클러스터예요.
  • Streaming Client: 서비스 디스커버리와 준비 상태 확인 같은 기본 기능을 캡슐화한, Milvus 내부 개발 클라이언트예요. 메시지 쓰기와 구독 같은 작업을 시작하는 데 사용돼요.

출처: Milvus 문서

본문

메시지 (Message)

Streaming Service는 로그 기반 스트리밍 시스템이므로 Milvus의 모든 쓰기 연산(DML, DDL 등)은 Message로 추상화돼요.

  • 모든 Message에는 Streaming Service가 Timestamp Oracle (TSO) 필드를 부여하며, 이는 WAL에서 메시지의 순서를 나타내요. 메시지의 순서는 Milvus에서 쓰기 연산의 순서를 결정해요. 이를 통해 최신 클러스터 상태를 로그에서 재구성할 수 있어요.
  • 각 Message는 특정 VChannel(가상 채널)에 속하며, 그 채널 안에서 연산 일관성을 보장하는 특정 불변 속성을 유지해요. 예를 들어 Insert 연산은 같은 채널의 DropCollection 연산보다 반드시 먼저 일어나야 해요.

Milvus의 메시지 순서는 다음과 비슷할 수 있어요.

Message Order

WAL 구성 요소 (WAL Component)

대규모 수평 확장을 지원하기 위해 Milvus의 WAL은 단일 로그 파일이 아니라 여러 로그의 합성물이에요. 각 로그는 여러 VChannel의 스트리밍 기능을 독립적으로 지원할 수 있어요. 주어진 시점에 WAL 구성 요소는 정확히 하나의 스트리밍 노드에서만 동작할 수 있으며, 이 제약은 기본 wal 저장소의 fencing 메커니즘과 스트리밍 코디네이터에 의해 보장돼요.

WAL 구성 요소의 추가 기능은 다음과 같아요.

  • 세그먼트 수명주기 관리: 메모리 상태/세그먼트 크기/세그먼트 유휴 시간 같은 정책을 기반으로 WAL이 모든 세그먼트의 수명주기를 관리해요.
  • 기본 트랜잭션 지원: 각 메시지에는 크기 제한이 있으므로, WAL 구성 요소는 VChannel 수준에서 원자적 쓰기를 보장하기 위한 간단한 트랜잭션 수준을 지원해요.
  • 고동시성 원격 로그 쓰기: Milvus는 서드파티 원격 메시지 큐를 WAL 저장소로 지원해요. 스트리밍 노드와 원격 WAL 저장소 사이의 왕복 지연(RTT)을 완화해 쓰기 처리량을 개선하기 위해 스트리밍 서비스는 동시 로그 쓰기를 지원해요. TSO와 TSO 동기화로 메시지 순서를 유지하며, WAL의 메시지는 TSO 순서대로 읽혀요.
  • Write-Ahead Buffer: 메시지가 WAL에 쓰여진 뒤 Write-Ahead Buffer에 임시 저장돼요. 이를 통해 원격 WAL 저장소에서 메시지를 가져오지 않고 로그의 꼬리를 읽을 수 있어요.
  • 여러 WAL 저장소 지원: Woodpecker, Pulsar, Kafka. zero-disk 모드로 woodpecker를 사용하면 원격 WAL 저장소 의존성을 제거할 수 있어요.

복구 저장소 (Recovery Storage)

Recovery Storage 구성 요소는 항상 해당 WAL 구성 요소가 위치한 스트리밍 노드에서 실행돼요.

  • 스트리밍 데이터를 영속화된 과거 데이터로 변환해 객체 저장소에 저장하는 역할을 해요.
  • 또한 스트리밍 노드의 WAL 구성 요소에 대한 인메모리 상태 복구를 처리해요.

Recovery Storage

쿼리 위임자 (Query Delegator)

Query Delegator는 각 스트리밍 노드에서 실행되며 단일 샤드에서 증분 쿼리(incremental queries) 를 실행하는 역할을 해요. 쿼리 계획을 생성해 관련 Query Node에 전달하고 결과를 집계해요.

또한 Query Delegator는 Delete 연산을 다른 Query Node에 브로드캐스트하는 역할을 해요.

Query Delegator는 항상 같은 스트리밍 노드의 WAL 구성 요소와 함께 존재해요. 하지만 컬렉션이 멀티레플리카로 구성되어 있다면 N-1개의 Delegator가 다른 스트리밍 노드에 배포돼요.

WAL 수명주기와 준비 대기 (WAL Lifetime and Wait for Ready)

컴퓨팅 노드를 저장소와 분리함으로써 Milvus는 WAL을 한 스트리밍 노드에서 다른 스트리밍 노드로 쉽게 옮길 수 있고, 이로써 스트리밍 서비스의 고가용성을 달성해요.

wal lifetime

준비 대기 (Wait for Ready)

WAL이 새 스트리밍 노드로 이동하려 할 때, 클라이언트는 이전 스트리밍 노드가 일부 요청을 거부하는 것을 발견해요. 그동안 WAL은 새 스트리밍 노드에서 복구되며, 클라이언트는 새 스트리밍 노드의 WAL이 서빙 준비가 될 때까지 기다려요.

wait for ready

더 알아보기 (Learn more)