On crée un agent IA local en streaming avec un flux SSE, des filtres simples, puis un LLM local réveillé seulement quand ça vaut le coup. Ici je pars d’un cas concret : surveiller les éditions publiques de Wikipedia et repérer les changements suspects sans cloud ni clé API.
C’est quoi un agent IA local en streaming ?
Un agent IA local en streaming, ce n’est pas juste un chatbot avec une interface plus jolie. C’est un programme qui reste actif, écoute ce qui se passe, trie les signaux, puis déclenche un raisonnement local quand quelque chose vaut vraiment le coup.
Dans ce projet, le mot streaming a deux sens. Le premier, c’est le flux entrant. On consomme en continu des événements publics, par exemple les éditions Wikipedia via EventStreams, le flux temps réel proposé par Wikimedia. À chaque modification d’une page, un événement arrive. L’agent ne pose pas une question, il observe. C’est pour ça qu’on est plus proche d’un agent always-on que d’un chatbot classique.
Le deuxième sens, c’est la sortie. Quand l’agent décide de répondre, il peut renvoyer son analyse progressivement, morceau par morceau. C’est utile si la réponse prend quelques secondes, ou si vous voulez afficher le raisonnement au fil de l’eau dans une interface locale. Mais ce n’est pas obligatoire. Le streaming de sortie est un confort, le streaming d’entrée est le cœur du système.
Entre nous, on le sait bien, faire appel à un consultant en automatisation intelligente et en agent IA, c’est souvent le raccourci le plus malin. On en parle ?
Le point important, et je le vois souvent chez les clients, c’est qu’il ne faut pas lancer un LLM sur chaque événement. C’est coûteux, lent, et souvent inutile. Les filtres simples doivent faire le sale boulot avant. Une règle sur la langue, le namespace Wikipedia, la taille du changement, certains mots-clés, une page sensible, ça suffit déjà à éliminer 95% du bruit. Le LLM local intervient seulement quand la décision devient ambiguë ou intéressante.
Le projet reste léger. Il faut Python 3.11 ou plus, Ollama installé en local, et un modèle disponible comme llama3.1:8b. Côté Python, j’utilise fastapi, uvicorn, httpx, pydantic, ollama et sse-starlette. Pas de clé API. Pas de compte cloud. Le seul point de sortie réseau utile, c’est le flux public Wikimedia.
| Flux entrant | Éditions publiques Wikipedia consommées en continu via Wikimedia EventStreams. |
| Filtres amont | Règles simples pour ignorer le bruit avant d’appeler l’IA. |
| LLM local | Modèle Ollama comme llama3.1:8b, utilisé seulement quand ça vaut le coup. |
| API locale | FastAPI expose l’état de l’agent et les réponses localement. |
| Sortie streaming éventuelle | Réponse renvoyée progressivement avec SSE si l’interface en a besoin. |
Pourquoi filtrer avant d’appeler le LLM ?
Je filtre avant d’appeler le LLM parce qu’un flux public peut sortir plusieurs événements par seconde. Ollama peut tourner très correctement en local, mais ça reste du CPU, du GPU, de la RAM, et surtout du temps. Si je demande au modèle d’analyser toutes les éditions, je fabrique une file d’attente. L’agent ne surveille plus le présent, il commente le passé. Et ça, je l’ai vu souvent dans des projets IA internes : le modèle est bon, mais on l’appelle beaucoup trop tôt.
Mon architecture tient dans un entonnoir simple. D’abord des filtres bêtes, rapides, déterministes. Ensuite seulement, un appel au LLM pour les événements qui ont assez de signal. Un filtre déterministe, c’est une règle qui donne toujours le même résultat pour la même entrée. Pas d’interprétation, pas de coût IA, pas de latence inutile.
Concrètement, je garde seulement les éditions pertinentes, j’écarte les bots quand l’info existe, je repère les utilisateurs anonymes via IPv4 ou IPv6, je regarde les grosses variations de taille, et je laisse passer ce qui ressemble vraiment à un cas à analyser.
# src/schemas.py
from pydantic import BaseModel, Field
from typing import Optional
class RecentChangeEvent(BaseModel):
# Schéma minimal pour normaliser un événement RecentChanges.
id: Optional[int] = Field(default=None)
type: Optional[str] = Field(default=None)
title: Optional[str] = Field(default=None)
user: Optional[str] = Field(default=None)
bot: bool = Field(default=False)
comment: Optional[str] = Field(default="")
namespace: Optional[int] = Field(default=None)
timestamp: Optional[int] = Field(default=None)
length_old: Optional[int] = Field(default=None)
length_new: Optional[int] = Field(default=None)
server_url: Optional[str] = Field(default=None)
# src/filters.py
import re
from .schemas import RecentChangeEvent
IPV4_RE = re.compile(r"^(?:\d{1,3}\.){3}\d{1,3}$")
IPV6_RE = re.compile(r"^(?:[0-9a-fA-F]{0,4}:){2,7}[0-9a-fA-F]{0,4}$")
SUSPICIOUS_COMMENT_RE = re.compile(
r"(blank|remove|delete|spam|test|vandal|insult|lol|haha)",
re.IGNORECASE,
)
def is_anonymous_user(user: str | None) -> bool:
# Un utilisateur anonyme MediaWiki apparaît souvent comme une adresse IP.
if not user:
return False
if IPV4_RE.match(user):
parts = user.split(".")
return all(0 <= int(part) <= 255 for part in parts)
return bool(IPV6_RE.match(user))
def size_delta(event: RecentChangeEvent) -> int:
# Une grosse variation de taille peut signaler une suppression ou un ajout massif.
if event.length_old is None or event.length_new is None:
return 0
return abs(event.length_new - event.length_old)
def is_relevant_edit(event: RecentChangeEvent) -> bool:
# Je garde les vraies éditions sur les pages de contenu.
if event.type not in {"edit", "new"}:
return False
if event.namespace != 0:
return False
if event.bot:
return False
return True
def should_wake_agent(event: RecentChangeEvent) -> bool:
# Le LLM ne doit se réveiller que si plusieurs signaux simples s’alignent.
if not is_relevant_edit(event):
return False
score = 0
delta = size_delta(event)
if is_anonymous_user(event.user):
score += 2
if delta >= 500:
score += 2
elif delta >= 100:
score += 1
if event.length_new == 0 and event.length_old and event.length_old > 0:
score += 3
if event.comment and SUSPICIOUS_COMMENT_RE.search(event.comment):
score += 2
if not event.comment:
score += 1
return score >= 2
Ce genre de filtre ne remplace pas le LLM. Il le protège. Il évite de brûler de la machine sur du bruit, et il garde l’agent en temps réel.
| Signal observé | Coût du filtre | Intérêt pour la détection de vandalisme |
| Événement de type edit ou new | Très faible | Évite d’analyser les logs inutiles |
| Bot déclaré | Très faible | Réduit beaucoup le bruit |
| Utilisateur anonyme IPv4 ou IPv6 | Faible | Bon signal de risque sur les wikis publics |
| Forte variation de taille | Très faible | Repère suppressions et ajouts massifs |
| Commentaire vide ou suspect | Faible | Ajoute du contexte avant l’analyse LLM |
Comment consommer le flux Wikipedia ?
Je consomme Wikimedia comme un vrai flux, pas comme une API classique qu’on interroge toutes les 5 secondes. Wikimedia EventStreams expose du Server-Sent Events, ou SSE. En clair, c’est une requête GET longue durée qui reste ouverte et qui envoie des lignes au fil de l’eau.
Le piège, c’est que toutes les lignes ne servent pas. Certaines sont vides. Certaines sont des commentaires SSE, souvent utilisées comme keep-alive pour garder la connexion ouverte. Les vraies données arrivent généralement derrière le préfixe data:. Sur ce type de flux, la robustesse compte plus que l’élégance. Le stream ne doit pas planter parce qu’une ligne est mal formée. J’ai déjà vu ça chez un client, un simple JSON cassé dans un flux temps réel, et tout le traitement s’arrêtait pendant la nuit. Pas acceptable.
Je crée donc un fichier src/stream_source.py qui lit le flux, filtre les lignes inutiles, parse le JSON, puis normalise les événements exploitables.
import json
from dataclasses import dataclass
from typing import Any, AsyncIterator
import httpx
WIKIMEDIA_RECENT_CHANGES_URL = "https://stream.wikimedia.org/v2/stream/recentchange"
@dataclass(frozen=True)
class RecentChangeEvent:
# Evénement propre utilisé par le reste de l'application.
wiki: str
title: str
change_type: str
user: str
timestamp: int
comment: str | None
url: str | None
def parse_sse_line(line: str) -> dict[str, Any] | None:
# Une ligne vide ne contient rien d'exploitable.
line = line.strip()
if not line:
return None
# En SSE, les lignes qui commencent par ":" sont des commentaires ou keep-alives.
if line.startswith(":"):
return None
# Je ne parse que les lignes data:, le reste ne m'intéresse pas ici.
if not line.startswith("data:"):
return None
raw_json = line.removeprefix("data:").strip()
try:
return json.loads(raw_json)
except json.JSONDecodeError:
# Le flux ne doit jamais tomber à cause d'un JSON mal formé.
return None
def to_event(payload: dict[str, Any]) -> RecentChangeEvent | None:
# Wikimedia envoie plusieurs types de changements.
# Ici je garde seulement les créations et modifications de pages.
change_type = payload.get("type")
if change_type not in {"edit", "new"}:
return None
title = payload.get("title")
wiki = payload.get("wiki")
user = payload.get("user")
timestamp = payload.get("timestamp")
# Si les champs minimums manquent, je préfère ignorer proprement.
if not title or not wiki or not user or not timestamp:
return None
return RecentChangeEvent(
wiki=wiki,
title=title,
change_type=change_type,
user=user,
timestamp=timestamp,
comment=payload.get("comment"),
url=payload.get("meta", {}).get("uri"),
)
async def stream_recent_changes() -> AsyncIterator[RecentChangeEvent]:
# Connexion HTTP persistante au flux public Wikimedia.
timeout = httpx.Timeout(connect=10.0, read=None, write=10.0, pool=10.0)
async with httpx.AsyncClient(timeout=timeout) as client:
async with client.stream("GET", WIKIMEDIA_RECENT_CHANGES_URL) as response:
response.raise_for_status()
# Lecture ligne par ligne, sans charger le flux en mémoire.
async for line in response.aiter_lines():
payload = parse_sse_line(line)
if payload is None:
continue
event = to_event(payload)
if event is None:
continue
yield event
Je couvre aussi les cas bêtes avec des tests unitaires. Ce sont eux qui empêchent le parseur de devenir fragile au premier keep-alive bizarre.
from src.stream_source import parse_sse_line, to_event, RecentChangeEvent
def test_empty_line_returns_none():
assert parse_sse_line("") is None
def test_sse_comment_returns_none():
assert parse_sse_line(": keep-alive") is None
def test_valid_data_line_returns_dict():
line = 'data: {"type":"edit","wiki":"frwiki","title":"IA","user":"Franck","timestamp":123}'
payload = parse_sse_line(line)
assert payload == {
"type": "edit",
"wiki": "frwiki",
"title": "IA",
"user": "Franck",
"timestamp": 123,
}
def test_invalid_data_line_does_not_crash():
assert parse_sse_line("data: {not valid json}") is None
def test_irrelevant_payload_is_ignored():
payload = {
"type": "log",
"wiki": "frwiki",
"title": "Some page",
"user": "Franck",
"timestamp": 123,
}
assert to_event(payload) is None
def test_relevant_payload_becomes_event():
payload = {
"type": "edit",
"wiki": "frwiki",
"title": "Agent IA",
"user": "Franck",
"timestamp": 123,
"comment": "Update",
"meta": {"uri": "https://fr.wikipedia.org/wiki/Agent_IA"},
}
event = to_event(payload)
assert isinstance(event, RecentChangeEvent)
assert event.title == "Agent IA"
assert event.url == "https://fr.wikipedia.org/wiki/Agent_IA"
Avec ça, j’ai une source de streaming simple, async, testable, et surtout tolérante. C’est exactement ce qu’il faut avant de brancher l’agent IA local derrière.
Comment brancher Ollama et diffuser les alertes ?
Je garde la même logique que depuis le début : le flux Wikimedia fournit beaucoup d’événements, les filtres évitent de réveiller l’IA pour rien, Ollama analyse seulement les cas suspects, puis FastAPI expose les alertes utiles en streaming. Côté sortie, j’utilise SSE avec sse-starlette. C’est simple, lisible, et cohérent avec l’entrée : on reçoit un flux, on renvoie un flux.
src/config.py
import os
WIKIMEDIA_STREAM_URL = os.getenv(
"WIKIMEDIA_STREAM_URL",
"https://stream.wikimedia.org/v2/stream/recentchange"
)
OLLAMA_URL = os.getenv("OLLAMA_URL", "http://localhost:11434/api/generate")
OLLAMA_MODEL = os.getenv("OLLAMA_MODEL", "llama3.1:8b")
MIN_COMMENT_LENGTH = int(os.getenv("MIN_COMMENT_LENGTH", "20"))
RECONNECT_DELAY_SECONDS = int(os.getenv("RECONNECT_DELAY_SECONDS", "5"))
src/agent.py
import json
import httpx
from src.config import OLLAMA_URL, OLLAMA_MODEL
async def analyze_event(event):
prompt = f"""
Tu analyses une modification Wikimedia.
Réponds uniquement en JSON valide.
Ne dramatise pas. Tu signales un risque, tu ne prouves pas un vandalisme.
Champs attendus:
suspicion: low, medium ou high
reason: phrase courte
recommendation: action simple
Événement:
title={event.get("title")}
user={event.get("user")}
comment={event.get("comment")}
wiki={event.get("wiki")}
"""
payload = {
"model": OLLAMA_MODEL,
"prompt": prompt,
"stream": False,
"format": "json"
}
async with httpx.AsyncClient(timeout=60) as client:
response = await client.post(OLLAMA_URL, json=payload)
response.raise_for_status()
data = response.json()
result = json.loads(data.get("response", "{}"))
return {
"event": event,
"analysis": {
"suspicion": result.get("suspicion", "low"),
"reason": result.get("reason", "Risque faible ou non déterminé"),
"recommendation": result.get("recommendation", "Surveiller sans action immédiate")
}
}
src/broadcaster.py
import asyncio
class Broadcaster:
def __init__(self):
self.clients = set()
async def subscribe(self):
queue = asyncio.Queue()
self.clients.add(queue)
try:
while True:
yield await queue.get()
finally:
self.clients.remove(queue)
async def publish(self, message):
for queue in list(self.clients):
await queue.put(message)
broadcaster = Broadcaster()
src/main.py
import asyncio
import json
from fastapi import FastAPI
from sse_starlette.sse import EventSourceResponse
from src.agent import analyze_event
from src.broadcaster import broadcaster
from src.config import RECONNECT_DELAY_SECONDS
from src.stream import stream_events
from src.filters import should_wake_agent
app = FastAPI()
async def pipeline():
while True:
try:
async for event in stream_events():
if not should_wake_agent(event):
continue
alert = await analyze_event(event)
if alert["analysis"]["suspicion"] in ["medium", "high"]:
await broadcaster.publish(alert)
except Exception as error:
await broadcaster.publish({
"type": "system",
"message": f"Reconnexion après erreur: {error}"
})
await asyncio.sleep(RECONNECT_DELAY_SECONDS)
@app.on_event("startup")
async def startup():
asyncio.create_task(pipeline())
@app.get("/events")
async def events():
async def event_generator():
async for message in broadcaster.subscribe():
yield {
"event": "alert",
"data": json.dumps(message, ensure_ascii=False)
}
return EventSourceResponse(event_generator())
requirements.txt
fastapi
uvicorn[standard]
httpx
sse-starlette
python-dotenv
.env.example
WIKIMEDIA_STREAM_URL=https://stream.wikimedia.org/v2/stream/recentchange
OLLAMA_URL=http://localhost:11434/api/generate
OLLAMA_MODEL=llama3.1:8b
MIN_COMMENT_LENGTH=20
RECONNECT_DELAY_SECONDS=5
Pour lancer en local, j’installe les dépendances, je vérifie qu’Ollama tourne, je télécharge le modèle, puis je démarre l’API.
pip install -r requirements.txt
ollama serve
ollama pull llama3.1:8b
uvicorn src.main:app --reload
Ensuite j’ouvre http://localhost:8000/events. Je ne cherche pas une machine magique. Je veux juste un agent local qui observe, trie, raisonne un minimum, et garde les données hors d’une API cloud.
Et si votre agent IA local devait rester éveillé toute la journée ?
Un agent IA local en streaming devient vraiment utile quand il ne gaspille pas son intelligence. Je garde le flux en continu, je nettoie et je normalise les événements, je filtre avec des règles simples, puis je réveille le LLM local seulement quand un signal mérite une analyse. Avec Python, FastAPI, SSE, Pydantic et Ollama, on obtient une architecture claire, testable, sans clé API et sans dépendance cloud pour l’inférence. Le vrai bénéfice pour vous, c’est là : vous pouvez construire une surveillance IA continue, maîtrisée, moins coûteuse, et assez robuste pour tourner longtemps.
FAQ
- Un agent IA local en streaming a-t-il besoin d’une API cloud ?
Non, pas pour l’inférence si vous utilisez Ollama avec un modèle installé localement. Dans ce projet, le seul accès réseau utile sert à lire le flux public Wikimedia. L’analyse IA, elle, tourne sur votre machine. - Pourquoi ne pas envoyer chaque événement au LLM local ?
Parce que le flux peut aller vite et qu’un LLM coûte cher en temps de calcul. Je préfère filtrer d’abord avec des règles simples, puis appeler le modèle seulement sur les événements suspects. C’est plus stable, plus rapide et beaucoup plus réaliste. - Quel modèle Ollama utiliser pour ce type d’agent ?
Un modèle local comme llama3.1:8b est un bon point de départ si votre machine suit. L’important n’est pas seulement le modèle, c’est aussi la qualité du tri avant l’appel IA et la façon dont vous structurez la réponse attendue. - À quoi servent FastAPI et sse-starlette dans ce projet ?
FastAPI expose l’application locale, et sse-starlette permet de diffuser les alertes progressivement via Server-Sent Events. Ça colle bien au principe du projet : un flux en entrée, et une sortie exploitable en continu. - Est-ce que ce projet détecte vraiment le vandalisme sur Wikipedia ?
Il repère des éditions qui ressemblent à des cas suspects. Je fais bien la différence : l’agent signale un risque, il ne rend pas un verdict définitif. Pour un usage sérieux, il faut tester, ajuster les filtres, relire les faux positifs et mesurer la qualité des alertes.
A propos de l’auteur
Je suis Franck Scandolera, expert et formateur en tracking avancé server-side, Analytics Engineering, automatisation No/Low Code avec n8n, intégration de l’IA en entreprise et SEO/GEO. Avec mon agence webAnalyste et l’organisme Formations Analytics, j’accompagne des équipes comme Logis Hôtel, Yelloh Village, BazarChic, la Fédération Française de Football ou Texdecor sur des sujets data, IA et automatisation très concrets. Si vous voulez mettre en place ce type d’agent IA local ou automatiser vos flux métier, contactez-moi, je peux vous aider à cadrer et construire proprement.
⭐ Data Analyst, Analytics Engineer et expert dans l’automatisation IA ⭐
Ref clients : Logis Hôtel, Yelloh Village, BazarChic, Fédération Football Français, Texdecor…
Mon terrain de jeu :
Data Analyst & Analytics engineering : tracking propre RGPD, entrepôt de données (GTM server, BigQuery…), modèles (dbt/Dataform), dashboards décisionnels (Looker, SQL, Python).
Automatisation IA des taches Data, Marketing, RH, compta etc : conception de workflows intelligents robustes (n8n, Make, App Script, scraping) connectés aux API de vos outils et LLM (OpenAI, Mistral, Claude…).
Engineering IA pour créer des applications et agent IA sur mesure : intégration de LLM (OpenAI, Mistral…), RAG, assistants métier, génération de documents complexes, APIs, backends Node.js/Python.





