KubeOps comme bus : CR = message, reconciler = consumer & Une file de travail réconciliée en C#
Réaliser en .NET le modèle « API server comme bus » : la Custom Resource comme message, le reconciler KubeOps comme consommateur, et une file de travail entièrement réconciliée — avec du YAML, kubectl et du C#/KubeOps.
KubeOps comme bus : CR = message, reconciler = consumer
L’idée en une phrase
KubeOps matérialise en .NET le modèle « l’API server comme bus de messages » (chapitre day 43-44) : une Custom Resource tient lieu de message durable, et un contrôleur IEntityController tient lieu de consommateur abonné à ce type. Le framework fournit l’abonnement (informer), la file interne (work queue) et le réacheminement en cas d’échec (requeue) ; le développeur n’écrit que le traitement — ReconcileAsync — qui réconcilie l’état désiré (spec) vers l’état réel, la boucle observe → diff → agit tenant lieu de consommation de message.
Analogie : Considérons un standard téléphonique d’astreinte doté d’une messagerie. Le standard — le framework — reçoit les appels, enregistre chaque demande comme un message durable, sonne chez le technicien de garde et redéclenche l’appel tant que la demande n’est pas close. Le technicien ne gère ni les lignes ni la file : il traite uniquement la demande courante, telle qu’elle est enregistrée. Si le technicien de garde change, le suivant retrouve les mêmes messages en attente, sans rien perdre.
Points clés
- Le mapping des rôles est direct : le producteur est le
kubectl applyqui crée la CR ; le broker est l’API server adossé à etcd ; le topic est le type CRD ; le message est une instance (Custom Resource) ; le consommateur est la méthodeReconcileAsyncdu contrôleur. Au démarrage, KubeOps enregistre un informer par type contrôlé — l’équivalent exact d’unsubscribe. - L’acquittement n’est pas explicite. Le retour normal de
ReconcileAsyncvaut accusé de traitement pour la version courante de la spec ; une exception vaut nack, et KubeOps réenfile l’objet avec un backoff exponentiel (règle du chapitre day 23-24). L’accusé durable n’est pas un offset committé mais lestatusécrit dans la CR. - Le modèle est level-based, donc sans offset. Un consommateur Kafka gère une position par partition ; un reconciler n’en a aucune : sa « position » est l’état lui-même. KubeOps relit l’objet courant depuis le cache de l’informer à chaque déclenchement, et le resync périodique (chapitre day 15-16) réachemine l’état complet — un événement watch manqué reste sans conséquence.
- Un seul consommateur actif traite un objet donné. La leader election (chapitre day 29-30) garantit qu’une seule instance de l’Operator réconcilie à un instant donné — l’équivalent d’un groupe de consommateurs réduit à un seul membre actif. La work queue interne sérialise par ailleurs les déclenchements d’un même objet (une clé à la fois), tout en autorisant le parallélisme entre objets distincts.
- Le développeur n’écrit que le handler.
AddKubernetesOperator()câble informer, work queue, requeue et leader election ; aucune connexion au broker, aucun offset, aucune file de lettres mortes à gérer : ces responsabilités appartiennent au plan de contrôle.
Exemple concret
Un type EmailRequest est défini comme CRD. Un utilisateur applique une CR welcome-42 portant spec.to et spec.template : le message est publié et persistant. KubeOps, abonné via l’informer, place la clé default/welcome-42 dans sa work queue et invoque ReconcileAsync. Le consommateur envoie le courriel via une passerelle idempotente indexée par le nom de la CR, puis écrit status.phase: Sent et status.observedGeneration: 1. Si la passerelle lève une exception, KubeOps réenfile avec backoff (1 s, 2 s, 4 s…) ; au réessai, l’idempotence évite un second envoi (rappel du chapitre day 45-46). Aucune notion d’offset n’intervient : la preuve du traitement est le status, relu à chaque cycle.
Tableau — consommateur de broker vs reconciler KubeOps
| Concept messaging | Kafka / broker | KubeOps |
|---|---|---|
| S’abonner | subscribe(topic) | Informer enregistré par type au démarrage |
| Message | Enregistrement dans un topic | Custom Resource (objet etcd) |
| Acquittement | Commit d’offset | Retour normal de ReconcileAsync (status = accusé durable) |
| Échec / rejeu | Nack + retry ou dead-letter | Exception → requeue avec backoff (day 23-24) |
| Position | Offset par partition | Aucune : l’état courant, relu (level-based) |
| Un seul consommateur | Groupe de consommateurs | Leader election (day 29-30) |
Code YAML — le message durable posté sur le « bus »
# La Custom Resource EST le message durable publie sur l'API server.
# La spec est la charge utile (l'intention) ; le status portera l'accuse de traitement.
apiVersion: notify.acme.io/v1
kind: EmailRequest
metadata:
name: welcome-42 # sert aussi de cle d'idempotence cote passerelle
spec:
to: alice@example.com
template: welcome
locale: fr
status:
phase: Sent # ecrit par le CONSOMMATEUR (le reconciler), jamais par l'emetteur
observedGeneration: 1 # suit la version de spec deja reconciliee (day 43-44)
Code kubectl — publier, puis observer l’accusé
# Publier le message sur le "bus" (equivalent d'un publish)
kubectl apply -f welcome-42.yaml
# emailrequest.notify.acme.io/welcome-42 created
# Observer l'accuse de traitement ecrit par le consommateur dans le status
kubectl get emailrequest welcome-42 -o jsonpath='{.status.phase}'
# Sent
Code C# / KubeOps — le consommateur abonné au type
// Program.cs — AddKubernetesOperator() cable informers, work queues,
// leader election et requeue. Chaque IEntityController<T> decouvert devient
// un CONSOMMATEUR abonne au type T : KubeOps enregistre un informer (le "subscribe").
var builder = WebApplication.CreateBuilder(args);
builder.Services.AddKubernetesOperator();
var app = builder.Build();
app.Run();
// La Custom Resource EmailRequest est le MESSAGE ; ce controleur est le CONSOMMATEUR.
// KubeOps l'abonne au type via un informer et le declenche par sa work queue.
public class EmailRequestController : IEntityController<V1EmailRequest>
{
private readonly IEmailGateway _gateway;
private readonly IKubernetesClient _client;
public EmailRequestController(IEmailGateway gateway, IKubernetesClient client)
{
_gateway = gateway;
_client = client;
}
public async Task ReconcileAsync(V1EmailRequest msg, CancellationToken token)
{
// Deja traite pour cette version de la spec ? -> rien a refaire (idempotence, day 43-44)
if (msg.Status.ObservedGeneration == msg.Generation())
return;
// AGIR : envoi idempotent, indexe par le nom de la CR (cle d'idempotence).
// Une passerelle idempotente ne renvoie pas deux fois le meme courriel (day 45-46).
await _gateway.SendOnceAsync(idempotencyKey: msg.Name(), msg.Spec, token);
// ACQUITTER dans le STATUS : accuse durable, aucun offset a committer.
msg.Status.Phase = "Sent";
msg.Status.ObservedGeneration = msg.Generation();
await _client.UpdateStatus(msg);
// Une exception non capturee ici -> KubeOps REQUEUE l'objet (nack) avec backoff.
}
}
Piège courant : « Puisque KubeOps est un bus,
ReconcileAsyncreçoit chaque événement une seule fois et dans l’ordre » est inexact. La work queue déduplique les déclenchements d’une même clé et peut coalescer plusieurs changements en un seul appel ; l’appel reçoit l’état courant, non l’historique des transitions. Concevoir le consommateur comme s’il voyait chaque transition — compter les appels, supposer un ordre — reproduit les fragilités d’un broker edge. Le contrôleur doit relire l’état complet et converger.
Une file de travail réconciliée en C#
L’idée en une phrase
Une file de travail peut être réifiée non comme une file de messages transitoires mais comme un ensemble de Custom Resources durables dont le champ status.phase marque l’avancement (Pending → Done ou Failed) ; le reconciler joue le rôle de worker qui vide la file en réconciliant chaque élément vers son état terminal. La « file » n’est alors rien d’autre que l’ensemble des objets non terminaux — une vue level-based, sans offset ni acquittement à outiller.
Analogie : Considérons une bannette « à traiter » sur un bureau. Chaque feuille est une unité de travail ; la file n’est que la pile de feuilles encore présentes. L’employé prend une feuille, exécute la tâche, puis la déplace dans la bannette « fait » : elle quitte la file sans disparaître, restant consultable. Si l’employé est interrompu, la feuille demeure « à traiter » et sera reprise ; rien n’est perdu, car l’avancement est porté par la feuille elle-même, non par la mémoire de l’employé.
Points clés
- La file n’est pas une structure séparée : c’est l’ensemble des CR dont la phase n’est pas terminale. Pour rendre cette file interrogeable, le worker maintient un label miroir de la phase (
workqueue.acme.io/phase) : lister la file revient alors à filtrer par label (kubectl get … -l phase=pending). - Le traitement est idempotent et terminal. Le worker lit
status.phase: un état terminal (Done,Failed) n’est pas retraité. Sinon, il exécute le travail de façon idempotente — indexé par le nom de l’objet — puis écritphase: Done. Rejouer un item déjà traité est sans effet : la file ne redistribue jamais à tort (règle d’or du chapitre day 23-24). - Le réessai est structurel. Un échec du traitement lève une exception → KubeOps réenfile avec backoff (chapitre day 23-24). Un compteur
status.retriespermet de basculer enphase: Failedau-delà d’un seuil : la file de lettres mortes est ainsi réifiée dans lestatus, sans dispositif externe. - La concurrence est maîtrisée. La work queue interne sérialise les cycles par objet (un seul
ReconcileAsyncà la fois pour une clé donnée), tandis que des objets distincts progressent en parallèle. La leader election (chapitre day 29-30) garantit qu’une seule instance draine la file : jamais deux workers sur le même item. - L’ordre n’est pas garanti. Un modèle level-based converge, il n’ordonnance pas. Si un ordre de traitement est requis, il s’encode dans la spec (
priority,notBefore) et se respecte dans le reconciler : l’ordonnancement devient une propriété de l’état désiré, non du transport.
Exemple concret
Une file d’export de rapports est modélisée par des CR ReportJob. Cent objets sont créés, tous en phase: Pending. L’Operator — leader unique — reçoit chacun via l’informer ; pour un objet donné, ReconcileAsync lit la phase : non terminale, il génère le rapport (opération idempotente indexée par nom), puis écrit phase: Done et l’URL du fichier dans le status. Si la génération échoue (délai dépassé du service de stockage), l’exception provoque un requeue avec backoff ; après cinq échecs, phase: Failed et status.message porte la cause — l’item quitte la file active sans bloquer les autres. Un redémarrage de l’Operator ne perd rien : au resync, tous les objets non terminaux sont relus et le drainage reprend. Aucune ligne de file d’attente, d’offset ou de dead-letter n’est écrite : la file est l’ensemble des objets, la reprise est structurelle.
Tableau — file de messages vs file de travail réconciliée
| Aspect | File de messages (broker) | File de travail réconciliée (CRs) |
|---|---|---|
| Élément de travail | Message en file | Custom Resource durable |
| « Dans la file » | En attente de consommation | Phase non terminale (Pending) |
| Retrait | Consommé et acquitté (disparaît) | Passage en phase: Done (l’objet demeure, traçable) |
| Réessai | Redelivery ou dead-letter à outiller | Requeue + backoff (day 23-24), phase: Failed |
| Reprise après crash | Selon rétention et offset | Structurelle : le resync relit les non-terminaux |
| Ordre | FIFO le plus souvent | Non garanti ; encodé dans la spec si requis |
Code YAML — une unité de travail comme Custom Resource
# Chaque unite de travail est une Custom Resource durable.
# La "file" n'est pas une structure a part : c'est l'ensemble des objets non terminaux.
apiVersion: batch.acme.io/v1
kind: ReportJob
metadata:
name: export-2026-07
labels:
workqueue.acme.io/phase: pending # miroir de status.phase, pour interroger la file
spec:
report: monthly-sales
format: pdf
priority: 5 # ordonnancement encode dans l'ETAT DESIRE, pas dans le transport
status:
phase: Pending # Pending -> Done | Failed ; ecrit par le worker
retries: 0
Code kubectl — empiler du travail, voir la file se vider
# Empiler du travail : creer plusieurs items (autant de messages durables)
kubectl apply -f report-jobs/
# reportjob.batch.acme.io/export-2026-06 created
# reportjob.batch.acme.io/export-2026-07 created
# "Voir la file" = lister les objets non terminaux, via le label miroir
kubectl get reportjobs -l workqueue.acme.io/phase=pending
# NAME AGE
# export-2026-06 5s
# export-2026-07 5s
# Observer la file se vider : chaque item converge vers Done
kubectl get reportjobs -o custom-columns=NAME:.metadata.name,PHASE:.status.phase
# NAME PHASE
# export-2026-06 Done
# export-2026-07 Done
Code C# / KubeOps — le worker qui draine la file
// Ce reconciler est le WORKER : il vide la file en reconciliant chaque ReportJob
// vers un etat terminal. Rejoue a chaque cycle, il doit converger sans double effet.
public class ReportJobController : IEntityController<V1ReportJob>
{
private const int MaxRetries = 5;
private readonly IReportExporter _exporter;
private readonly IKubernetesClient _client;
private readonly ILogger<ReportJobController> _logger;
public ReportJobController(
IReportExporter exporter,
IKubernetesClient client,
ILogger<ReportJobController> logger)
{
_exporter = exporter;
_client = client;
_logger = logger;
}
public async Task ReconcileAsync(V1ReportJob job, CancellationToken token)
{
// Un etat terminal ne se retraite pas : le drainage est idempotent.
if (job.Status.Phase is "Done" or "Failed")
return;
try
{
// AGIR de facon idempotente : "ensure" indexe par le nom de la CR.
// Rejouer apres un crash ne produit pas un second rapport (day 45-46).
var url = await _exporter.EnsureReportAsync(job.Name(), job.Spec, token);
job.Status.Phase = "Done";
job.Status.Url = url;
job.Status.ObservedGeneration = job.Generation();
await _client.UpdateStatus(job); // l'"accuse" durable
}
catch (Exception ex)
{
// Echec transitoire : compter l'essai et laisser KubeOps requeue avec backoff.
job.Status.Retries += 1;
job.Status.Message = ex.Message;
if (job.Status.Retries >= MaxRetries)
{
job.Status.Phase = "Failed"; // "dead-letter" reifie dans le status
await _client.UpdateStatus(job);
_logger.LogError(ex, "ReportJob {Name} en echec definitif", job.Name());
return; // sortie propre : plus de requeue
}
await _client.UpdateStatus(job);
throw; // exception -> requeue + backoff (day 23-24)
}
}
}
Piège courant : « Passer
status.phaseàRunningen tête de la boucle verrouille l’item et empêche tout double traitement » est un raisonnement fragile. Le champphaseest un indicateur, pas un verrou : rien n’empêche techniquement une relecture de retrouverRunningaprès un crash, et la correction ne repose pas sur lui. Ce qui garantit l’absence de double effet, c’est la leader election — un seul worker actif — et l’idempotence de l’action, indexée par nom (chapitres day 29-30 et day 45-46). Traiterphasecomme un verrou reproduit le raisonnement de verrou distribué que le modèle level-based rend inutile : on ne verrouille pas la file, on converge, et l’action est rejouable. Ainsi se referme le fil rouge du livre : partout, un état désiré, un état réel, et une boucle qui les réconcilie.