EIP (Enterprise Integration Patterns)

Les patterns d’intégration enterprise implémentés dans GoCamel.

Split

Divise un message en plusieurs parties traitées individuellement.

builder.From("direct:start").
    Split(func(e *gocamel.Exchange) (any, error) {
        body := e.GetIn().GetBody().(string)
        return strings.Split(body, ","), nil
    }).
    Log("Partie: ${body}").
    To("direct:process").
    End()

Propriétés de l’Exchange:

  • CamelSplitIndex — Index actuel (0-based)
  • CamelSplitSize — Nombre total de parties
  • CamelSplitComplete — Dernière partie ?
Ces propriétés sont limitées à la branche

Elles décrivent une partie, pas le message entier, et sont retirées au moment où le résultat agrégé est publié sur l’exchange parent. Un processeur placé après le Split ne voit donc pas un CamelSplitComplete=true résiduel de la dernière partie, et un Split imbriqué n’écrase pas la comptabilité de celui qui l’englobe. La même règle s’applique à CamelMulticast*, CamelRecipientList* et CamelRoutingSlip*.

Un corps []byte n’est pas une collection

Découper un corps []byte produit une seule partie, comme Apache Camel qui ne découpe pas les byte[]. Le traiter comme une liste produisait un exchange par octet — une charge utile de 10 Ko devenait 10 000 exchanges. Pour découper délibérément du contenu binaire, découpez sur une valeur décodée :

Split(func(e *gocamel.Exchange) (any, error) {
    return strings.Split(string(e.GetBody().([]byte)), "\n"), nil
})

Avec agrégation :

type StringJoinStrategy struct{}

func (s *StringJoinStrategy) Aggregate(oldEx, newEx *gocamel.Exchange) *gocamel.Exchange {
    if oldEx == nil {
        return newEx
    }
    old := oldEx.GetIn().GetBody().(string)
    new := newEx.GetIn().GetBody().(string)
    oldEx.GetIn().SetBody(old + "," + new)
    return oldEx
}

strategy := &StringJoinStrategy{}

builder.From("direct:start").
    Split(splitter).Strategy(strategy).
        To("direct:process").
    End()

Aggregate

Combine plusieurs messages en un seul.

strategy := &MyAggregationStrategy{}
repo := gocamel.NewMemoryAggregationRepository()

builder.From("direct:start").
    Aggregate(gocamel.NewAggregator(correlationExpr, strategy, repo).
        SetCompletionSize(3)). // Completer après 3 messages
    Log("Agrégation terminée: ${body}")

Stratégie personnalisée:

type MyAggregationStrategy struct{}

func (s *MyAggregationStrategy) Aggregate(
    oldExchange, 
    newExchange *gocamel.Exchange
) *gocamel.Exchange {
    // Logique de fusion
    return oldExchange
}

Conditions de complétion :

MéthodeDescription
SetCompletionSize(n)Compléter après n messages
SetCompletionTimeout(ms)Compléter après un délai (timeout)
SetCompletionPredicate(fn)Compléter quand le prédicat retourne vrai

Options de stockage :

// En mémoire (par défaut)
repo := gocamel.NewMemoryAggregationRepository()

// Persistance SQLite
repo := gocamel.NewSQLAggregationRepository(db, "schema")

Multicast

Envoie une copie à plusieurs destinations.

builder.From("direct:start").
    Multicast().
        Pipeline().
            Log("Branche 1: ${body}").
            To("direct:out1").
        End().
        Pipeline().
            Log("Branche 2: ${body}").
            To("direct:out2").
        End().
    End()

Options:

  • ParallelProcessing() — Exécute les branches en parallèle
  • AggregationStrategy — Fusionne les résultats

Propriétés:

  • CamelMulticastIndex
  • CamelMulticastSize
  • CamelMulticastComplete

Wire Tap

Envoie une copie de l’exchange à un endpoint secondaire de manière asynchrone (fire-and-forget), sans bloquer la route principale. La copie est indépendante : les modifications sur l’exchange tapé n’affectent pas la route principale et vice-versa.

