Module 13 — Piloter Elasticsearch et Neo4j depuis Python
Karim, le développeur de Veille, veut brancher l'API du produit sur les deux moteurs sans réécrire les requêtes en HTTP brut. Sami, lui, veut automatiser un rapport mensuel : « donne-moi le top des auteurs sur les six derniers mois et leurs catégories favorites ». Ce module montre comment tout cela s'écrit en quinze lignes de Python, dans un conteneur du kit où les deux clients officiels sont déjà installés.
Le conteneur veille-python : Python sans Python
Le kit fournit un conteneur veille-python prêt à l'emploi. Son Dockerfile installe deux paquets et rien d'autre : elasticsearch>=9,<10 et neo4j>=5.28,<6. Le dossier python/ du kit est monté dans /work : tout script python/mon_script.py que vous écrivez est immédiatement exécutable dans le conteneur. Les variables d'environnement ES_URL, ES_USER, ES_PASSWORD, NEO4J_URI (bolt://neo4j:7687), NEO4J_USER, NEO4J_PASSWORD sont injectées par docker-compose.yml : le script les lit sans rien coder en dur.
Deux commandes suffisent :
./lab.sh python veille.py "climate change" # exécute python/veille.py avec un argument
./lab.sh python-shell # ouvre un interpréteur Python interactif
L'équivalent Windows PowerShell est .\lab.ps1 python veille.py "climate change". Aucun pip install sur la machine : c'est le principe du kit, et c'est ce qui rend le cours reproductible d'un poste à l'autre.
Vos fichiers .py vont dans le dossier python/ du kit. Le conteneur voit ce dossier sous /work. Enregistrez le fichier depuis votre éditeur, relancez ./lab.sh python … : la nouvelle version est prise en compte immédiatement, sans reconstruction d'image.
Le client elasticsearch 9 en cinq gestes
Le client Python d'Elasticsearch calque l'API REST. Vous écrivez du Python idiomatique, la bibliothèque fabrique les requêtes HTTP.
Se connecter
from elasticsearch import Elasticsearch
import os
es = Elasticsearch(
os.environ["ES_URL"],
basic_auth=(os.environ["ES_USER"], os.environ["ES_PASSWORD"]),
request_timeout=30,
)
info = es.info()
print(f"Elasticsearch {info['version']['number']} — cluster « {info['cluster_name']} »")
Sortie attendue :
Elasticsearch 9.5.3 — cluster « veille »
Elasticsearch(url, basic_auth=(u, p)) est l'appel canonique en version 9. Le client réutilise une pool de connexions HTTP ; ne créez pas un client par requête.
Rechercher
Toute requête Query DSL passe par es.search(index=..., query=..., size=..., source=[...]). Notez : à partir du client 9, les arguments s'écrivent directement (query=...), plus besoin d'envelopper dans body={"query": ...}.
rep = es.search(
index="news",
size=3,
query={
"multi_match": {
"query": "climate change",
"fields": ["headline^3", "short_description"],
}
},
source=["headline", "date", "category"],
)
for h in rep["hits"]["hits"]:
print(round(h["_score"], 2), h["_source"]["headline"])
Sortie attendue (les scores peuvent varier légèrement selon BM25) :
19.12 Change Is Here. Climate Change.
15.83 Ellen DeGeneres Warns Climate Change Will Be 'Dangerous'
14.44 Climate Change: Time for Action Is Now
Le nombre total de résultats se lit dans rep["hits"]["total"]["value"] — 2 834 pour climate change.
Lire, indexer, mettre à jour
doc = es.get(index="news", id="1") # un document par _id
es.index(index="news", id="200854", document={ # créer ou remplacer
"headline": "Veille lance sa v2",
"short_description": "Un moteur combinant Elasticsearch et Neo4j.",
"category": "TECH",
"authors": "Inès Bouraoui",
"date": "2026-09-09",
})
es.update(index="news", id="200854", doc={"category": "BUSINESS"})
Chaque méthode renvoie un dictionnaire avec _id, _version et result (created, updated, noop).
Indexer en masse : helpers.bulk
Écrire cinquante mille appels es.index séparés coûte trop cher. Le paquet elasticsearch.helpers fournit bulk, qui prépare le format NDJSON et gère les nouveaux essais.
from elasticsearch import helpers
actions = [
{"_index": "news", "_id": str(200_855 + i), "_source": {
"headline": f"Article de test numéro {i}",
"short_description": "Généré par le module 13.",
"category": "TECH",
"authors": "Sami Karray",
"date": "2026-09-09",
}}
for i in range(200)
]
ok, erreurs = helpers.bulk(es, actions, chunk_size=100, request_timeout=60)
print(f"{ok} documents indexés, {len(erreurs)} erreurs")
Sortie attendue :
200 documents indexés, 0 erreurs
chunk_size=100 envoie les documents par paquets de cent. Sur le corpus News, l'importateur du kit (importer/import_news.py) suit exactement cette logique, avec des lots de deux mille.
Gérer les erreurs
Deux exceptions reviennent plus souvent que les autres. Ne les laissez pas remonter sans message clair.
from elasticsearch import Elasticsearch, AuthenticationException, NotFoundError, TransportError
try:
es = Elasticsearch(os.environ["ES_URL"], basic_auth=(os.environ["ES_USER"], os.environ["ES_PASSWORD"]))
es.info()
except AuthenticationException:
print("401 : identifiants invalides. Vérifiez ELASTIC_PASSWORD dans .env,")
print("et si vous l'avez changé après le premier démarrage, faites ./lab.sh reset.")
raise
except TransportError as e:
print(f"Elasticsearch injoignable ({e}). Le service est-il « healthy » ? ./lab.sh status")
raise
AuthenticationException correspond au HTTP 401 « security_exception ... unable to authenticate user [elastic] ». NotFoundError s'utilise pour es.get sur un _id absent. TransportError couvre les problèmes réseau.
Le client neo4j 5 en cinq gestes
Le client officiel de Neo4j suit un modèle un peu différent : un pilote (driver), une session, et des transactions — implicites via session.run, ou explicites via execute_read et execute_write.
Se connecter
from neo4j import GraphDatabase
import os
pilote = GraphDatabase.driver(
os.environ["NEO4J_URI"],
auth=(os.environ["NEO4J_USER"], os.environ["NEO4J_PASSWORD"]),
)
pilote.verify_connectivity()
print("Neo4j accessible via", os.environ["NEO4J_URI"])
verify_connectivity() teste la connexion sans lancer de requête ; utile au démarrage d'une application pour tomber tout de suite si le pilote n'arrive pas à joindre le serveur.
Une requête, des paramètres
with pilote.session() as session:
resultat = session.run(
"MATCH (au:Auteur {nom: $nom})<-[:ECRIT_PAR]-(a:Article) "
"RETURN a.titre AS titre, a.date AS date "
"ORDER BY a.date DESC LIMIT 3",
nom="Lee Moran",
)
for enregistrement in resultat:
print(enregistrement["date"], enregistrement["titre"])
Sortie attendue (votre chiffre peut différer selon les articles récents) :
2018-05-25 Trump Cracks Terrible Joke About Assassinated Journalist
2018-05-24 Ivanka Trump's Latest Bit Of Advice Doesn't Sit Well With Some Twitter Users
2018-05-23 Rudy Giuliani Sparks Confusion With Baffling Answers
Toujours utiliser des paramètres ($nom), jamais du f"…{nom}…" : c'est plus rapide (Cypher met la requête en cache) et cela évite toute injection.
Transactions explicites
Pour un code robuste, la bonne pratique est d'envelopper une requête dans une fonction et de la passer à execute_read ou execute_write. Le pilote gère les nouveaux essais en cas de coupure réseau.
def top_auteurs(tx, limite: int):
rep = tx.run(
"MATCH (au:Auteur)<-[:ECRIT_PAR]-(a:Article) "
"RETURN au.nom AS auteur, count(a) AS articles "
"ORDER BY articles DESC LIMIT $limite",
limite=limite,
)
return [dict(r) for r in rep]
with pilote.session() as session:
for ligne in session.execute_read(top_auteurs, 5):
print(f"{ligne['articles']:>5} {ligne['auteur']}")
Sortie attendue :
4954 Reuters
2433 Lee Moran
1915 Ron Dicker
1328 Ed Mazza
1145 Cole Delbyck
Ces cinq lignes sont les auteurs les plus prolifiques du corpus News, tels que chargés par ./lab.sh cypher 11-charger-news.cypher.
Fermer proprement
pilote.close()
Ou, préférable, utiliser with GraphDatabase.driver(...) as pilote: — la connexion se ferme même en cas d'exception.
Gérer les erreurs
from neo4j.exceptions import ServiceUnavailable, AuthError
try:
pilote = GraphDatabase.driver(os.environ["NEO4J_URI"], auth=(os.environ["NEO4J_USER"], os.environ["NEO4J_PASSWORD"]))
pilote.verify_connectivity()
except AuthError:
print("Neo4j : mot de passe refusé. Si vous avez modifié NEO4J_PASSWORD après le premier démarrage,")
print("le mot de passe vit dans le volume : ./lab.sh reset puis ./lab.sh up.")
raise
except ServiceUnavailable as e:
print(f"Neo4j injoignable ({e}). Vérifiez ./lab.sh status.")
raise
AuthError couvre le HTTP 401 renvoyé par le serveur Bolt. ServiceUnavailable recouvre les timeouts et le service pas encore prêt.
Lecture guidée de python/veille.py
Le kit inclut python/veille.py, un script qui combine les deux moteurs : Elasticsearch trouve les articles pertinents pour une requête, Neo4j part de leurs auteurs et remonte d'autres articles à recommander. C'est le noyau du module 16.
rechercher() : la partie Elasticsearch
def rechercher(es: Elasticsearch, texte: str, taille: int = 5) -> list[dict]:
rep = es.search(
index="news",
size=taille,
query={
"multi_match": {
"query": texte,
"fields": ["headline^3", "short_description"],
"type": "best_fields",
}
},
source=["headline", "authors", "category", "date"],
)
return [{"id": int(h["_id"]), "score": round(h["_score"], 2), **h["_source"]} for h in rep["hits"]["hits"]]
Trois points à noter. headline^3 donne un poids trois fois plus élevé au titre qu'à la description : un mot dans le titre pèse davantage dans le score. type: "best_fields" garde, pour chaque document, le meilleur des deux champs plutôt que d'additionner. source=[...] filtre les champs renvoyés pour économiser de la bande passante — le script n'a pas besoin du link ni du texte intégral.
Le retour aplati fusionne _score, _id (le numéro de ligne, converti en int pour l'utiliser côté Cypher) et le _source en un seul dictionnaire par article.
recommander() : la partie Neo4j
def recommander(session, ids: list[int], limite: int = 5) -> list[dict]:
cypher = """
MATCH (a:Article)-[:ECRIT_PAR]->(au:Auteur)<-[:ECRIT_PAR]-(reco:Article)-[:PUBLIE_DANS]->(c:Categorie)
WHERE a.id IN $ids AND NOT reco.id IN $ids AND au.nom <> 'Reuters'
RETURN reco.titre AS titre, au.nom AS auteur, c.nom AS categorie, reco.date AS date
ORDER BY date DESC
LIMIT $limite
"""
return [dict(r) for r in session.run(cypher, ids=ids, limite=limite)]
Le motif se lit comme une phrase : « je pars des articles trouvés, je remonte leurs auteurs, je redescends vers d'autres articles écrits par ces mêmes auteurs, et je récupère la catégorie ». Les deux garde-fous — NOT reco.id IN $ids et au.nom <> 'Reuters' — évitent respectivement de recommander l'article que l'utilisateur vient de lire et d'inonder le résultat avec l'agence Reuters, cinq mille articles à elle seule.
main() : la colle entre les deux
def main() -> None:
es = Elasticsearch(ES_URL, basic_auth=(ES_USER, ES_PASSWORD))
info = es.info()
print(f"Elasticsearch {info['version']['number']} — cluster « {info['cluster_name']} »")
resultats = rechercher(es, REQUETE)
print(f"\n{len(resultats)} articles pour « {REQUETE} » :")
for r in resultats:
print(f" [{r['score']:>5}] {r['date']} {r['headline']} — {r['authors'] or 'sans auteur'} ({r['category']})")
pilote = GraphDatabase.driver(NEO4J_URI, auth=(NEO4J_USER, NEO4J_PASSWORD))
with pilote.session() as session:
recos = recommander(session, [r["id"] for r in resultats])
pilote.close()
print(f"\n{len(recos)} recommandations (mêmes auteurs, autres articles) :")
for r in recos:
print(f" {r['date']} {r['titre']} — {r['auteur']} ({r['categorie']})")
Rien de plus : deux appels aux clients, une boucle d'affichage. Toute la logique métier tient en deux fonctions et un point d'entrée.
Exécution réelle
./lab.sh python veille.py climate change
Sortie observée sur le corpus complet (les scores peuvent varier légèrement) :
Elasticsearch 9.5.3 — cluster « veille »
5 articles pour « climate change » :
[19.12] 2015-06-24 Change Is Here. Climate Change. — Dominique Browning (ENVIRONMENT)
[15.83] 2016-04-25 Ellen DeGeneres Warns Climate Change Will Be 'Dangerous' — Kate Sheppard (ENVIRONMENT)
[15.44] 2013-12-04 Climate Change Isn't Real, Say Fewer Than Ever — Kate Sheppard (POLITICS)
[14.44] 2014-11-19 Climate Change: Time for Action Is Now — Karen Feridun (POLITICS)
[13.98] 2016-08-15 Louisiana Flooding Is Climate Change — Kate Sheppard (ENVIRONMENT)
5 recommandations (mêmes auteurs, autres articles) :
2018-04-19 The E.P.A. Rolls Back A Regulation On Coal Ash — Kate Sheppard (ENVIRONMENT)
2017-11-08 A Vote For Common Sense And Public Health — Karen Feridun (POLITICS)
2017-06-01 Trump Just Pulled Out Of The Paris Agreement — Kate Sheppard (POLITICS)
2016-11-14 What Do We Tell The Children? — Dominique Browning (ENVIRONMENT)
2016-09-08 The State Of The Air We Breathe — Karen Feridun (ENVIRONMENT)
Le premier article a un score presque quatre points au-dessus du deuxième : les deux mots « climate » et « change » sont dans son titre, headline^3 pèse. Les recommandations remontent des articles postérieurs à ceux trouvés, écrits par les mêmes auteurs — c'est la valeur du couplage Elasticsearch + Neo4j.
Écrire son propre script : python/top_auteurs.py
Sami veut, pour un rapport mensuel, connaître les auteurs qui ont le plus écrit sur les six derniers mois de 2017, et pour chacun ses trois catégories principales. Elasticsearch fait l'agrégation, Neo4j complète avec le graphe.
Créez le fichier python/top_auteurs.py :
#!/usr/bin/env python3
"""Top auteurs sur une période + leurs catégories favorites."""
from __future__ import annotations
import os
from elasticsearch import Elasticsearch
from neo4j import GraphDatabase
ES = Elasticsearch(os.environ["ES_URL"], basic_auth=(os.environ["ES_USER"], os.environ["ES_PASSWORD"]))
PILOTE = GraphDatabase.driver(os.environ["NEO4J_URI"], auth=(os.environ["NEO4J_USER"], os.environ["NEO4J_PASSWORD"]))
def top_auteurs_es(du: str, au: str, taille: int = 5) -> list[tuple[str, int]]:
"""Agrégation terms sur authors.raw filtrée par date."""
rep = ES.search(
index="news",
size=0,
query={"range": {"date": {"gte": du, "lte": au}}},
aggs={
"auteurs": {
"terms": {"field": "authors.raw", "size": taille + 1, "exclude": ["", "Reuters"]}
}
},
)
seaux = rep["aggregations"]["auteurs"]["buckets"]
return [(s["key"], s["doc_count"]) for s in seaux[:taille]]
def categories_neo4j(nom: str, limite: int = 3) -> list[tuple[str, int]]:
"""Pour un auteur donné, ses catégories les plus fréquentes."""
cypher = """
MATCH (au:Auteur {nom: $nom})<-[:ECRIT_PAR]-(a:Article)-[:PUBLIE_DANS]->(c:Categorie)
RETURN c.nom AS categorie, count(a) AS articles
ORDER BY articles DESC LIMIT $limite
"""
with PILOTE.session() as session:
return [(r["categorie"], r["articles"]) for r in session.run(cypher, nom=nom, limite=limite)]
def main() -> None:
du, au = "2017-07-01", "2017-12-31"
print(f"Top 5 auteurs entre {du} et {au} (hors Reuters)\n")
for auteur, articles in top_auteurs_es(du, au):
cats = categories_neo4j(auteur)
etiquettes = ", ".join(f"{c} ({n})" for c, n in cats)
print(f" {articles:>4} {auteur:<25} {etiquettes}")
PILOTE.close()
if __name__ == "__main__":
main()
Lancez-le :
./lab.sh python top_auteurs.py
Sortie attendue (les compteurs des seconds semestres 2017 peuvent différer légèrement selon la version du corpus) :
Top 5 auteurs entre 2017-07-01 et 2017-12-31 (hors Reuters)
312 Lee Moran COMEDY (779), WEIRD NEWS (391), ENTERTAINMENT (387)
246 Ed Mazza POLITICS (546), ENTERTAINMENT (240), COMEDY (155)
198 Ron Dicker SPORTS (612), COMEDY (355), ENTERTAINMENT (198)
187 Cole Delbyck ENTERTAINMENT (886), STYLE & BEAUTY (76), POLITICS (54)
142 Nina Golgowski POLITICS (298), U.S. NEWS (188), CRIME (161)
Les chiffres à droite entre parenthèses viennent du graphe : ils comptent tous les articles de l'auteur, pas seulement ceux du second semestre 2017. Selon votre besoin métier, vous pouvez filtrer côté Cypher en ajoutant une clause WHERE a.date >= date($du).
Elasticsearch excelle à filtrer et agréger sur des millions de documents ; Neo4j excelle à parcourir des relations. Cette division est stable : quand une question mélange les deux, la faire résoudre en deux étapes est presque toujours plus lisible et plus rapide que d'entasser toute la logique dans un seul moteur.
À vous
Exercice 1 — Convertir une requête Kibana en Python
Prenez la requête suivante, écrite en Query DSL dans Kibana Dev Tools :
GET news/_search
{
"size": 3,
"query": {
"bool": {
"must": [ { "match": { "headline": "election" } } ],
"filter": [ { "range": { "date": { "gte": "2016-11-01", "lte": "2016-11-30" } } } ]
}
},
"_source": ["headline", "date"]
}
Réécrivez-la dans un script python/election_novembre.py, exécutez avec ./lab.sh python election_novembre.py et affichez les trois titres.
Solution
import os
from elasticsearch import Elasticsearch
es = Elasticsearch(os.environ["ES_URL"], basic_auth=(os.environ["ES_USER"], os.environ["ES_PASSWORD"]))
rep = es.search(
index="news",
size=3,
query={
"bool": {
"must": [{"match": {"headline": "election"}}],
"filter": [{"range": {"date": {"gte": "2016-11-01", "lte": "2016-11-30"}}}],
}
},
source=["headline", "date"],
)
for h in rep["hits"]["hits"]:
print(h["_source"]["date"], h["_source"]["headline"])
L'ordre du bool en Python est identique à l'ordre en JSON : le client se contente de sérialiser. Le seul changement : _source devient l'argument source (le préfixe _ est réservé aux métadonnées côté serveur).
Exercice 2 — Recommandation avec execute_read
Écrivez python/reco_auteur.py qui prend un nom d'auteur en argument et affiche les trois articles les plus récents des auteurs qui partagent au moins deux catégories avec lui. Utilisez session.execute_read.
Solution
import os, sys
from neo4j import GraphDatabase
nom = " ".join(sys.argv[1:]) or "Lee Moran"
pilote = GraphDatabase.driver(os.environ["NEO4J_URI"], auth=(os.environ["NEO4J_USER"], os.environ["NEO4J_PASSWORD"]))
def voisins(tx, nom_auteur, limite=3):
cypher = """
MATCH (moi:Auteur {nom: $nom})<-[:ECRIT_PAR]-(:Article)-[:PUBLIE_DANS]->(c:Categorie)
<-[:PUBLIE_DANS]-(:Article)-[:ECRIT_PAR]->(voisin:Auteur)
WHERE voisin <> moi
WITH voisin, count(DISTINCT c) AS communes
WHERE communes >= 2
MATCH (voisin)<-[:ECRIT_PAR]-(a:Article)
RETURN voisin.nom AS auteur, a.titre AS titre, a.date AS date
ORDER BY a.date DESC LIMIT $limite
"""
return [dict(r) for r in tx.run(cypher, nom=nom_auteur, limite=limite)]
with pilote.session() as session:
for r in session.execute_read(voisins, nom):
print(f"{r['date']} [{r['auteur']}] {r['titre']}")
pilote.close()
Le double MATCH traverse deux fois la relation :PUBLIE_DANS pour trouver les auteurs qui partagent des catégories avec la personne demandée. execute_read gère les nouveaux essais si Neo4j répond lentement.
Exercice 3 — Indexer un mini-corpus avec helpers.bulk
Écrivez python/ajouter_notes.py qui insère dans l'index news cinq articles fictifs correspondant à des notes internes de Veille (catégorie INTERNE, auteur Inès Bouraoui), puis vérifiez avec GET news/_count que l'index a bien grandi de cinq.
Solution
import os
from elasticsearch import Elasticsearch, helpers
es = Elasticsearch(os.environ["ES_URL"], basic_auth=(os.environ["ES_USER"], os.environ["ES_PASSWORD"]))
sujets = [
"Point hebdomadaire produit — semaine 36",
"Ateliers Kibana : nouveaux tableaux pour Léa",
"Revue de sécurité mensuelle",
"Feuille de route Q4 : reco temps réel",
"Onboarding Sami : bilan à 30 jours",
]
actions = [
{"_index": "news", "_id": f"veille-{i}", "_source": {
"headline": s,
"short_description": f"Note interne du {i} septembre 2026.",
"category": "INTERNE",
"authors": "Inès Bouraoui",
"date": "2026-09-09",
}}
for i, s in enumerate(sujets, start=1)
]
ok, erreurs = helpers.bulk(es, actions)
print(f"{ok} documents indexés")
es.indices.refresh(index="news")
print("Total news :", es.count(index="news")["count"])
Sortie attendue (le total dépend de votre état ; sur un corpus fraîchement importé, il passe de 200 853 à 200 858) :
5 documents indexés
Total news : 200858
Le refresh explicite garantit que le count immédiat voit bien les documents ; en production, l'intervalle de refresh d'une seconde suffit.
Points à retenir
- Le conteneur
veille-pythona les clientselasticsearch9 etneo4j5 préinstallés ; vos scripts vivent dans le dossierpython/du kit et se lancent avec./lab.sh python <script>.py. - Toutes les variables (
ES_URL,ES_USER,ES_PASSWORD,NEO4J_URI,NEO4J_USER,NEO4J_PASSWORD) sont déjà injectées ; ne codez jamais un mot de passe en dur. - Côté Elasticsearch, les gestes clés sont
es.info(),es.search(index=, query=, size=, source=),es.get,es.index,es.update, ethelpers.bulkpour un import massif. - Côté Neo4j, un pilote (
GraphDatabase.driver), une session (with pilote.session()), une requête (session.runousession.execute_read/execute_write), et toujours des paramètres nommés$x. - Trois exceptions à traiter proprement :
AuthenticationException(401 Elasticsearch),AuthError(Neo4j),ServiceUnavailable(Neo4j pas prêt). - La combinaison des deux moteurs, illustrée par
veille.py, tient en deux fonctions : Elasticsearch trouve, Neo4j élargit.
Si ça ne marche pas
elasticsearch.AuthenticationException: 401 security_exception→ mot de passe changé après création du volume../lab.sh resetpuis./lab.sh up.neo4j.exceptions.ServiceUnavailable→ Neo4j n'est pas prêt ou le pilote pointe verslocalhostau lieu deneo4j. Vérifiez./lab.sh statuset confirmez queNEO4J_URIvautbolt://neo4j:7687dans le conteneur.ImportError: No module named elasticsearch→ vous exécutez le script sur votre machine hôte au lieu du conteneur. Utilisez toujours./lab.sh python <script>.py.ModuleNotFoundError: No module named 'elasticsearch.helpers'alors que l'import passe → conflit avec un ancien paquet Python installé dans un venv de votre machine. Le conteneur n'a pas le problème : passez par./lab.sh python.