Creación de pipelines de búsqueda listos para producción con Spark y Milvus
Crear un pipeline de búsqueda vectorial escalable en producción no es tan fácil como construir su prototipo. Al trabajar en un prototipo, a menudo tratamos solo con una pequeña cantidad de datos no estructurados. Sin embargo, al llevar el prototipo a producción, normalmente necesitamos manejar millones o miles de millones de datos no estructurados y altos volúmenes de consultas. Por lo tanto, se necesita una solución robusta para la ejecución eficiente de operaciones comunes del pipeline de búsqueda vectorial, como la ingesta de datos y la recuperación de información.
En una charla reciente, Jiang Chen, Head of Ecosystem & AI Platform en Zilliz, presentó un proceso paso a paso para construir un pipeline de búsqueda vectorial eficiente y listo para producción. Este artículo analizará los puntos principales de la charla, que constan de tres temas:
Flujo de trabajo de la búsqueda de información en entornos tradicionales y de Retrieval Augmented Generation (RAG).
Construcción de un pipeline de búsqueda vectorial escalable en un sistema RAG con Milvus y Spark.
Consejos para mejorar la calidad de tu sistema RAG.
Sin más preámbulos, hablemos del primer tema principal y empecemos con el flujo de trabajo de la búsqueda de información en un entorno tradicional.
Flujo de trabajo de una búsqueda de información tradicional
Antes de los avances del aprendizaje profundo, los sistemas tradicionales de recuperación o búsqueda de información dependían en gran medida de etiquetas y etiquetado manual. Considera el caso de las tiendas en línea: dependían de etiquetas de productos para ofrecer a los clientes los productos más adecuados según sus necesidades. Por lo tanto, las tiendas en línea necesitan un pipeline de búsqueda escalable y eficiente que les permita atender una gran cantidad de consultas de clientes todos los días.
Para satisfacer esta demanda, la arquitectura del sistema de búsqueda tradicional normalmente se divide en dos componentes: uno para la ingesta de datos sin conexión y otro para atender consultas en línea.
Ingesta de datos sin conexión
El objetivo principal de la ingesta de datos sin conexión es cargar todos los datos en una base de datos. Como primer paso de este proceso, los datos se recopilan de una o múltiples fuentes, como documentos internos o Internet. Una vez que podemos obtener los datos, podemos continuar con el etiquetado de datos. Finalmente, los datos etiquetados pueden indexarse y cargarse en la base de datos.
Usemos el ejemplo de la tienda en línea para ilustrar el flujo de trabajo. Podemos obtener descripciones de productos de Internet rastreando la web. A continuación, una vez que tenemos las descripciones, creamos etiquetas que representan esas descripciones de productos, como 'ropa', 'vestido', 'vestido formal', 'vestido de fiesta', etc. Luego, construimos el índice de las etiquetas, descripciones de productos, precios y otros metadatos antes de cargarlos en una base de datos no estructurada. Finalmente, enviamos esta base de datos al entorno de servicio.
Dos componentes de una arquitectura tradicional de búsqueda de información
Servicio de consultas en línea
El objetivo principal del segundo componente es atender las consultas de los clientes y realizar la recuperación de información desde la base de datos que hemos creado en el flujo de trabajo anterior.
El proceso comienza desde el front-end y el compilador de consultas para sintetizar la consulta del usuario en un conjunto de etiquetas. A continuación, el sistema utilizará las etiquetas generadas como entradas para la búsqueda por similitud. Las primeras k entradas de la base de datos con las palabras clave o etiquetas más similares se recuperarán y clasificarán utilizando un algoritmo que varía según los casos de uso. El resultado clasificado finalmente se devolverá al usuario.
El principal inconveniente de los sistemas tradicionales de recuperación y búsqueda de información es la falta de comprensión semántica. Las etiquetas o rótulos creados manualmente no podían capturar el significado semántico ni la intención de la consulta del usuario, lo que podría conducir a resultados de búsqueda inexactos. Además, etiquetar manualmente cada entrada sería engorroso si tenemos una enorme cantidad de datos.
El flujo de trabajo de RAG como la nueva búsqueda de información
Los rápidos avances del aprendizaje profundo han cambiado significativamente el panorama de los procesos de recuperación de información. Con la ayuda de modelos de embeddings y modelos de lenguaje grandes (LLMs) como GPT, Claude, LLAMA y Mistral, el significado semántico de la consulta de un usuario puede capturarse eficazmente, eliminando la necesidad de crear manualmente etiquetas o tags para cada entrada de datos.
Con los modelos de embeddings, los datos no estructurados pueden transformarse en embeddings vectoriales, que consisten en vectores n-dimensionales. La dimensionalidad del embedding depende del modelo utilizado. Estos embeddings contienen el significado semántico de los datos que representan y, por lo tanto, la similitud entre dos embeddings cualesquiera puede calcularse fácilmente utilizando métricas como la distancia coseno. La intuición es que los embeddings que contienen significados similares se colocarán más cerca unos de otros en el espacio vectorial.
Ejemplo de embeddings vectoriales que contienen un significado semántico similar en un espacio vectorial 2D
Una vez que tenemos los embeddings, pueden ingerirse directamente en una base de datos vectorial como Milvus, completando la parte de ingestión de datos. Después, puede llevarse a cabo el proceso de recuperación de información.
Cuando se recibe una consulta de usuario, se transformará en un embedding utilizando el mismo modelo de embeddings durante la parte de ingestión de datos. A continuación, el pipeline realiza una operación de búsqueda vectorial y obtiene los k embeddings más similares de la base de datos. En el contexto de Retrieval Augmented Generation (RAG), estos embeddings similares se utilizan luego como contextos para que el LLM responda a la consulta del usuario.
Flujo de trabajo de RAG
Pipeline de búsqueda vectorial con Spark y Milvus
RAG es un enfoque novedoso para mejorar la precisión de las respuestas generadas por un LLM proporcionándole contextos relevantes obtenidos de una búsqueda vectorial. Sin embargo, crear una aplicación RAG lista para producción es un desafío debido a problemas de escalabilidad.
Al desplegar una aplicación RAG en producción, lo más probable es que estés tratando con millones o incluso miles de millones de datos no estructurados. Además, tu sistema RAG recibirá miles o incluso más consultas de clientes. Por lo tanto, se necesita una solución eficiente y escalable para manejar estos problemas de manera efectiva, y aquí es donde Apache Spark puede ser útil.
En esta sección, construiremos un pipeline de búsqueda utilizando Milvus y Spark. Milvus es una base de datos vectorial de código abierto que nos permite realizar búsquedas vectoriales en cantidades masivas de datos en segundos. Mientras tanto, Spark es un potente framework de computación distribuida de código abierto, particularmente útil para procesar y analizar grandes conjuntos de datos de manera rápida y eficiente.
Comencemos instalando Milvus. Hay varias formas de instalar Milvus, pero si quieres usar Milvus en un entorno de producción, lo mejor es instalar y ejecutar Milvus en Docker con los siguientes comandos:
# Descargar el script de instalación
$ curl -sfL <https://raw.githubusercontent.com/milvus-io/milvus/master/scripts/standalone_embed.sh> -o standalone_embed.sh
# Iniciar el contenedor Docker
$ bash standalone_embed.sh start
Como base de datos vectorial de código abierto, Milvus ofrece integraciones fluidas con muchas herramientas y frameworks de IA, lo que facilita la creación de aplicaciones impulsadas por IA listas para producción como RAG. Apache Spark es uno de los frameworks que pueden utilizarse junto con Milvus para escalar eficientemente los procesos de ingestión de datos y recuperación de consultas.
Ejemplo de un flujo de trabajo de pipeline de búsqueda vectorial con Milvus y Spark
Dado que Spark es un sistema de procesamiento distribuido, puede distribuir tareas de procesamiento de datos entre varios equipos en un lote. Esta característica acelera el procesamiento de datos cuando se trabaja con cantidades masivas de datos, como al desplegar una aplicación RAG en producción. Gracias a esta integración, también podemos mover datos entre Milvus y otros servicios de bases de datos como MySQL.
Para instalar Apache Spark, consulta su documentación de instalación más reciente. Una vez que hayas instalado Spark, también tendrás que instalar el archivo jar spark-milvus.
wget <https://github.com/zilliztech/spark-milvus/raw/1.0.0-SNAPSHOT/output/spark-milvus-1.0.0-SNAPSHOT.jar>
Una vez que hayas descargado el archivo jar spark-milvus, puedes agregarlo como dependencia siguiendo estos pasos:
# 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
Y ahora estamos listos para integrar Milvus con Spark. En el siguiente ejemplo, te mostraremos cómo puedes ingerir datos desde un dataframe de Spark directamente en 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()
}
En el fragmento de código proporcionado arriba, ingerimos un DataFrame de Spark con tres campos: un ID, un texto y un embedding vectorial en una colección llamada "hello_spark_milvus". Los embeddings consisten en vectores de 5 dimensiones, y usamos el ID como clave primaria de nuestra colección.
También necesitamos proporcionar varias piezas de configuración sobre nuestra base de datos Milvus dentro del mapa milvusOptions:
milvus.hostymilvus.port: servidor y puerto de Milvus. Si ejecutas Milvus en Docker, el puerto predeterminado es 19530.milvus.collection.name: el nombre de la colección dentro de la base de datos Milvus donde se ingerirán los datos.milvus.collection.vectorField: el nombre de la columna de nuestros datos que contiene el embedding vectorial.milvus.collection.vectorDim: la dimensionalidad de nuestro embedding vectorial.milvus.collection.primaryKeyField: el nombre de la columna de nuestros datos que contiene la clave primaria.
Si deseas saber más sobre los diferentes tipos de opciones de Milvus que puedes ajustar, consulta la página de documentación de Milvus.
Ahora que hemos ingerido datos en la base de datos de Milvus, necesitamos especificar el método de indexación para nuestra colección. Milvus admite varios métodos de indexación, como el índice Flat normal, Inverted Flat Index (IVF), Hierarchical Navigable Small World (HNSW) y muchos más.
En el siguiente ejemplo, usaremos AUTOINDEX, que es una versión personalizada de HNSW. Como métrica durante la búsqueda vectorial, usaremos la distancia 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)
A continuación, necesitamos cargar nuestra colección “hello_spark_milvus” antes de poder realizar una búsqueda vectorial en ella.
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)
Ahora que hemos ingerido los datos en Milvus y creado un índice, finalmente estamos listos para realizar una operación de búsqueda vectorial. Usaremos la primera fila del Spark DataFrame que creamos anteriormente como vector de entrada.
// 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)
Como puedes ver, necesitamos proporcionar varias llamadas a métodos dentro del método SearchParam.newBuilder() para realizar una operación de búsqueda vectorial, como:
.withCollectionName(): el nombre de la colección donde se realiza la búsqueda vectorial..withMetricType(): la métrica utilizada para realizar la búsqueda vectorial..withOutFields(): los campos de salida en una colección para devolver el resultado..withVectors(): el vector de entrada o de consulta..withVectorFieldName(): el campo en una colección que contiene embeddings vectoriales..withTopK(): devuelve las primeras k entradas que tienen los embeddings más similares al vector de consulta.
En una aplicación RAG, las k entradas más similares se usarán como contextos para pasarlos junto con la consulta al LLM. De este modo, el LLM puede usar los contextos para generar una respuesta precisa a la consulta.
Hay muchos casos de uso más avanzados en los que puedes aprovechar la integración de Milvus con Spark. Por ejemplo, puedes leer datos de tu base de datos habitual como MySQL, transformarlos en incrustaciones vectoriales e ingerir esas incrustaciones en Milvus. Puedes explorar estos casos de uso en este repositorio de GitHub o en esta demostración de notebook de Milvus.
Un buen RAG proviene de buenos datos
La integración de Milvus con muchos kits de herramientas y frameworks de IA simplifica el desarrollo de aplicaciones RAG (Retrieval Augmented Generation) listas para producción. Sin embargo, después de desplegar un sistema RAG en producción, es importante supervisar continuamente la calidad de las respuestas generadas por el sistema.
Si es necesario mejorar la calidad de las respuestas, es crucial examinar primero los fundamentos antes de profundizar en algoritmos más complejos. El enfoque clave debe estar en la calidad de la fuente de datos utilizada por el sistema RAG.
Al evaluar la fuente de datos, considera las siguientes preguntas:
¿Tenemos los datos necesarios para responder a las consultas del usuario en nuestra base de datos?
¿Hemos recopilado todos los datos necesarios en nuestra base de datos?
¿Hemos realizado los pasos correctos de preprocesamiento de datos antes de ingerir los datos en nuestra base de datos (por ejemplo, análisis de datos, limpieza de datos, fragmentación, uso de modelos de incrustación adecuados)?
Una vez que hayas verificado la calidad de la fuente de datos, puedes considerar mejorar la calidad del sistema RAG desde una perspectiva algorítmica. Hay varias formas de mejorar el rendimiento de un sistema RAG, tales como:
Usar modelos de incrustación más potentes: Experimenta con diferentes modelos de incrustación preentrenados o entrenados a medida para encontrar el que mejor capture las relaciones semánticas en tus datos.
Implementar enrutamiento de consultas e integración con herramientas de terceros: Si el modelo de incrustación no es el problema, puedes mejorar el sistema RAG aplicando un agente para el enrutamiento de consultas e integrándolo con herramientas o fuentes de datos adicionales.
Al centrarte en los fundamentos e iterar continuamente sobre la fuente de datos y los componentes algorítmicos, puedes asegurarte de que tu aplicación RAG lista para producción entregue respuestas de alta calidad a tus usuarios.
Conclusión
La integración fluida de Milvus con varios frameworks como Spark nos facilita crear y desplegar una aplicación escalable impulsada por LLM. La capacidad de Spark para distribuir tareas de procesamiento de datos entre varios ordenadores en lotes acelera mucho las operaciones de procesamiento de datos. Esta característica es especialmente útil cuando queremos ingerir cantidades masivas de datos en nuestra base de datos vectorial o cuando nuestra aplicación está gestionando una gran cantidad de consultas de usuarios al mismo tiempo.
Una vez que ingerimos datos en nuestra base de datos vectorial Milvus y recibimos la consulta de un usuario, podemos realizar una búsqueda vectorial. Este proceso es crucial para obtener los contextos más relevantes entre los datos dentro de nuestra base de datos, que se pasarán a un LLM para generar respuestas altamente contextualizadas.
Sigue leyendo

Introducing Functions and Model Inference on Zilliz Cloud: Automatic Embedding and Reranking with Hosted Models
Zilliz Cloud Functions auto-generate embeddings via OpenAI, Voyage AI, Cohere, or Zilliz Hosted Models. Built-in reranking — just insert text and search.

Zilliz Cloud Launches in AWS Australia, Expanding Global Reach to Australia and Neighboring Markets
We're thrilled to announce that Zilliz Cloud is now available in the AWS Sydney, Australia region (ap-southeast-2).

The Real Bottlenecks in Autonomous Driving — And How AI Infrastructure Can Solve Them
Autonomous driving faces a data bottleneck. Learn how AI-native vector databases like Zilliz solve scale, cost, and insight challenges across AV pipelines.