builder.From("direct:start").
    WireTap("direct:audit"). // Copie asynchrone vers audit
    To("direct:result")

Fonctionnement :

  1. Crée une copie profonde (deep copy) de l’exchange.
  2. Lance un goroutine qui envoie la copie à l’endpoint spécifié.
  3. Retourne immédiatement — la route principale continue sans attendre.
  4. L’exchange tapé utilise un contexte indépendant (context.Background()) : il n’est pas annulé quand la route principale se termine.
  5. Les erreurs lors du tap sont loggées mais ne sont pas propagées à la route principale.

Notes :

  • Idéal pour l’audit, la journalisation ou la surveillance sans impact sur les performances.
  • Route.Stop() attend que tous les taps en cours soient terminés avant de s’arrêter (aucun message n’est perdu lors du shutdown).
  • Disponible également dans les contextes Multicast, Pipeline, Split et Loop.

Pipeline

Chaîne séquentielle de processors.

builder.From("direct:start").
    Pipeline().
        Log("Étape 1").
        Transform(transformer).
        Log("Étape 2").
        To("direct:end").
    End()

Choice

Content-Based Router - route les messages vers différentes destinations selon des conditions sur le contenu.

builder.From("direct:start").
    Choice().
        When("${header.priority == 'high'}").
            SimpleSetBody("🔴 Priority: ${body}").
            To("direct:urgent").
        When("${header.priority == 'medium'}").
            SimpleSetBody("🟡 Priority: ${body}").
            To("direct:normal").
        When("${body['count'] > 100}").
            To("direct:large-batch").
        Otherwise().
            To("direct:default").
    EndChoice()

Expressions supportées :

ExpressionDescriptionExemple
${header.name}Valeur d’en-tête${header.Content-Type == 'json'}
${body}Corps du message${body == 'active'}
${body['key']}Accès par clé map${body['status'] == 'pending'}
${body[0]}Accès par index${body[0] > 50}
${exchangeProperty.prop}Propriété Exchange${exchangeProperty.userId != ''}
${date:now}Date/heure actuelle
${random(100)}Nombre aléatoire
${uuid}UUID unique

Opérateurs de comparaison :

  • == — Égal à
  • != — Différent de
  • > — Supérieur à
  • < — Inférieur à
  • >= — Supérieur ou égal
  • <= — Inférieur ou égal

Méthodes disponibles dans When/Otherwise :

// Définir le corps/header
When(expression).SetBody(value)
When(expression).SetHeader(key, value)
When(expression).SimpleSetBody("template ${body}")
When(expression).SimpleSetHeader(key, "template ${header.name}")

// Chainer des processors
When(expression).
    Log("Message").
    SetHeader("X-Processed", "true").
    To("direct:output")
Ordre d’évaluation

Les clauses When sont évaluées dans l’ordre. La première condition vraie déclenche son processor, les autres sont ignorées. Si aucune condition ne correspond et qu’Otherwise est présent, il est exécuté.


Filter

Filtre les messages selon une condition. Les exchanges qui ne correspondent pas sont silencieusement supprimés ( ErrStopRouting).

// Expression Simple Language (comparaison à l'intérieur de ${...})
builder.From("direct:start").
    Filter("${header.type == 'order'}").
    To("direct:commandes").
    Build()

// Prédicat Go
builder.From("direct:start").
    Filter(func(e *gocamel.Exchange) (bool, error) {
        body, ok := e.GetBodyAsString()
        return ok && body != "", nil
    }).
    To("direct:nonVide").
    Build()

Throttle

Limite le débit de traitement à N messages maximum par fenêtre de temps. Les exchanges excédentaires sont retardés ( jamais perdus) jusqu’à ce que la fenêtre se réinitialise.

// Maximum 100 messages par seconde
builder.From("direct:start").
    Throttle(100, time.Second).
    To("direct:aval").
    Build()

// Maximum 10 messages par 500ms (fenêtre plus fine)
builder.From("direct:start").
    Throttle(10, 500*time.Millisecond).
    To("direct:api").
    Build()

