Créer des pipelines de recherche prêts pour la production avec Spark et Milvus
Construire un pipeline de recherche vectorielle évolutif en production n’est pas aussi facile que construire son prototype. Lorsque nous travaillons sur un prototype, nous traitons souvent seulement une petite quantité de données non structurées. Cependant, lors du passage du prototype en production, nous devons généralement gérer des millions ou des milliards de données non structurées et des volumes de requêtes élevés. Par conséquent, une solution robuste est nécessaire pour l’exécution efficace des opérations courantes d’un pipeline de recherche vectorielle, telles que l’ingestion de données et la récupération d’informations.
Dans une conférence récente, Jiang Chen, Head of Ecosystem & AI Platform chez Zilliz, a présenté un processus étape par étape pour construire un pipeline de recherche vectorielle efficace et prêt pour la production. Cet article abordera les points principaux de la conférence, qui se composent de trois sujets :
Flux de travail de la recherche d’informations dans des contextes traditionnels et de Retrieval Augmented Generation (RAG).
Construire un pipeline de recherche vectorielle évolutif dans un système RAG avec Milvus et Spark.
Conseils pour améliorer la qualité de votre système RAG.
Sans plus attendre, parlons du premier sujet principal, et commençons par le flux de travail de la recherche d’informations dans un contexte traditionnel.
Flux de travail d’une recherche d’informations traditionnelle
Avant les avancées du deep learning, les systèmes traditionnels de récupération d’informations ou de recherche reposaient fortement sur les tags et l’étiquetage manuel. Prenons le cas des boutiques en ligne : elles s’appuyaient sur les tags de produits pour proposer aux clients les produits les plus adaptés à leurs besoins. Par conséquent, les boutiques en ligne ont besoin d’un pipeline de recherche évolutif et efficace qui leur permet de traiter chaque jour un grand nombre de requêtes clients.
Pour répondre à cette demande, l’architecture du système de recherche traditionnel est normalement divisée en deux composants : l’un pour l’ingestion de données hors ligne et l’autre pour le service des requêtes en ligne.
Ingestion de données hors ligne
L’objectif principal de l’ingestion de données hors ligne est de charger toutes les données dans une base de données. Comme première étape de ce processus, les données sont collectées depuis une ou plusieurs sources, telles que des documents internes ou Internet. Une fois que nous pouvons récupérer les données, nous pouvons poursuivre avec le tagging des données. Enfin, les données taguées peuvent être indexées et chargées dans la base de données.
Utilisons l’exemple de la boutique en ligne pour illustrer le flux de travail. Nous pouvons obtenir des descriptions de produits depuis Internet en explorant le web. Ensuite, une fois que nous avons les descriptions, nous créons des tags qui représentent ces descriptions de produits, tels que 'vêtements', 'robe', 'robe formelle', 'robe de soirée', etc. Puis, nous construisons l’index des tags, des descriptions de produits, des prix et d’autres métadonnées avant de les charger dans une base de données non structurée. Enfin, nous transférons cette base de données vers l’environnement de service.
Deux composants d’une architecture de recherche d’informations traditionnelle
Service des requêtes en ligne
L’objectif principal du second composant est de servir les requêtes des clients et d’effectuer la récupération d’informations depuis la base de données que nous avons créée dans le flux de travail précédent.
Le processus commence par le front-end et le compilateur de requêtes afin de synthétiser la requête de l’utilisateur en un ensemble de tags. Ensuite, le système utilisera les tags générés comme entrées pour la recherche de similarité. Les k premières entrées de la base de données ayant les mots-clés ou tags les plus similaires seront alors récupérées et classées à l’aide d’un algorithme qui varie selon les cas d’utilisation. Le résultat classé sera enfin renvoyé à l’utilisateur.
Le principal inconvénient des systèmes traditionnels de récupération d’informations et de recherche est le manque de compréhension sémantique. Les tags ou les étiquettes créées manuellement ne pouvaient pas capturer le sens sémantique ni l’intention de la requête de l’utilisateur, ce qui pouvait conduire à des résultats de recherche inexacts. De plus, taguer manuellement chaque entrée serait fastidieux si nous avons une quantité massive de données.
Le flux de travail de RAG comme nouvelle recherche d’informations
Les progrès rapides du deep learning ont considérablement changé le paysage des processus de recherche d’informations. Avec l’aide de modèles d’embedding et de grands modèles de langage (LLM) comme GPT, Claude, LLAMA et Mistral, le sens sémantique de la requête d’un utilisateur peut être capturé efficacement, éliminant ainsi la nécessité de créer manuellement des libellés ou des tags pour chaque entrée de données.
Avec les modèles d’embedding, les données non structurées peuvent être transformées en embeddings vectoriels, qui consistent en des vecteurs à n dimensions. La dimensionnalité de l’embedding dépend du modèle utilisé. Ces embeddings portent le sens sémantique des données qu’ils représentent et, par conséquent, la similarité entre deux embeddings peut être facilement calculée à l’aide de métriques comme la distance cosinus. L’intuition est que les embeddings qui portent des significations similaires seront placés plus près les uns des autres dans l’espace vectoriel.
Exemple d’embeddings vectoriels qui portent un sens sémantique similaire dans un espace vectoriel 2D
Une fois que nous avons les embeddings, ils peuvent être directement ingérés dans une base de données vectorielle comme Milvus, complétant ainsi la partie ingestion des données. Ensuite, le processus de recherche d’informations peut être mené.
Lorsqu’une requête utilisateur est reçue, elle sera transformée en embedding à l’aide du même modèle d’embedding que lors de la partie ingestion des données. Ensuite, le pipeline effectue une opération de recherche vectorielle et récupère les k embeddings les plus similaires depuis la base de données. Dans le contexte de la génération augmentée par récupération (RAG), ces embeddings similaires sont ensuite utilisés comme contextes pour que le LLM réponde à la requête de l’utilisateur.
Flux de travail RAG
Pipeline de recherche vectorielle avec Spark et Milvus
RAG est une approche nouvelle pour améliorer la précision des réponses générées par un LLM en lui fournissant des contextes pertinents récupérés à partir d’une recherche vectorielle. Cependant, créer une application RAG prête pour la production est difficile en raison des problèmes de scalabilité.
Lors du déploiement d’une application RAG en production, vous aurez très probablement affaire à des millions, voire des milliards, de données non structurées. De plus, votre système RAG recevra des milliers, voire davantage, de requêtes de la part des clients. Par conséquent, une solution efficace et scalable est nécessaire pour gérer ces problèmes efficacement, et c’est là qu’Apache Spark peut être utile.
Dans cette section, nous allons créer un pipeline de recherche à l’aide de Milvus et de Spark. Milvus est une base de données vectorielle open source qui nous permet d’effectuer une recherche vectorielle sur d’énormes volumes de données en quelques secondes. Pendant ce temps, Spark est un puissant framework de calcul distribué open source particulièrement utile pour traiter et analyser de grands ensembles de données de manière rapide et efficace.
Commençons par installer Milvus. Il existe plusieurs façons d’installer Milvus, mais si vous souhaitez utiliser Milvus dans un environnement de production, il est préférable d’installer et d’exécuter Milvus dans Docker avec les commandes suivantes :
# 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
En tant que base de données vectorielle open source, Milvus offre des intégrations transparentes avec de nombreux outils et frameworks d’IA, ce qui facilite la création d’applications alimentées par l’IA prêtes pour la production, comme RAG. Apache Spark est l’un des frameworks qui peuvent être utilisés avec Milvus pour mettre à l’échelle efficacement les processus d’ingestion des données et de récupération des requêtes.
Exemple de flux de travail d’un pipeline de recherche vectorielle avec Milvus et Spark
Comme Spark est un système de traitement distribué, il est capable de répartir les tâches de traitement des données sur plusieurs ordinateurs dans un lot. Cette fonctionnalité accélère le traitement des données lorsqu’il s’agit de volumes massifs de données, par exemple lors du déploiement d’une application RAG en production. Grâce à cette intégration, nous pouvons également déplacer des données entre Milvus et d’autres services de bases de données comme MySQL.
Pour installer Apache Spark, veuillez consulter leur dernière documentation d’installation. Une fois Spark installé, vous devrez également installer le fichier jar spark-milvus.
wget <https://github.com/zilliztech/spark-milvus/raw/1.0.0-SNAPSHOT/output/spark-milvus-1.0.0-SNAPSHOT.jar>
Une fois que vous avez téléchargé le fichier jar spark-milvus, vous pouvez l’ajouter comme dépendance en suivant ces étapes :
# Pour pyspark
./bin/pyspark --jars spark-milvus-1.0.0-SNAPSHOT.jar
# Pour spark-shell
./bin/spark-shell --jars spark-milvus-1.0.0-SNAPSHOT.jar
Et maintenant, nous sommes prêts à intégrer Milvus avec Spark. Dans l’exemple suivant, nous vous montrerons comment ingérer des données d’un dataframe Spark directement dans 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()
}
Dans l’extrait de code fourni ci-dessus, nous avons ingéré un DataFrame Spark avec trois champs : un ID, un texte et un embedding vectoriel dans une collection appelée "hello_spark_milvus". Les embeddings se composent de vecteurs à 5 dimensions, et nous avons utilisé l’ID comme clé primaire de notre collection.
Nous devons également fournir plusieurs éléments de configuration concernant notre base de données Milvus dans la map milvusOptions :
milvus.hostetmilvus.port: serveur et port Milvus. Si vous exécutez Milvus dans Docker, le port par défaut est 19530.milvus.collection.name: le nom de la collection dans la base de données Milvus où les données seront ingérées.milvus.collection.vectorField: le nom de la colonne de nos données qui contient l’embedding vectoriel.milvus.collection.vectorDim: la dimensionnalité de notre embedding vectoriel.milvus.collection.primaryKeyField: le nom de la colonne de nos données qui contient la clé primaire.
Si vous souhaitez en savoir plus sur les différents types d’options Milvus que vous pouvez ajuster, consultez la page de documentation Milvus.
Maintenant que nous avons ingéré des données dans la base de données Milvus, nous devons spécifier la méthode d’indexation pour notre collection. Milvus prend en charge diverses méthodes d’indexation, telles que l’index Flat standard, l’Inverted Flat Index (IVF), Hierarchical Navigable Small World (HNSW), et bien d’autres.
Dans l’exemple suivant, nous utiliserons AUTOINDEX, qui est une version personnalisée de HNSW. Comme métrique lors de la recherche vectorielle, nous utiliserons la distance 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)
Ensuite, nous devons charger notre collection “hello_spark_milvus” avant de pouvoir effectuer une recherche vectorielle dessus.
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)
Maintenant que nous avons ingéré les données dans Milvus et créé un index, nous sommes enfin prêts à effectuer une opération de recherche vectorielle. Nous utiliserons la première ligne du Spark DataFrame que nous avons créé précédemment comme vecteur d’entrée.
// 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)
Comme vous pouvez le voir, nous devons fournir plusieurs appels de méthode dans la méthode SearchParam.newBuilder() pour effectuer une opération de recherche vectorielle, tels que :
.withCollectionName(): le nom de la collection où la recherche vectorielle est effectuée..withMetricType(): la métrique utilisée pour effectuer la recherche vectorielle..withOutFields(): les champs de sortie d’une collection à renvoyer dans le résultat..withVectors(): le vecteur d’entrée ou de requête..withVectorFieldName(): le champ d’une collection qui contient les embeddings vectoriels..withTopK(): renvoie les k premières entrées dont les embeddings sont les plus similaires au vecteur de requête.
Dans une application RAG, les k entrées les plus similaires seront utilisées comme contextes à transmettre avec la requête au LLM. De cette façon, le LLM peut utiliser les contextes pour générer une réponse précise à la requête.
Il existe de nombreux cas d’utilisation plus avancés dans lesquels vous pouvez tirer parti de l’intégration de Milvus avec Spark. Par exemple, vous pouvez lire des données depuis votre base de données habituelle comme MySQL, les transformer en embeddings vectoriels, et ingérer ces embeddings dans Milvus. Vous pouvez explorer ces cas d’utilisation dans ce dépôt GitHub ou dans cette démo notebook Milvus.
Un bon RAG provient de bonnes données
L’intégration de Milvus avec de nombreuses boîtes à outils et frameworks d’IA simplifie le développement d’applications RAG (Retrieval Augmented Generation) prêtes pour la production. Cependant, après avoir déployé un système RAG en production, il est important de surveiller en continu la qualité des réponses générées par le système.
Si la qualité des réponses doit être améliorée, il est essentiel d’examiner d’abord les fondamentaux avant de se lancer dans des algorithmes plus complexes. L’accent principal doit être mis sur la qualité de la source de données utilisée par le système RAG.
Lors de l’évaluation de la source de données, tenez compte des questions suivantes :
Disposons-nous des données requises pour répondre aux requêtes de l’utilisateur dans notre base de données ?
Avons-nous collecté toutes les données nécessaires dans notre base de données ?
Avons-nous effectué les bonnes étapes de prétraitement des données avant d’ingérer les données dans notre base de données (par exemple, analyse des données, nettoyage des données, découpage en chunks, utilisation de modèles d’embedding appropriés) ?
Une fois que vous avez vérifié la qualité de la source de données, vous pouvez ensuite envisager d’améliorer la qualité du système RAG d’un point de vue algorithmique. Il existe plusieurs façons d’améliorer les performances d’un système RAG, telles que :
Utiliser des modèles d’embedding plus puissants : Expérimentez avec différents modèles d’embedding pré-entraînés ou entraînés sur mesure pour trouver celui qui capture le mieux les relations sémantiques dans vos données.
Mettre en œuvre le routage des requêtes et l’intégration d’outils tiers : Si le modèle d’embedding n’est pas en cause, vous pouvez améliorer le système RAG en appliquant un agent pour le routage des requêtes et en l’intégrant à des outils ou sources de données supplémentaires.
En vous concentrant sur les fondamentaux et en itérant en continu sur la source de données et les composants algorithmiques, vous pouvez garantir que votre application RAG prête pour la production fournit des réponses de haute qualité à vos utilisateurs.
Conclusion
L’intégration transparente de Milvus avec divers frameworks comme Spark nous permet de créer et de déployer facilement une application évolutive propulsée par un LLM. La capacité de Spark à distribuer les tâches de traitement des données sur plusieurs ordinateurs par lots accélère réellement les opérations de traitement des données. Cette fonctionnalité est particulièrement utile lorsque nous voulons ingérer d’énormes volumes de données dans notre base de données vectorielle ou lorsque notre application traite simultanément un grand nombre de requêtes utilisateur.
Une fois que nous avons ingéré les données dans notre base de données vectorielle Milvus et reçu la requête d’un utilisateur, nous pouvons alors effectuer une recherche vectorielle. Ce processus est crucial pour récupérer les contextes les plus pertinents parmi les données présentes dans notre base de données afin de les transmettre à un LLM pour générer des réponses hautement contextualisées.
Continuer à lire

Zilliz Skills Breakdown: How AI Agents Master Vector Databases
Zilliz's Milvus Skill (pymilvus, 7 files) and Zilliz Cloud Skill (zilliz-cli, 14 modules) bring vector-DB dev and ops into one Claude Code session.

Why I’m Against Claude Code’s Grep-Only Retrieval? It Just Burns Too Many Tokens
Learn how vector-based code retrieval cuts Claude Code token consumption by 40%. Open-source solution with easy MCP integration. Try claude-context today.

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.


