قيادة Elasticsearch وNeo4j من بايثون
يريد Karim، مطوِّر Veille، أن يربط واجهة المنتج بالمحرِّكَين دون إعادة كتابة الاستعلامات بـHTTP خامًا. أمّا Sami فيريد أن يُؤتمت تقريرًا شهريًّا: «أعطني قائمة أعلى المؤلِّفين خلال الأشهر الستّة الأخيرة وفئاتهم المفضَّلة». تُبيِّن هذه الوحدة كيف يُكتَب كلّ ذلك في خمسة عشر سطرًا من بايثون، داخل إحدى حاويات الحقيبة حيث العميلان الرّسميّان مثبَّتان سلفًا.
الحاوية veille-python: بايثون دون بايثون
توفّر الحقيبة حاوية veille-python جاهزة للاستعمال. يُثبِّت ملفّها Dockerfile حزمتَين فحسب: elasticsearch>=9,<10 وneo4j>=5.28,<6. مجلّد python/ الخاصّ بالحقيبة مركَّب على /work: أيّ سكربت python/mon_script.py تكتبه ينفَّذ فورًا داخل الحاوية. متغيّرات البيئة ES_URL وES_USER وES_PASSWORD وNEO4J_URI (bolt://neo4j:7687) وNEO4J_USER وNEO4J_PASSWORD يحقنها docker-compose.yml: يقرؤها السكربت بلا أيّ ترميز صلب.
يكفي أمران اثنان:
./lab.sh python veille.py "climate change" # ينفّذ python/veille.py بوسيطة
./lab.sh python-shell # يفتح مفسِّر بايثون تفاعليًّا
المكافئ على Windows PowerShell هو .\lab.ps1 python veille.py "climate change". لا pip install على الجهاز: هذا مبدأ الحقيبة، وهو ما يجعل الدّورة قابلة للتّكرار من مقعد إلى آخر.
ملفّاتك .py تُوضَع في مجلّد python/ من الحقيبة. تراه الحاوية تحت /work. احفظ الملفّ من محرِّرك، ثمّ أعِد تشغيل ./lab.sh python …: تُؤخَذ النّسخة الجديدة في الحسبان فورًا، دون إعادة بناء الصّورة.
عميل elasticsearch 9 في خمس حركات
يُحاكي عميل بايثون لـElasticsearch واجهة REST. تكتب بايثون اصطلاحيًّا، والمكتبة تُصنّع طلبات HTTP.
الاتّصال
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']} »")
المخرَج المنتظَر:
Elasticsearch 9.5.3 — cluster « veille »
Elasticsearch(url, basic_auth=(u, p)) هي الصّيغة الاعتياديّة في الإصدار 9. يُعيد العميل استعمال تجمُّع اتّصالات HTTP؛ لا تُنشئ عميلًا لكلّ طلب.
البحث
كلّ استعلام Query DSL يمرّ عبر es.search(index=..., query=..., size=..., source=[...]). لاحِظ: بدءًا من العميل 9، تُكتَب الوسائط مباشرةً (query=...)، ولم يعد لازمًا تغليفها في 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"])
المخرَج المنتظَر (قد تتفاوت الدّرجات قليلًا حسب 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
يُقرأ إجماليّ النّتائج في rep["hits"]["total"]["value"] — 2 834 لعبارة climate change.
القراءة والفهرسة والتّحديث
doc = es.get(index="news", id="1") # وثيقة واحدة عبر _id
es.index(index="news", id="200854", document={ # إنشاء أو استبدال
"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"})
تُعيد كلّ دالّة قاموسًا يحوي _id و_version وresult (created، updated، noop).
الفهرسة الجماعيّة: helpers.bulk
إرسال خمسين ألف استدعاء es.index منفصلًا مكلف جدًّا. توفّر الحزمة elasticsearch.helpers الدّالّة bulk التي تُعدّ صيغة NDJSON وتدير المحاولات المُعادة.
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")
المخرَج المنتظَر:
200 documents indexés, 0 erreurs
chunk_size=100 يُرسل الوثائق في حزم من مئة. على مجموعة News، يتّبع مُستورِد الحقيبة (importer/import_news.py) هذا المنطق بالضّبط، بدفعات من ألفين.
إدارة الأخطاء
يعود استثناءان أكثر من غيرهما. لا تدعهما يصعدان بلا رسالة واضحة.
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 رمزَ HTTP 401 «security_exception ... unable to authenticate user [elastic]». يُستعمل NotFoundError عند es.get على _id غير موجود. أمّا TransportError فيغطّي مشكلات الشّبكة.
عميل neo4j 5 في خمس حركات
يتّبع العميل الرّسميّ لـNeo4j نموذجًا مختلفًا قليلًا: سائق (driver)، وجلسة (session)، ومعاملات — ضمنيّة عبر session.run، أو صريحة عبر execute_read وexecute_write.
الاتّصال
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() الاتّصال دون إطلاق استعلام؛ مفيد عند بدء تشغيل تطبيق للفشل فورًا إن عجز السّائق عن الوصول إلى الخادم.
استعلام مع مُعامِلات
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"])
المخرَج المنتظَر (قد يختلف رقمك بحسب المقالات الحديثة):
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
استعمل المُعامِلات دائمًا ($nom)، وليس أبدًا f"…{nom}…": أسرع (يحفظ Cypher الاستعلام في التّخبئة) ويجنّبك أيّ حقن.
المعاملات الصّريحة
للحصول على شيفرة متينة، الممارسة الحسنة هي تغليف الاستعلام في دالّة وتمريرها إلى execute_read أو execute_write. يُدير السّائق المحاولات المُعادة عند انقطاع الشّبكة.
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']}")
المخرَج المنتظَر:
4954 Reuters
2433 Lee Moran
1915 Ron Dicker
1328 Ed Mazza
1145 Cole Delbyck