Notes :

  • Utilise un compteur à fenêtre fixe — pour un lissage plus fin, réduire la durée de la fenêtre
  • Bloquant : le goroutine attend que la fenêtre autorise le passage
  • maxRequests et timePeriod doivent être supérieurs à zéro ; les valeurs invalides retournent une erreur au lieu de bloquer indéfiniment
  • Si le contexte de l’exchange est annulé pendant l’attente de la fenêtre suivante, Process retourne une erreur

Delay

Met le traitement de l’exchange en pause pour une durée fixe ou calculée dynamiquement. Le sleep respecte l’annulation du contexte.

// Délai fixe
builder.From("direct:start").
    Delay(500 * time.Millisecond).
    To("direct:suivant").
    Build()

// Délai dynamique depuis un en-tête
builder.From("direct:start").
    DelayFunc(func(e *gocamel.Exchange) (time.Duration, error) {
        ms, ok := e.GetHeaderAsInt("X-Delay-Ms")
        if !ok {
            return 0, nil  // pas de délai si l'en-tête est absent
        }
        return time.Duration(ms) * time.Millisecond, nil
    }).
    To("direct:suivant").
    Build()

Notes :

  • Bloquant : le goroutine dort pendant la durée calculée
  • Si le contexte de l’exchange est annulé pendant le sleep, Process retourne une erreur
  • Retourner 0 depuis DelayFunc pour ignorer le délai pour un exchange donné
  • Passer nil à DelayFunc fait échouer la construction de la route ; NewDelayerFunc(nil) reste constructible mais retourne une erreur de traitement au lieu de paniquer

Content Enricher

Enrichit l’exchange courant avec des données provenant d’une source externe.

// Avec stratégie d'agrégation
type FusionStrategy struct{}

func (s *FusionStrategy) Aggregate(original, enriched *gocamel.Exchange) *gocamel.Exchange {
    origBody, _ := original.GetBodyAsString()
    enrichBody, _ := enriched.GetBodyAsString()
    original.GetIn().SetBody(origBody + " | " + enrichBody)
    return original
}

builder.From("direct:start").
    Enrich("http://api.example.com/data", &FusionStrategy{}).
    To("direct:suivant").
    Build()

// Sans stratégie — remplace le corps par la réponse d'enrichissement
builder.From("direct:start").
    Enrich("direct:lookup", nil).
    To("direct:suivant").
    Build()

Fonctionnement :

  1. Une copie de l’exchange courant est envoyée à l’URI d’enrichissement
  2. La réponse est récupérée
  3. AggregationStrategy.Aggregate(original, réponse) fusionne les deux
  4. Sans stratégie, le corps de la réponse remplace le corps original

Idempotent Consumer

Empêche le traitement de messages en double.

repo := gocamel.NewMemoryIdempotentRepository()

builder.From("direct:start").
    IdempotentConsumer("${header.MessageId}", repo).
    To("direct:suivant").
    Build()

Fonctionnement :

  1. Évalue l’expression pour obtenir une clé unique.
  2. Vérifie si la clé existe déjà dans le dépôt (repository).
  3. Si elle n’est pas présente, ajoute la clé immédiatement (bloquant les doublons concurrents) et continue le traitement.
  4. Si elle est déjà présente, arrête le routage (ErrStopRouting).
  5. Si l’exchange échoue ensuite en aval, la clé est retirée du dépôt (via une Synchronization de l’exchange) : une relivraison du même message sera donc traitée au lieu d’être perdue — sémantique at-least-once.

Dépôts (Repositories) :

  • NewMemoryIdempotentRepository() : Stockage en mémoire vive (par défaut).
  • NewRedisIdempotentRepository() : Stockage distribué via Redis.
  • Des dépôts basés sur SQL peuvent être implémentés via l’interface IdempotentRepository.

Circuit Breaker

Protège les routes contre les pannes en cascade lors de l’appel de services externes instables.

