Создание готовых к промышленной эксплуатации поисковых конвейеров с Spark и Milvus
Создание масштабируемого конвейера векторного поиска в production не так просто, как создание его прототипа. При работе над прототипом мы часто имеем дело лишь с небольшим объемом неструктурированных данных. Однако при переносе прототипа в production нам обычно нужно обрабатывать миллионы или миллиарды неструктурированных данных и большие объемы запросов. Поэтому требуется надежное решение для эффективного выполнения распространенных операций конвейера векторного поиска, таких как загрузка данных и извлечение информации.
В недавнем выступлении Jiang Chen, Head of Ecosystem & AI Platform в Zilliz, представил пошаговый процесс создания эффективного и готового к production конвейера векторного поиска. В этой статье будут рассмотрены основные моменты выступления, которые состоят из трех тем:
Рабочий процесс информационного поиска в традиционных настройках и в Retrieval Augmented Generation (RAG).
Создание масштабируемого конвейера векторного поиска в системе RAG с Milvus и Spark.
Советы по улучшению качества вашей системы RAG.
Без лишних слов перейдем к первой основной теме и начнем с рабочего процесса информационного поиска в традиционной настройке.
Рабочий процесс традиционного информационного поиска
До достижений в области глубокого обучения традиционные системы извлечения информации или поиска в значительной степени полагались на теги и ручную разметку. Рассмотрим пример онлайн-магазинов: они полагались на теги продуктов, чтобы предлагать клиентам наиболее подходящие товары в соответствии с их потребностями. Поэтому онлайн-магазинам нужен масштабируемый и эффективный поисковый конвейер, который позволяет им ежедневно обслуживать большое количество клиентских запросов.
Чтобы удовлетворить этот спрос, архитектура традиционной поисковой системы обычно делится на два компонента: один для офлайн-загрузки данных и один для онлайн-обслуживания запросов.
Офлайн-загрузка данных
Основная цель офлайн-загрузки данных — загрузить все данные в базу данных. В качестве первого шага этого процесса данные собираются из одного или нескольких источников, таких как внутренние документы или Интернет. Как только мы можем получить данные, мы можем продолжить разметку данных тегами. Наконец, размеченные тегами данные могут быть проиндексированы и загружены в базу данных.
Давайте используем пример онлайн-магазина, чтобы проиллюстрировать рабочий процесс. Мы можем получать описания товаров из Интернета путем веб-скрейпинга. Затем, когда у нас есть описания, мы создаем теги, представляющие эти описания товаров, такие как 'clothes', 'dress', 'formal dress', 'party dress' и т. д. Затем мы строим индекс тегов, описаний товаров, цен и других метаданных перед загрузкой их в неструктурированную базу данных. Наконец, мы отправляем эту базу данных в среду обслуживания.
Два компонента архитектуры традиционного информационного поиска
Онлайн-обслуживание запросов
Основная цель второго компонента — обслуживать запросы клиентов и выполнять извлечение информации из базы данных, которую мы создали в предыдущем рабочем процессе.
Процесс начинается с front-end и компилятора запросов, чтобы преобразовать запрос пользователя в набор тегов. Затем система будет использовать сгенерированные теги в качестве входных данных для поиска по сходству. Верхние k записей в базе данных с наиболее похожими ключевыми словами или тегами затем будут извлечены и ранжированы с использованием алгоритма, который варьируется в зависимости от сценариев использования. Ранжированный результат в итоге будет возвращен пользователю.
Главный недостаток традиционных систем извлечения информации и поиска — отсутствие семантического понимания. Теги или вручную созданные метки не могли уловить семантическое значение и намерение запроса пользователя, что могло приводить к неточным результатам поиска. Кроме того, ручная разметка каждой записи была бы трудоемкой, если у нас есть огромный объем данных.
Рабочий процесс RAG как нового информационного поиска
Стремительное развитие глубокого обучения значительно изменило ландшафт процессов поиска информации. С помощью моделей эмбеддингов и больших языковых моделей (LLMs), таких как GPT, Claude, LLAMA и Mistral, семантическое значение запроса пользователя может быть эффективно зафиксировано, что устраняет необходимость вручную создавать метки или теги для каждой записи данных.
С помощью моделей эмбеддингов неструктурированные данные могут быть преобразованы в векторные эмбеддинги, которые состоят из n-мерных векторов. Размерность эмбеддинга зависит от используемой модели. Эти эмбеддинги несут семантическое значение данных, которые они представляют, и поэтому сходство между любыми двумя эмбеддингами можно легко вычислить с использованием таких метрик, как косинусное расстояние. Интуитивно это означает, что эмбеддинги, несущие схожие значения, будут расположены ближе друг к другу в векторном пространстве.
Пример векторных эмбеддингов, несущих схожее семантическое значение в 2D-векторном пространстве
После того как у нас есть эмбеддинги, они могут быть напрямую загружены в векторную базу данных, такую как Milvus, завершая часть, связанную с загрузкой данных. После этого может быть выполнен процесс поиска информации.
Когда поступает пользовательский запрос, он будет преобразован в эмбеддинг с использованием той же модели эмбеддингов, что и на этапе загрузки данных. Затем конвейер выполняет операцию векторного поиска и извлекает из базы данных k наиболее похожих эмбеддингов. В контексте Retrieval Augmented Generation (RAG) эти похожие эмбеддинги затем используются как контексты для LLM при ответе на запрос пользователя.
Рабочий процесс RAG
Конвейер векторного поиска со Spark и Milvus
RAG — это новый подход к повышению точности ответов, генерируемых LLM, за счет предоставления ей релевантных контекстов, полученных из векторного поиска. Однако создание готового к эксплуатации RAG-приложения является сложной задачей из-за проблем масштабируемости.
При развертывании RAG-приложения в production-среде вы, скорее всего, будете иметь дело с миллионами или даже миллиардами неструктурированных данных. Кроме того, ваша RAG-система будет получать тысячи или даже больше запросов от клиентов. Поэтому для эффективного решения этих проблем необходимо производительное и масштабируемое решение, и именно здесь может быть полезен Apache Spark.
В этом разделе мы построим поисковый конвейер с использованием Milvus и Spark. Milvus — это open-source векторная база данных, которая позволяет нам выполнять векторный поиск по огромным объемам данных за секунды. В то же время Spark — это мощный open-source фреймворк распределенных вычислений, особенно полезный для быстрой и эффективной обработки и анализа больших наборов данных.
Начнем с установки Milvus. Существует несколько способов установить Milvus, но если вы хотите использовать Milvus в production-среде, лучше всего установить и запустить Milvus в Docker с помощью следующих команд:
# Download the installation script
$ curl -sfL <https://raw.githubusercontent.com/milvus-io/milvus/master/scripts/standalone_embed.sh> -o standalone_embed.sh
# Start the Docker container
$ bash standalone_embed.sh start
Как open-source векторная база данных, Milvus предлагает бесшовные интеграции со многими инструментами и AI-фреймворками, что упрощает создание готовых к эксплуатации AI-приложений, таких как RAG. Apache Spark — один из фреймворков, который можно использовать вместе с Milvus для эффективного масштабирования процессов загрузки данных и извлечения запросов.
Пример рабочего процесса конвейера векторного поиска с Milvus и Spark
Поскольку Spark является распределённой системой обработки, он способен распределять задачи обработки данных между несколькими компьютерами в пакетном режиме. Эта возможность ускоряет обработку данных при работе с огромными объёмами данных, например при развёртывании RAG-приложения в production. Благодаря этой интеграции мы также можем перемещать данные между Milvus и другими сервисами баз данных, такими как MySQL.
Чтобы установить Apache Spark, обратитесь к их последней документации по установке. После установки Spark вам также потребуется установить jar-файл spark-milvus.
wget <https://github.com/zilliztech/spark-milvus/raw/1.0.0-SNAPSHOT/output/spark-milvus-1.0.0-SNAPSHOT.jar>
После загрузки jar-файла spark-milvus вы можете добавить его как зависимость, выполнив следующие шаги:
# Для pyspark
./bin/pyspark --jars spark-milvus-1.0.0-SNAPSHOT.jar
# Для 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()
}
В приведённом выше фрагменте кода мы загрузили Spark DataFrame с тремя полями: ID, текстом и векторным embedding в коллекцию под названием "hello_spark_milvus". Embeddings состоят из 5-мерных векторов, и мы использовали ID в качестве первичного ключа нашей коллекции.
Нам также необходимо предоставить несколько параметров конфигурации нашей базы данных Milvus внутри map milvusOptions:
milvus.hostиmilvus.port: сервер и порт Milvus. Если вы запускаете Milvus в Docker, порт по умолчанию — 19530.milvus.collection.name: имя коллекции внутри базы данных Milvus, куда будут загружены данные.milvus.collection.vectorField: имя столбца наших данных, который содержит векторный embedding.milvus.collection.vectorDim: размерность нашего векторного embedding.milvus.collection.primaryKeyField: имя столбца наших данных, который содержит первичный ключ.
Если вы хотите узнать больше о различных видах параметров Milvus, которые можно настроить, ознакомьтесь со страницей документации Milvus.
Теперь, когда мы загрузили данные в базу данных Milvus, нам нужно указать метод индексирования для нашей коллекции. Milvus поддерживает различные методы индексирования, такие как обычный индекс Flat, Inverted Flat Index (IVF), Hierarchical Navigable Small World (HNSW) и многие другие.
В следующем примере мы будем использовать AUTOINDEX, который является настроенной версией HNSW. В качестве метрики при векторном поиске мы будем использовать расстояние 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}
// Load collection, only loaded collection can be searched
val loadCollectionParam = LoadCollectionParam.newBuilder().withCollectionName("hello_spark_milvus").build()
val loadCollectionR = client.loadCollection(loadCollectionParam)
println(loadCollectionR)
Теперь, когда мы загрузили данные в Milvus и создали индекс, наконец мы готовы выполнить операцию векторного поиска. Мы будем использовать первую строку Spark DataFrame, созданного ранее, в качестве входного вектора.
// Search, use the first row of input dataframe as search vector
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 может использовать контексты для генерации точного ответа на запрос.
Существует гораздо больше продвинутых сценариев использования, в которых вы можете использовать интеграцию Milvus со Spark. Например, вы можете читать данные из своей обычной базы данных, такой как MySQL, преобразовывать их в векторные эмбеддинги и загружать эти эмбеддинги в Milvus. Вы можете изучить эти сценарии использования в этом репозитории GitHub или в этом демонстрационном ноутбуке Milvus.
Хороший RAG начинается с хороших данных
Интеграция Milvus со многими AI-инструментами и фреймворками упрощает разработку готовых к промышленной эксплуатации приложений RAG (Retrieval Augmented Generation). Однако после развертывания RAG-системы в production важно постоянно отслеживать качество ответов, генерируемых системой.
Если качество ответов необходимо улучшить, крайне важно сначала изучить основы, прежде чем переходить к более сложным алгоритмам. Основное внимание следует уделить качеству источника данных, используемого RAG-системой.
При оценке источника данных рассмотрите следующие вопросы:
Есть ли в нашей базе данных необходимые данные для ответа на запросы пользователя?
Собрали ли мы все необходимые данные в нашу базу данных?
Выполнили ли мы правильные этапы предварительной обработки данных перед загрузкой данных в нашу базу данных (например, парсинг данных, очистка данных, разбиение на фрагменты, использование подходящих embedding-моделей)?
После того как вы проверили качество источника данных, вы можете рассмотреть улучшение качества RAG-системы с алгоритмической точки зрения. Существует несколько способов повысить производительность RAG-системы, например:
Использование более мощных embedding-моделей: Экспериментируйте с различными предварительно обученными или обученными на заказ embedding-моделями, чтобы найти ту, которая лучше всего отражает семантические связи в ваших данных.
Реализация маршрутизации запросов и интеграции со сторонними инструментами: Если проблема не в embedding-модели, вы можете улучшить RAG-систему, применив агента для маршрутизации запросов и интегрировав дополнительные инструменты или источники данных.
Сосредоточившись на основах и постоянно итерируя источник данных и алгоритмические компоненты, вы можете гарантировать, что ваше готовое к промышленной эксплуатации RAG-приложение предоставляет пользователям высококачественные ответы.
Заключение
Бесшовная интеграция Milvus с различными фреймворками, такими как Spark, позволяет нам легко создавать и развертывать масштабируемое приложение на базе LLM. Способность Spark распределять задачи обработки данных между несколькими компьютерами пакетами действительно ускоряет операции обработки данных. Эта функция особенно полезна, когда мы хотим загрузить огромные объемы данных в нашу векторную базу данных или когда наше приложение одновременно обрабатывает большое количество пользовательских запросов.
После того как мы загрузим данные в нашу векторную базу данных Milvus и получим запрос пользователя, мы можем выполнить векторный поиск. Этот процесс крайне важен для извлечения наиболее релевантных контекстов среди данных внутри нашей базы данных, чтобы передать их в LLM для генерации ответов с высокой степенью контекстуализации.
Читать далее

Bringing AI to Legal Tech: The Role of Vector Databases in Enhancing LLM Guardrails
Discover how vector databases enhance AI reliability in legal tech, ensuring accurate, compliant, and trustworthy AI-powered legal solutions.

Legal Document Analysis: Harnessing Zilliz Cloud's Semantic Search and RAG for Legal Insights
Enhance legal document analysis with Zilliz Cloud’s Semantic Search and RAG. Improve accuracy, efficiency, and scalability for contracts, case law, and compliance.

DeepSeek-VL2: Mixture-of-Experts Vision-Language Models for Advanced Multimodal Understanding
Explore DeepSeek-VL2, the open-source MoE vision-language model. Discover its architecture, efficient training pipeline, and top-tier performance.


