Spark와 Milvus로 프로덕션 준비가 완료된 검색 파이프라인 구축하기
프로덕션에서 확장 가능한 벡터 검색 파이프라인을 구축하는 것은 프로토타입을 구축하는 것만큼 쉽지 않습니다. 프로토타입 작업을 할 때는 대개 소량의 비정형 데이터만 다룹니다. 그러나 프로토타입을 프로덕션으로 옮기면 일반적으로 수백만 또는 수십억 건의 비정형 데이터와 높은 쿼리 볼륨을 처리해야 합니다. 따라서 데이터 수집 및 정보 검색과 같은 일반적인 벡터 검색 파이프라인 작업을 효율적으로 실행하기 위한 강력한 솔루션이 필요합니다.
최근 발표에서, Zilliz의 Head of Ecosystem & AI Platform인 Jiang Chen은 효율적이고 프로덕션에 바로 사용할 수 있는 벡터 검색 파이프라인을 구축하는 단계별 프로세스를 소개했습니다. 이 글에서는 세 가지 주제로 구성된 발표의 주요 내용을 다룹니다:
기존 방식과 Retrieval Augmented Generation (RAG) 환경에서의 정보 검색 워크플로.
Milvus와 Spark를 사용하여 RAG 시스템에서 확장 가능한 벡터 검색 파이프라인 구축.
RAG 시스템의 품질을 개선하기 위한 조언.
더 이상 지체하지 않고, 첫 번째 주요 주제에 대해 이야기해 보겠습니다. 기존 방식에서의 정보 검색 워크플로부터 시작하겠습니다.
기존 정보 검색의 워크플로
딥러닝의 발전 이전에는 기존 정보 검색 또는 검색 시스템이 태그와 수동 라벨링에 크게 의존했습니다. 온라인 쇼핑몰의 사례를 생각해 보겠습니다. 온라인 쇼핑몰은 고객의 요구에 따라 가장 적합한 제품을 제공하기 위해 제품 태그에 의존했습니다. 따라서 온라인 쇼핑몰은 매일 많은 고객 쿼리를 처리할 수 있도록 해주는 확장 가능하고 효율적인 검색 파이프라인이 필요합니다.
이러한 요구를 수용하기 위해 기존 검색 시스템의 아키텍처는 일반적으로 두 가지 구성 요소로 나뉩니다: 하나는 오프라인 데이터 수집을 위한 것이고, 다른 하나는 온라인 쿼리 서빙을 위한 것입니다.
오프라인 데이터 수집
오프라인 데이터 수집의 주요 목표는 모든 데이터를 데이터베이스에 로드하는 것입니다. 이 프로세스의 첫 단계로, 데이터는 내부 문서나 인터넷과 같은 하나 또는 여러 소스에서 수집됩니다. 데이터를 가져올 수 있게 되면 데이터 태깅을 계속할 수 있습니다. 마지막으로, 태그가 지정된 데이터는 인덱싱되어 데이터베이스에 로드될 수 있습니다.
워크플로를 설명하기 위해 온라인 쇼핑몰 예시를 사용해 보겠습니다. 웹을 크롤링하여 인터넷에서 제품 설명을 얻을 수 있습니다. 다음으로, 설명을 확보한 후에는 'clothes', 'dress', 'formal dress', 'party dress' 등과 같이 해당 제품 설명을 나타내는 태그를 생성합니다. 그런 다음, 비정형 데이터베이스에 로드하기 전에 태그, 제품 설명, 가격 및 기타 메타데이터의 인덱스를 구축합니다. 마지막으로, 이 데이터베이스를 서빙 환경으로 푸시합니다.
기존 정보 검색 아키텍처의 두 가지 구성 요소
온라인 쿼리 서빙
두 번째 구성 요소의 주요 목표는 고객의 쿼리를 처리하고 이전 워크플로에서 생성한 데이터베이스로부터 정보 검색을 수행하는 것입니다.
프로세스는 프론트엔드와 쿼리 컴파일러에서 시작되어 사용자의 쿼리를 태그 집합으로 합성합니다. 다음으로, 시스템은 생성된 태그를 유사도 검색의 입력으로 사용합니다. 데이터베이스에서 가장 유사한 키워드 또는 태그를 가진 상위 k개 항목을 가져온 뒤, 사용 사례에 따라 달라지는 알고리즘을 사용하여 순위를 매깁니다. 순위가 매겨진 결과는 최종적으로 사용자에게 반환됩니다.
기존 정보 검색 및 검색 시스템의 주요 단점은 의미론적 이해가 부족하다는 점입니다. 태그나 수동으로 생성된 라벨은 사용자의 쿼리에 담긴 의미와 의도를 포착하지 못할 수 있으며, 이는 부정확한 검색 결과로 이어질 수 있습니다. 또한 데이터가 방대하다면 각 항목에 수동으로 태그를 지정하는 일은 번거로울 것입니다.
새로운 정보 검색으로서의 RAG 워크플로
딥러닝의 급속한 발전은 정보 검색 프로세스의 지형을 크게 변화시켰습니다. 임베딩 모델과 GPT, Claude, LLAMA, Mistral과 같은 대규모 언어 모델(LLM)의 도움으로 사용자의 쿼리에 담긴 의미를 효과적으로 포착할 수 있어, 각 데이터 항목에 대해 수동으로 레이블이나 태그를 생성할 필요가 없어졌습니다.
임베딩 모델을 사용하면 비정형 데이터를 n차원 벡터로 구성된 벡터 임베딩으로 변환할 수 있습니다. 임베딩의 차원 수는 사용된 모델에 따라 달라집니다. 이러한 임베딩은 자신이 나타내는 데이터의 의미를 담고 있으므로, 두 임베딩 간의 유사도는 코사인 거리와 같은 지표를 사용해 쉽게 계산할 수 있습니다. 직관적으로는, 유사한 의미를 담은 임베딩들이 벡터 공간에서 서로 더 가깝게 배치된다는 것입니다.
2D 벡터 공간에서 유사한 의미를 담은 벡터 임베딩의 예
임베딩을 확보하면, 이를 Milvus와 같은 벡터 데이터베이스에 직접 수집하여 데이터 수집 단계를 완료할 수 있습니다. 그 후 정보 검색 프로세스를 수행할 수 있습니다.
사용자 쿼리가 수신되면, 데이터 수집 단계에서 사용한 것과 동일한 임베딩 모델을 사용해 임베딩으로 변환됩니다. 다음으로, 파이프라인은 벡터 검색 작업을 수행하고 데이터베이스에서 가장 유사한 상위 k개의 임베딩을 가져옵니다. 검색 증강 생성(RAG) 맥락에서는, 이러한 유사 임베딩이 사용자의 쿼리에 답변하기 위한 LLM의 컨텍스트로 사용됩니다.
RAG 워크플로
Spark와 Milvus를 사용한 벡터 검색 파이프라인
RAG는 벡터 검색에서 가져온 관련 컨텍스트를 LLM에 제공함으로써 LLM이 생성하는 답변의 정확도를 향상시키는 새로운 접근 방식입니다. 그러나 프로덕션에 적합한 RAG 애플리케이션을 구축하는 것은 확장성 문제로 인해 어렵습니다.
프로덕션 환경에 RAG 애플리케이션을 배포할 때는 수백만 또는 수십억 개의 비정형 데이터를 처리하게 될 가능성이 큽니다. 또한 RAG 시스템은 고객으로부터 수천 개 또는 그 이상의 쿼리를 받게 됩니다. 따라서 이러한 문제를 효과적으로 처리하려면 효율적이고 확장 가능한 솔루션이 필요하며, 바로 이 지점에서 Apache Spark가 유용할 수 있습니다.
이 섹션에서는 Milvus와 Spark를 사용해 검색 파이프라인을 구축하겠습니다. Milvus는 방대한 양의 데이터에 대해 몇 초 만에 벡터 검색을 수행할 수 있게 해주는 오픈 소스 벡터 데이터베이스입니다. 한편 Spark는 대규모 데이터셋을 빠르고 효율적인 방식으로 처리하고 분석하는 데 특히 유용한 강력한 오픈 소스 분산 컴퓨팅 프레임워크입니다.
먼저 Milvus를 설치하는 것부터 시작하겠습니다. Milvus를 설치하는 방법은 여러 가지가 있지만, 프로덕션 환경에서 Milvus를 사용하려면 다음 명령어로 Docker에서 Milvus를 설치하고 실행하는 것이 가장 좋습니다.
# 설치 스크립트 다운로드
$ curl -sfL <https://raw.githubusercontent.com/milvus-io/milvus/master/scripts/standalone_embed.sh> -o standalone_embed.sh
# Docker 컨테이너 시작
$ bash standalone_embed.sh start
오픈 소스 벡터 데이터베이스인 Milvus는 여러 도구 및 AI 프레임워크와의 원활한 통합을 제공하여, RAG와 같은 프로덕션에 적합한 AI 기반 애플리케이션을 쉽게 구축할 수 있게 해줍니다. Apache Spark는 Milvus와 함께 사용하여 데이터 수집 및 쿼리 검색 프로세스를 효율적으로 확장할 수 있는 프레임워크 중 하나입니다.
Milvus와 Spark를 사용한 벡터 검색 파이프라인 워크플로의 예
Spark는 분산 처리 시스템이므로 데이터 처리 작업을 여러 컴퓨터에 일괄적으로 분산할 수 있습니다. 이 기능은 프로덕션에서 RAG 애플리케이션을 배포할 때처럼 대량의 데이터를 다룰 때 데이터 처리를 가속화합니다. 이 통합 덕분에 Milvus와 MySQL 같은 다른 데이터베이스 서비스 간에도 데이터를 이동할 수 있습니다.
Apache Spark를 설치하려면 해당 최신 설치 문서를 참조하세요. Spark를 설치한 후에는 spark-milvus jar 파일도 설치해야 합니다.
wget <https://github.com/zilliztech/spark-milvus/raw/1.0.0-SNAPSHOT/output/spark-milvus-1.0.0-SNAPSHOT.jar>
spark-milvus jar 파일을 다운로드한 후에는 다음 단계에 따라 종속성으로 추가할 수 있습니다.
# For pyspark
./bin/pyspark --jars spark-milvus-1.0.0-SNAPSHOT.jar
# For spark-shell
./bin/spark-shell --jars spark-milvus-1.0.0-SNAPSHOT.jar
이제 Milvus를 Spark와 통합할 준비가 되었습니다. 다음 예제에서는 Spark dataframe에서 Milvus로 직접 데이터를 수집하는 방법을 보여드리겠습니다.
import org.apache.spark.sql.{SaveMode, SparkSession}
import io.milvus.client.{MilvusClient, MilvusServiceClient}
import io.milvus.grpc.{DataType, FlushResponse, ImportResponse}
import io.milvus.param.bulkinsert.{BulkInsertParam, GetBulkInsertStateParam}
import io.milvus.param.collection.{CreateCollectionParam, DescribeCollectionParam, FieldType, FlushParam, LoadCollectionParam}
import io.milvus.param.dml.SearchParam
import io.milvus.param.index.CreateIndexParam
import io.milvus.param.{ConnectParam, IndexType, MetricType, R, RpcStatus}
import zilliztech.spark.milvus.{MilvusOptions, MilvusUtils}
import zilliztech.spark.milvus.MilvusOptions._
import org.apache.spark.SparkConf
import org.apache.spark.sql.types._
import org.apache.spark.sql.{SaveMode, SparkSession}
import org.apache.spark.sql.util.CaseInsensitiveStringMap
import org.apache.hadoop.fs.{FileSystem, Path}
import org.apache.log4j.Logger
import org.slf4j.LoggerFactory
import java.util
import scala.collection.JavaConverters._
object Hello extends App {
val spark = SparkSession.builder().master("local[*]")
.appName("HelloSparkMilvus")
.getOrCreate()
import spark.implicits._
// Create DataFrame
val sampleDF = Seq(
(1, "a", Seq(1.0,2.0,3.0,4.0,5.0)),
(2, "b", Seq(1.0,2.0,3.0,4.0,5.0)),
(3, "c", Seq(1.0,2.0,3.0,4.0,5.0)),
(4, "d", Seq(1.0,2.0,3.0,4.0,5.0))
).toDF("id", "text", "vec")
// set milvus options
val milvusOptions = Map(
"milvus.host" -> "localhost" -> uri,
"milvus.port" -> "19530",
"milvus.collection.name" -> "hello_spark_milvus",
"milvus.collection.vectorField" -> "vec",
"milvus.collection.vectorDim" -> "5",
"milvus.collection.primaryKeyField", "id"
)
sampleDF.write.format("milvus")
.options(milvusOptions)
.mode(SaveMode.Append)
.save()
}
위에 제공된 코드 스니펫에서는 ID, 텍스트, 벡터 임베딩이라는 세 개의 필드를 가진 Spark DataFrame을 "hello_spark_milvus"라는 컬렉션에 수집했습니다. 임베딩은 5차원 벡터로 구성되어 있으며, ID를 컬렉션의 기본 키로 사용했습니다.
또한 milvusOptions map 안에 Milvus 데이터베이스에 대한 몇 가지 구성 정보를 제공해야 합니다.
milvus.host및milvus.port: Milvus 서버와 포트입니다. Docker에서 Milvus를 실행하는 경우 기본 포트는 19530입니다.milvus.collection.name: 데이터가 수집될 Milvus 데이터베이스 내부의 컬렉션 이름입니다.milvus.collection.vectorField: 벡터 임베딩을 포함하는 데이터의 열 이름입니다.milvus.collection.vectorDim: 벡터 임베딩의 차원 수입니다.milvus.collection.primaryKeyField: 기본 키를 포함하는 데이터의 열 이름입니다.
조정할 수 있는 다양한 Milvus 옵션에 대해 더 자세히 알고 싶다면 Milvus documentation page를 확인하세요.
이제 Milvus 데이터베이스에 데이터를 수집했으므로, 컬렉션의 인덱싱 방법을 지정해야 합니다. Milvus는 일반 Flat 인덱스, Inverted Flat Index (IVF), Hierarchical Navigable Small World (HNSW) 등 다양한 인덱싱 방법을 지원합니다.
다음 예제에서는 HNSW의 맞춤형 버전인 AUTOINDEX를 사용하겠습니다. 벡터 검색 중 메트릭으로는 L2 거리를 사용하겠습니다.
val username = <YOUR_MILVUS_USER>
val password = <YOUR_MILVUS_PASSWORD>
val connectParam: ConnectParam = ConnectParam.newBuilder
.withHost("localhost")
.withPort("19530")
.withAuthorization(username, password)
.build
val client: MilvusClient = new MilvusServiceClient(connectParam)
val createIndexParam = CreateIndexParam.newBuilder()
.withCollectionName("hello_spark_milvus")
.withIndexName("index_name")
.withFieldName("vec")
.withMetricType(MetricType.L2)
.withIndexType(IndexType.AUTOINDEX)
.build()
val createIndexR = client.createIndex(createIndexParam)
println(createIndexR)
다음으로, 벡터 검색을 수행할 수 있으려면 먼저 “hello_spark_milvus” 컬렉션을 로드해야 합니다.
import io.milvus.param.collection.{CreateCollectionParam, DescribeCollectionParam, FieldType, FlushParam, LoadCollectionParam}
// 컬렉션 로드, 로드된 컬렉션만 검색할 수 있음
val loadCollectionParam = LoadCollectionParam.newBuilder().withCollectionName("hello_spark_milvus").build()
val loadCollectionR = client.loadCollection(loadCollectionParam)
println(loadCollectionR)
이제 데이터를 Milvus에 수집하고 인덱스를 생성했으므로, 드디어 벡터 검색 작업을 수행할 준비가 되었습니다. 앞서 생성한 Spark DataFrame의 첫 번째 행을 입력 벡터로 사용하겠습니다.
// 검색, 입력 dataframe의 첫 번째 행을 검색 벡터로 사용
val fieldList: util.List[String] = new util.ArrayList[String]()
fieldList.add("vec")
val searchVectors = util.Arrays.asList(sampleDF.first().getList(2))
val searchParam = SearchParam.newBuilder()
.withCollectionName("hello_spark_milvus")
.withMetricType(MetricType.L2)
.withOutFields("text")
.withVectors(searchVectors)
.withVectorFieldName("vec")
.withTopK(2)
.build()
val searchParamR = client.search(searchParam)
println(searchParamR)
보시다시피, 벡터 검색 작업을 수행하려면 SearchParam.newBuilder() 메서드 안에 다음과 같은 여러 메서드 호출을 제공해야 합니다.
.withCollectionName(): 벡터 검색이 수행되는 컬렉션 이름입니다..withMetricType(): 벡터 검색을 수행하는 데 사용되는 메트릭입니다..withOutFields(): 결과를 반환할 컬렉션의 출력 필드입니다..withVectors(): 입력 또는 쿼리 벡터입니다..withVectorFieldName(): 벡터 임베딩을 포함하는 컬렉션의 필드입니다..withTopK(): 쿼리 벡터와 가장 유사한 임베딩을 가진 상위 k개 항목을 반환합니다.
RAG 애플리케이션에서는 가장 유사한 상위 k개 항목이 컨텍스트로 사용되어 쿼리와 함께 LLM에 전달됩니다. 이렇게 하면 LLM이 해당 컨텍스트를 사용하여 쿼리에 대한 정확한 답변을 생성할 수 있습니다.
Spark와의 Milvus 통합을 활용할 수 있는 훨씬 더 많은 고급 사용 사례가 있습니다. 예를 들어, MySQL과 같은 일반 데이터베이스에서 데이터를 읽고, 이를 벡터 임베딩으로 변환한 다음, 해당 임베딩을 Milvus에 수집할 수 있습니다. 이러한 사용 사례는 이 GitHub 리포지토리 또는 이 Milvus 노트북 데모에서 살펴볼 수 있습니다.
좋은 RAG는 좋은 데이터에서 비롯됩니다
Milvus가 여러 AI 툴킷 및 프레임워크와 통합되어 있어 프로덕션 준비가 된 RAG(Retrieval Augmented Generation) 애플리케이션 개발이 간소화됩니다. 그러나 RAG 시스템을 프로덕션에 배포한 후에는 시스템이 생성하는 응답의 품질을 지속적으로 모니터링하는 것이 중요합니다.
응답 품질을 개선해야 한다면, 더 복잡한 알고리즘으로 들어가기 전에 먼저 기본 사항을 검토하는 것이 중요합니다. 핵심 초점은 RAG 시스템에서 사용하는 데이터 소스의 품질에 맞춰져야 합니다.
데이터 소스를 평가할 때는 다음 질문을 고려하세요.
사용자의 쿼리에 답변하는 데 필요한 데이터가 데이터베이스에 있나요?
필요한 모든 데이터를 데이터베이스에 수집했나요?
데이터를 데이터베이스에 수집하기 전에 올바른 데이터 전처리 단계를 수행했나요(예: 데이터 파싱, 데이터 정제, 청킹, 적절한 임베딩 모델 사용)?
데이터 소스의 품질을 확인한 후에는 알고리즘 관점에서 RAG 시스템의 품질을 개선하는 것을 고려할 수 있습니다. RAG 시스템의 성능을 향상하는 방법은 여러 가지가 있습니다. 예를 들면 다음과 같습니다.
더 강력한 임베딩 모델 사용: 다양한 사전 학습 또는 사용자 지정 학습 임베딩 모델을 실험하여 데이터의 의미적 관계를 가장 잘 포착하는 모델을 찾으세요.
쿼리 라우팅 및 타사 도구 통합 구현: 임베딩 모델이 문제가 아니라면, 쿼리 라우팅을 위한 에이전트를 적용하고 추가 도구 또는 데이터 소스와 통합하여 RAG 시스템을 개선할 수 있습니다.
기본 사항에 집중하고 데이터 소스와 알고리즘 구성 요소를 지속적으로 반복 개선함으로써, 프로덕션 준비가 된 RAG 애플리케이션이 사용자에게 고품질 응답을 제공하도록 보장할 수 있습니다.
결론
Milvus와 Spark 같은 다양한 프레임워크의 원활한 통합 덕분에 확장 가능한 LLM 기반 애플리케이션을 쉽게 구축하고 배포할 수 있습니다. Spark가 데이터 처리 작업을 여러 컴퓨터에 배치로 분산하는 기능은 데이터 처리 작업을 실제로 가속화합니다. 이 기능은 방대한 양의 데이터를 벡터 데이터베이스에 수집하려는 경우나 애플리케이션이 동시에 많은 수의 사용자 쿼리를 처리하는 경우에 특히 유용합니다.
데이터를 Milvus 벡터 데이터베이스에 수집하고 사용자의 쿼리를 받으면, 벡터 검색을 수행할 수 있습니다. 이 과정은 데이터베이스 내부 데이터 중 가장 관련성 높은 컨텍스트를 가져와 LLM에 전달하고 고도로 맥락화된 답변을 생성하는 데 매우 중요합니다.
계속 읽기

Announcing the General Availability of Zilliz Cloud BYOC on Google Cloud Platform
Zilliz Cloud BYOC on GCP offers enterprise vector search with full data sovereignty and seamless integration.

Proactive Monitoring for Vector Database: Zilliz Cloud Integrates with Datadog
we're excited to announce Zilliz Cloud's integration with Datadog, enabling comprehensive monitoring and observability for your vectorDB deployments.

DeepRAG: Thinking to Retrieval Step by Step for Large Language Models
Discover DeepRAG, an advanced retrieval-augmented generation (RAG) model that improves LLM accuracy by retrieving only essential data through step-by-step reasoning.