builder.From("direct:start").
    CircuitBreaker().
        FailureThreshold(5).           // Ouvre après 5 échecs
        OpenTimeout(30 * time.Second).  // Reste ouvert pendant 30s
        SuccessThreshold(2).           // Ferme après 2 succès en état Half-Open
        Process(unstableProcessor).
        OnFallback().
            Log("Circuit ouvert ou échec, appel du fallback").
            To("direct:fallback").
        End().
    End().
    To("direct:resultat").
    Build()

États :

  • Closed (Fermé) : Les requêtes passent normalement. Les échecs sont comptabilisés.
  • Open (Ouvert) : Les requêtes échouent immédiatement ou passent par le fallback.
  • Half-Open (Semi-ouvert) : Après le OpenTimeout, une seule requête d’essai à la fois sonde le service aval ; les requêtes concurrentes sont redirigées vers le fallback comme si le circuit était encore ouvert. Quand SuccessThreshold sondes réussissent, le circuit se referme. Si une sonde échoue, il s’ouvre à nouveau.

Load Balancer

Répartit les messages sur plusieurs endpoints.

builder.From("direct:start").
    LoadBalance().RoundRobin().
        To("direct:serveur1").
        To("direct:serveur2").
        To("direct:serveur3").
    End().
    Build()

Stratégies :

  • RoundRobin() : Distribution circulaire (par défaut).
  • Random() : Distribution aléatoire.
  • Custom(strategy) : Fournir une implémentation personnalisée de l’interface LoadBalancerStrategy.

SEDA

Staged Event-Driven Architecture. Découple les producteurs et les consommateurs via une file d’attente interne (Go channel).

// Consommateur avec workers concurrents
builder.From("seda:traitement?concurrentConsumers=5&size=1000").
    Log("Traitement asynchrone : ${body}").
    To("direct:suivant")

// Producteur
builder.To("seda:traitement")

Fonctionnement :

  1. Le producteur envoie une copie de l’exchange dans une file d’attente interne.
  2. Le producteur rend la main immédiatement (asynchrone).
  3. Un ou plusieurs workers lisent la file et traitent l’exchange de manière indépendante.

Paramètres :

  • size : La capacité de la file d’attente (défaut 1000).
  • concurrentConsumers : Le nombre de goroutines esclaves traitant la file (défaut 1).

Les valeurs invalides (size négatif, concurrentConsumers non positif) retombent sur les valeurs par défaut avec un avertissement. Les paramètres sont appliqués à la première création de l’endpoint ; les références ultérieures au même nom de file réutilisent l’endpoint existant.


Resequencer

Remet dans l’ordre les messages arrivés de manière désordonnée en se basant sur un numéro de séquence ou un horodatage.

builder.From("direct:start").
    Resequence("${header.seq}", 1). // Commence à la séquence 1
        Timeout(500 * time.Millisecond).
        Log("Traitement dans l'ordre : ${header.seq}").
        To("direct:suivant").
    End()

Fonctionnement :

  1. Évalue l’expression pour obtenir le numéro de séquence.
  2. Si la séquence est celle attendue, elle est traitée immédiatement.
  3. S’il s’agit d’une séquence future, elle est mise en tampon et triée.
  4. Si un trou (gap) persiste plus longtemps que le Timeout, le resequenceur saute la séquence manquante et continue avec la suivante disponible.

Notes :

  • Les messages sont délivrés en aval via le pipeline interne du resequenceur (les étapes entre Resequence(...) et End()). L’exchange original s’arrête là (ErrStopRouting) : les étapes placées après End() ne voient jamais de messages hors ordre.
  • La livraison est sérialisée pour garantir l’ordre ; les étapes en aval ne doivent pas rerouter vers le même resequenceur.
  • Les numéros de séquence dupliqués ou obsolètes sont ignorés (avec un avertissement dans les logs).

Claim Check

Réduit la taille des messages en stockant le contenu (payload) dans un dépôt externe et en ne faisant circuler qu’une référence (une clé).

repo := gocamel.NewMemoryClaimCheckRepository()

