Déployer un système de recherche perfromant qui s'adapte dynamiquement aux nouveaux documents constitue un défi technique majeur. Ce guide détaille la conception d'une architecture permettant à un service utilisant le modèle d'embedding BGE-M3 de traiter les données entrantes en continu, sans interruption de service ni retraitement complet de l'ensemble du corpus.
Problématique des mises à jour en temps réel
Les déploiements traditionnels des modèles d'embedding fonctionnent de manière statique. L'index et les vecteurs sont générés une fois au démarrage. L'ajout de nouveaux documents nécessite un recalcul complet de tous les vecteurs et la reconstruction de l'index, ce qui immobilise le service. Cette approche est incompatible avec les sources de données à flux élevé comme les fils d'actualité, les journaux en temps réel ou les bases de connaissances collaboratives.
L'objectif est de passer d'un modèle "tout statique" à un système "incrémental dynamique". Le principe fondamental consiste à :
- Traiter exclusivement les nouveaux documents pour la vectorisation.
- Mettre à jour l'index de manière incrémentielle.
- Maintenir le service de recherche disponible en permanence.
Architecture du système de mise à jour
L'architecture repose sur une chaîne de traitement asynchrone et découplée, garantissant la résilience et l'évolutivité.
Composants principaux
- Ingestion et file d'attente : Les documents issus de différentes sources (API, CDC, monitoring de fichiers) sont stockés dans une base documentaire (ex: MongoDB) et une notification est envoyée à une file de messages (ex: RabbitMQ, Kafka).
- Workers d'embedding incrémental : Ces services sans état consomment les messages, récupèrent le texte correspondant, et utilisent BGE-M3 pour générer le vecteur associé. Ils peuvent être mis à l'échelle horizontalement et optimisés par traitement par lots.
- Gestionnaire d'index vectoriel : Reçoit les nouveaux vecteurs et les insère dans la base vectorielle (ex: Milvus, Weaviate). Il déclenche ou planifie le rafraîchissement de l'index pour rendre les vecteurs recherchables.
- Service de recherche : Expose l'API de recherche. Il interroge l'index qui intègre automatiquement les dernières données validées.
Flux de données
Le flux est asynchrone et garantit une cohérence éventuelle :
- Nouveau document → Stockage persistant → Message dans la file d'attente.
- Worker traite le message → Génère le vecteur → L'insère dans la base vectorielle.
- La base vectorielle fusionne périodiquement les nouveaux vecteurs dans l'index principal.
- Les requêtes de recherche accèdent à la vue la plus récente de l'index.
Le délai entre l'ingestion d'un document et sa recherche est de l'ordre de quelques secondes à une minute, selon le débit et la politique de rafraîchissement.
Implémentation des composants clés
1. Producteur de tâches (Surveillance de fichiers)
Ce script surveille un répertoire, enregistre les nouveaux fichiers texte dans MongoDB et publie un message de tâche dans RabbitMQ.
import os, json, time, uuid
from watchdog.observers import Observer
from watchdog.events import FileSystemEventHandler
import pika
from pymongo import MongoClient
class TraitementFichier(FileSystemEventHandler):
def __init__(self, canal_mq, collection_mongo):
self.canal = canal_mq
self.docs = collection_mongo
def on_created(self, event):
if not event.is_directory and event.src_path.endswith('.txt'):
self.traiter(event.src_path)
def traiter(self, chemin_fichier):
try:
with open(chemin_fichier, 'r', encoding='utf-8') as f:
texte = f.read()
identifiant = str(uuid.uuid4())
document = {
"_id": identifiant,
"source": chemin_fichier,
"contenu": texte,
"horodatage": time.time()
}
self.docs.insert_one(document)
message = {
"id_doc": identifiant,
"operation": "vectoriser"
}
self.canal.basic_publish(
exchange='',
routing_key='file_attente_vectorisation',
body=json.dumps(message),
properties=pika.BasicProperties(delivery_mode=2)
)
except Exception as e:
print(f"Erreur sur {chemin_fichier}: {e}")
client_mongo = MongoClient('localhost', 27017)
db = client_mongo['bd_connaissances']
collection_docs = db['documents_bruts']
conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
canal = conn.channel()
canal.queue_declare(queue='file_attente_vectorisation', durable=True)
gestionnaire = TraitementFichier(canal, collection_docs)
observateur = Observer()
observateur.schedule(gestionnaire, path='./dossier_surveillance', recursive=False)
observateur.start()
try:
while True:
time.sleep(1)
except KeyboardInterrupt:
observateur.stop()
observateur.join()
conn.close()
2. Worker d'embedding incrémental
Le worker charge le modèle une seule fois, consomme les messages, récupère le texte depuis MongoDB et génère le vecteur dense via BGE-M3 avant de l'envoyer vers la base vectorielle.
import json, pika, torch
from pymongo import MongoClient
from FlagEmbedding import BGEM3FlagModel
from pymilvus import connections, Collection
class OuvrierVectorisation:
def __init__(self, chemin_modele, uri_mongo, hote_milvus, port_milvus):
self.modele = BGEM3FlagModel(chemin_modele, use_fp16=True)
self.modele.eval()
self.client_mongo = MongoClient(uri_mongo)
self.collection_docs = self.client_mongo['bd_connaissances']['documents_bruts']
connections.connect("default", host=hote_milvus, port=port_milvus)
self.collection_vecteurs = Collection("vecteurs_documents")
self.tampon = []
self.taille_lot = 20
def traiter_message(self, canal, methode, proprietes, corps):
donnees = json.loads(corps)
doc = self.collection_docs.find_one({"_id": donnees['id_doc']})
if doc:
self.tampon.append((donnees['id_doc'], doc['contenu']))
if len(self.tampon) >= self.taille_lot:
self._traiter_lot()
canal.basic_ack(delivery_tag=methode.delivery_tag)
def _traiter_lot(self):
if not self.tampon:
return
ids, textes = zip(*self.tampon)
with torch.no_grad():
resultats = self.modele.encode(
list(textes),
batch_size=4,
max_length=4096,
return_dense=True,
return_sparse=False,
return_colbert_vecs=False
)
vecteurs = resultats['dense_vecs']
donnees_insertion = [
{"id_vector": ids[i], "embedding": vecteurs[i].tolist()}
for i in range(len(ids))
]
self.collection_vecteurs.insert(donnees_insertion)
self.tampon.clear()
def demarrer(self):
conn = pika.BlockingConnection(pika.ConnectionParameters('localhost'))
canal = conn.channel()
canal.queue_declare(queue='file_attente_vectorisation', durable=True)
canal.basic_qos(prefetch_count=1)
canal.basic_consume(queue='file_attente_vectorisation',
on_message_callback=self.traiter_message)
try:
canal.start_consuming()
except KeyboardInterrupt:
self._traiter_lot()
canal.stop_consuming()
conn.close()
if __name__ == "__main__":
ouvrier = OuvrierVectorisation(
chemin_modele="/chemin/vers/bge-m3",
uri_mongo="mongodb://localhost:27017",
hote_milvus="localhost",
port_milvus="19530"
)
ouvrier.demarrer()
3. Stratégie de rafraîchissement de l'index
Les bases vectorielles modernes gèrent l'insertion incrémentale via des segments. Un processus périodique consolide ces segments pour optimiser les performances de recherche.
from pymilvus import connections, Collection
import schedule, time
class GestionnaireIndex:
def __init__(self, hote, port, nom_collection):
connections.connect("default", host=hote, port=port)
self.collection = Collection(nom_collection)
def rafraichir_index(self):
print("Rafraîchissement de l'index...")
self.collection.flush()
# La compaction et l'optimisation des segments se gèrent selon la politique de la base
print("Rafraîchissement terminé.")
def planifier_maintenance(self, intervalle_heures=4):
schedule.every(intervalle_heures).hours.do(self.rafraichir_index)
while True:
schedule.run_pending()
time.sleep(60)
# Exemple d'utilisation comme service d'arrière-plan
# gestionnaire = GestionnaireIndex("localhost", "19530", "vecteurs_documents")
# gestionnaire.planifier_maintenance(intervalle_heures=6)
Déploiement et optimisation
Déploiement via Docker Compose
L'orchestration des services garantit un environnement reproductible.
version: '3.8'
services:
mongodb:
image: mongo:6
volumes: ["mongo_data:/data/db"]
ports: ["27017:27017"]
milvus:
image: milvusdb/milvus:v2.3
volumes: ["milvus_data:/var/lib/milvus"]
ports: ["19530:19530"]
depends_on: [etcd]
etcd:
image: quay.io/coreos/etcd:v3.5
rabbitmq:
image: rabbitmq:3-management
ports: ["5672:5672", "15672:15672"]
producteur-docs:
build: ./producteur
volumes: ["./docs:/app/observation"]
depends_on: [mongodb, rabbitmq]
ouvrier-vectorisation:
build: ./ouvrier
deploy:
replicas: 2
environment:
- CUDA_VISIBLE_DEVICES=0
depends_on: [mongodb, milvus, rabbitmq]
api-recherche:
build: ./api
ports: ["8000:8000"]
depends_on: [milvus, mongodb]
volumes:
mongo_data:
milvus_data:
Recommandations de performance
- Taille des lots : Ajuster la taille du lot de vectorisation en fonction de la mémoire GPU et de la longueur moyenne des documents.
- Type d'index vectoriel : Choisir HNSW pour un bon compromis vitesse/recouvrement, ou IVF pour les très grands ensembles de données.
- Fréquence de rafraîchissement : Définir un intervalle court (1-5 minutes) pour les applications critiquse en temps réel, ou plus long (plusieurs heures) pour les mises à jour de connaissances.
- Mise en cache : Implémenter un cache au niveau de l'API de recherche pour les requêtes fréquentes.
Cette architecture transforme un service d'embedding statique en un système capable d'absorber et de rendre recherchables les nouveaux documents au fil de l'eau, répondant ainsi aux exigences des applications modernes à données dynamiques.