builder.From("direct:start").
    // Stocke le corps et le remplace par une clé
    ClaimCheck(gocamel.ClaimCheckPush, "mon-ticket", repo).
    To("direct:autre-systeme"). // Seule la clé "mon-ticket" est envoyée
    // Restaure le corps original
    ClaimCheck(gocamel.ClaimCheckPop, "mon-ticket", repo).
    To("direct:suivant")

Opérations :

  • ClaimCheckPush : Stocke le corps et le remplace par une clé. Si aucune clé n’est fournie, elle est générée.
  • ClaimCheckPop : Récupère les données par clé, restaure le corps et les supprime du dépôt.
  • ClaimCheckGet : Récupère les données par clé et restaure le corps (les garde dans le dépôt).
  • ClaimCheckDiscard : Supprime les données du dépôt sans les restaurer.

Stop

Arrête le routage actuel sans erreur.

builder.From("direct:start").
    ProcessFunc(func(e *gocamel.Exchange) error {
        if shouldStop(e) {
            e.SetProperty("stopped", true)
            return nil
        }
        return nil
    }).
    Stop(). // Arrête ici si condition remplie
    Log("Jamais atteint si stop")

ToD (Dynamic To)

URI résolue dynamiquement à l’exécution.

builder.From("direct:start").
    SetHeader("CamelFileName", "data.txt").
    ToD("file://output/${header.CamelFileName}")

Expressions supportées:

  • ${header.<name>} — Valeur d’en-tête
  • ${property.<name>} — Valeur de propriété
  • ${body} — Corps du message

Recipient List

Envoie un message à plusieurs destinataires calculés dynamiquement à l’exécution.

// Destinataires depuis un en-tête
builder.From("direct:start").
    SetHeader("recipients", "direct:a,direct:b,direct:c").
    RecipientList("${header.recipients}")

Dynamic Router

Route un exchange à travers une séquence d’endpoints résolus dynamiquement en appelant une fonction de contrôle après chaque étape. Le routage s’arrête quand la fonction retourne une chaîne vide ou nil.

// Utilisation d'une fonction Go
builder.From("direct:start").
    DynamicRouter(func(e *gocamel.Exchange) (string, error) {
        // Décide de la prochaine étape selon l'état/contenu
        val, _ := e.GetProperty("step")
        count, _ := val.(int)
        e.SetProperty("step", count + 1)
        
        if count == 0 { return "direct:a", nil }
        if count == 1 { return "direct:b", nil }
        return "", nil // Fin du routage
    })

// Utilisation du Simple Language
builder.From("direct:start").
    DynamicRouter("${header.nextEndpoint}")

Fonctionnement :

  1. Appelle la fonction ou l’expression de contrôle.
  2. Si le résultat n’est pas vide, envoie l’exchange à cet URI.
  3. Propage la sortie vers l’entrée et recommence à l’étape 1.

Routing Slip

Route un exchange à travers une séquence dynamique d’endpoints. Contrairement à Recipient List, chaque étape reçoit la sortie de l’étape précédente.

// Routing slip depuis un en-tête
builder.From("direct:start").
    RoutingSlip("${header.routingSlip}")

// Fonction Go
builder.From("direct:start").
    RoutingSlip(func(e *gocamel.Exchange) ([]string, error) {
        return []string{"direct:valider", "direct:transformer", "direct:envoyer"}, nil
    })

Propriétés Exchange pendant Routing Slip :

PropriétéTypeDescription
CamelRoutingSlipIndexintIndex de l’étape courante (base 0)
CamelRoutingSlipSizeintNombre total d’étapes
CamelRoutingSlipCompleteboolIndique la dernière étape

Headers & Properties

Set/Remove Headers

// Définir
builder.SetHeader("X-Correlation-ID", uuid.New().String())
builder.SetHeaders(map[string]any{
    "X-Client-Version": "1.0",
    "X-Request-Time": time.Now(),
})

// Supprimer
builder.RemoveHeader("X-Temp-Data")
builder.RemoveHeaders("X-Debug*", "X-Debug-Rare") // Wildcard

Set/Remove Properties

// Définir
builder.SetProperty("correlationId", "abc-123")
builder.SetPropertyFunc("status", func(e *gocamel.Exchange) (any, error) {
    return e.GetIn().GetHeader("status"), nil
})

// Supprimer
builder.RemoveProperty("temp-token")
builder.RemoveProperties("cache-*")

Gestion des erreurs

Do-Try-Catch-Finally

Gestion structurée des erreurs au sein d’une route.

builder.From("direct:start").
    DoTry().
        To("http://service-instable").
    DoCatch("connection refused").
        Log("Le service est hors ligne, application du mode dégradé").
        To("direct:recovery").
    DoCatch(""). // Rattrape toutes les autres erreurs
        Log("Erreur inattendue : ${exchangeProperty.CamelExceptionCaught}").
        To("direct:error-handler").
    DoFinally().
        Log("Traitement terminé (succès ou échec)").
    End().
    To("direct:suivant")

Fonctionnement :

  1. Exécute les processeurs du bloc DoTry.
  2. Si une erreur survient, recherche un bloc DoCatch correspondant (une chaîne vide intercepte tout).
  3. Si une correspondance est trouvée, l’erreur est considérée comme “gérée” et le routage continue après le bloc DoTry (sauf si le bloc catch échoue lui-même).
  4. Le bloc DoFinally est toujours exécuté, qu’il y ait eu une erreur ou non.

Process

Exécute un processeur personnalisé.

// Fonction simple
builder.ProcessFunc(func(e *gocamel.Exchange) error {
    body := e.GetIn().GetBody().(string)
    e.GetOut().SetBody(process(body))
    return nil
})

// Struct implémentant Processor
builder.Process(&MyCustomProcessor{config: cfg})

Loop

Répète un segment de route un nombre de fois statique ou dynamique.

// Nombre d'itérations fixe (statique)
builder.From("direct:start").
    Loop(3).
        ProcessFunc(func(e *gocamel.Exchange) error {
            // Évalué 3 fois
            idx, _ := e.GetPropertyAsInt(gocamel.CamelLoopIndex)
            size, _ := e.GetPropertyAsInt(gocamel.CamelLoopSize)
            fmt.Printf("Traitement de l'itération %d sur %d\n", idx, size)
            return nil
        }).
    End().
    To("direct:output")

// Nombre d'itérations dynamique via une expression Simple Language
builder.From("direct:start").
    LoopSimple("${header.LoopCount}").
        To("direct:process-item").
    End()

// Nombre d'itérations dynamique via une fonction Go
builder.From("direct:start").
    LoopFunc(func(e *gocamel.Exchange) (int, error) {
        val, _ := e.GetHeaderAsInt("X-Max-Retries")
        return val, nil
    }).
        To("direct:retry-attempt").
    End()

Propriétés de l’Exchange pendant la boucle :

PropriétéTypeDescription
CamelLoopIndexintIndex de l’itération actuelle (commence à 0)
CamelLoopSizeintNombre total d’itérations de la boucle
CamelLoopCompleteboolBooléen indiquant s’il s’agit de la dernière itération

Saga

Le pattern Saga fournit un mécanisme pour coordonner une série de transactions locales à travers des services distribués sans verrouiller les bases de données. Il définit un ensemble d’actions et leurs actions de compensation correspondantes (annulations).

Si une action dans le bloc Saga échoue ou retourne une erreur, l’orchestrateur de la Saga déclenche automatiquement les étapes de compensation pour toutes les actions terminées, dans l’ordre inverse (LIFO - Last In, First Out).

builder.From("direct:saga-route").
    Saga().
        // Action 1: Réserver Hôtel. Compensation 1: Annuler Hôtel.
        Action("direct:reserve-hotel", "direct:cancel-hotel").
        // Action 2: Réserver Vol. Compensation 2: Annuler Vol.
        Action("direct:reserve-flight", "direct:cancel-flight").
    End().
    To("direct:saga-completed")

Vous pouvez également utiliser des fonctions Go personnalisées pour les actions et les compensations directement avec ActionFunc :

builder.From("direct:saga-route").
    Saga().
        ActionFunc(
            func(e *gocamel.Exchange) error {
                // Réserver l'hébergement
                return nil
            },
            func(e *gocamel.Exchange) error {
                // Annuler/compenser la réservation
                return nil
            },
        ).
    End()

Transformer

Transforme le contenu du message.

// Définir le corps directement
builder.From("direct:start").
    SetBody("Hello World")

// Définir le corps via une fonction
builder.From("direct:start").
    SetBodyFunc(func(e *gocamel.Exchange) (any, error) {
        input := e.GetIn().GetBody().(string)
        return strings.ToUpper(input), nil
    })

// Transformer via Simple Language
builder.From("direct:start").
    SimpleSetBody("Processed: ${body} at ${date:now}")

Transform

Un EIP dédié à la transformation de message qui modifie dynamiquement le corps (body) du message à l’aide d’une fonction Go ou d’une expression Simple Language.

// Transformation via une fonction Go personnalisée
builder.From("direct:start").
    Transform(func(e *gocamel.Exchange) (any, error) {
        input, _ := e.GetBodyAsString()
        return "Transformed: " + input, nil
    }).
    To("direct:output")

// Transformation via une expression Simple Language
builder.From("direct:start").
    TransformSimple("Modified: ${body} and type is ${header.Type}").
    To("direct:output")

Message History

L’EIP Message History permet de retracer le parcours chronologique d’un message lorsqu’il traverse les différents nœuds et processeurs de la route.

Pour activer le suivi d’une route, appelez .EnableMessageHistory() sur votre route builder. Vous pouvez ensuite récupérer et inspecter l’historique des messages à l’aide de gocamel.GetMessageHistory(exchange).

// 1. Activer le suivi dans la définition de la route
builder.From("direct:start").
    EnableMessageHistory().
    SetBody("Hello GoCamel").
    TransformSimple("Modified: ${body}").
    To("direct:output")

// 2. Récupérer l'historique dans un processeur en aval
builder.From("direct:output").
    ProcessFunc(func(e *gocamel.Exchange) error {
        history := gocamel.GetMessageHistory(e)
        for i, entry := range history {
            fmt.Printf("[%d] Route : %s | Processor : %T | Time : %s\n", 
                i, entry.RouteID, entry.Processor, entry.Timestamp.Format(time.RFC3339))
        }
        return nil
    })

Récapitulatif des patterns EIP

PatternCatégorieDescription
ChoiceRoutageRoutage basé sur le contenu
FilterRoutageFiltrage conditionnel
Idempotent ConsumerRoutageEmpêche les messages en double
Circuit BreakerRoutagePattern de résilience pour les pannes
Load BalancerRoutageRépartition de charge entre endpoints
Dynamic RouterRoutageRoutage dynamique basé sur décision
ResequencerRoutageRéordonnancer les messages par séquence
Claim CheckRoutageStocker le payload et passer une référence
ThrottleRoutageLimitation de débit
DelayRoutagePause avant la prochaine étape
MulticastRoutageDestinations multiples
Recipient ListRoutageRoutage dynamique vers plusieurs destinataires
Routing SlipRoutageRoutage séquentiel dynamique
Wire TapRoutageCopie asynchrone fire-and-forget
SplitTransformationDécoupage de message
AggregateTransformationAgrégation de messages
Content EnricherTransformationEnrichissement avec données externes
TransformTransformationTransformation de contenu
ToDEndpointEndpoint dynamique
StopContrôleArrêter le routage
DoTryContrôleGestion d’erreurs Try-Catch-Finally
SetHeaderEn-têtesManipulation d’en-têtes
SetPropertyPropriétésPropriétés d’exchange
LoopContrôleItération avec compteur statique/dynamique
SagaContrôlePattern transactionnel Saga
Message HistoryDiagnosticsSuivi du parcours de traitement des messages