Apache Kafka

trigger_kafka · trigger · Triggers · Disponibile · v1.0.0

Descrizione

Avvia il workflow ogni volta che arriva un messaggio su un topic Apache Kafka. Il runtime apre un consumer con consumer group e commit dell'offset: ideale per pipeline event-driven ad alto throughput dove più servizi pubblicano eventi su un log distribuito e il workflow li elabora. Differenza con i sibling: trigger_rabbitmq = work queue con ack per-messaggio e requeue; trigger_kafka = log distribuito partizionato con offset e consumer group (scalabilità orizzontale, replay dal passato); trigger_websocket = stream push; trigger_webhook = HTTP in ingresso. Scegli Kafka quando hai grandi volumi di eventi, più consumer che devono scalare in parallelo, o vuoi poter rileggere lo storico. Consumer group: il "Group ID" identifica il gruppo di consumer che si spartiscono le partizioni del topic. Più container/istanze con lo STESSO group id scalano il consumo in parallelo (ognuno prende un sottoinsieme di partizioni). Lascialo vuoto per un id dedicato a questo workflow. Consegna at-least-once: l'offset viene committato SOLO dopo che il run è partito con successo. Se il run fallisce, l'offset non avanza e il messaggio viene ri-consumato: nessun evento perso su un crash. "Leggi dall'inizio" ricomincia dall'offset più vecchio disponibile (utile per un backfill iniziale); di default riparte dall'ultimo offset committato del gruppo. Sicurezza: TLS (ssl) e autenticazione SASL (PLAIN o SCRAM-SHA-256/512) per i cluster gestiti (Confluent Cloud, AWS MSK, Aiven, Redpanda). Usa le espressioni {{secrets.X}} per non incollare le credenziali in chiaro nel nodo. Output per ogni messaggio: { data } = payload parsato JSON quando possibile (altrimenti stringa), { raw } = testo originale, { topic }, { partition }, { receivedAt }. Con "JSON Pointer di filtro" processi solo i messaggi che hanno un certo campo. Use case: (1) eventi di dominio da microservizi → proiezione su DB/CRM, (2) click/telemetria ad alto volume → aggregazione → alert, (3) change-data-capture (Debezium) → sync verso sistemi esterni, (4) pipeline di elaborazione ordini con replay in caso di errore.

⚙️ Parametri di configurazione

Campi mostrati nell’editor quando si configura il nodo. Generati direttamente dal NodeDefconfigFields.

CampoTipoRequiredDefaultDescrizione
brokers
Broker (host:porta, separati da virgola)
stringsi
kafka1.example.com:9092, kafka2.example.com:9092
Lista dei bootstrap broker del cluster, separati da virgola. Ne basta uno raggiungibile: il client scopre gli altri. Formato host:porta.
topic
Topic
stringsi
orders.events
Nome del topic da consumare. Il consumer si sottoscrive a tutte le sue partizioni.
groupId
Consumer Group ID
stringno
flowforge-orders-processor
Identifica il gruppo di consumer che si spartiscono le partizioni. Più istanze con lo stesso id scalano in parallelo. Vuoto = id dedicato a questo workflow (flowforge-<workflowId>).
fromBeginning
Leggi dall'inizio (primo avvio)
booleannofalseOn: al primo avvio del gruppo consuma dall'offset più vecchio disponibile (backfill dello storico). Off (default): riparte dall'ultimo offset committato dal gruppo (solo i nuovi messaggi).
ssl
TLS (ssl)
booleannofalseOn per i cluster che richiedono connessione cifrata (quasi tutti i gestiti: Confluent, MSK, Aiven).
saslMechanism
Autenticazione SASL
enum
noneplainscram-sha-256scram-sha-512
nononeMeccanismo SASL. none = nessuna auth (broker aperti / reti fidate). plain / scram-sha-256 / scram-sha-512 per i cluster gestiti. Richiede username e password qui sotto.
saslUsername
SASL username / API key
stringno
{{secrets.KAFKA_KEY}}
Username SASL (o API key per Confluent Cloud). Usa {{secrets.X}} per non incollarlo in chiaro.
saslPassword
SASL password / API secret
stringno
{{secrets.KAFKA_SECRET}}
Password SASL (o API secret). Usa {{secrets.X}} dalle Variabili tenant: non incollare segreti nel nodo.
jsonParse
Parsa i messaggi come JSON
booleannotrueOn (default): il payload viene parsato come JSON in "data" (fallback alla stringa grezza se non è JSON valido). Off: "data" resta la stringa. "raw" contiene sempre il testo originale.
messagePointer
JSON Pointer di filtro/estrazione
stringno
/eventType oppure /payload/status
RFC 6901 JSON Pointer. Se valorizzato, il run parte SOLO se il puntatore risolve a un valore non-undefined (esposto come "matched"). I messaggi senza match avanzano comunque l'offset (scartati per scelta). Vuoto = ogni messaggio fa partire un run.
maxMessagesPerSec
Budget anti-flood (messaggi/sec)
numberno0Tetto di run avviati al secondo. Oltre il budget i messaggi avanzano l'offset senza far partire un run (scartati) per proteggere il runtime. 0 = nessun limite. In Kafka il consumer group + partizioni già distribuiscono il carico.
reconnect
Riconnessione automatica
booleannotrueOn (default): su crash del consumer riconnette con backoff esponenziale (1s→2s→…→30s). Off: alla prima caduta il consumer si ferma finché non riabiliti/salvi il workflow.

💡 Esempio configurazione

Snippet JSON del nodo come compare nel workflow. I valori sono derivati daidefaultValue e dai parametri required.

{
  "id": "node-trigger_kafka-1",
  "defId": "trigger_kafka",
  "label": "Apache Kafka",
  "config": {
    "brokers": "kafka1.example.com:9092, kafka2.example.com:9092",
    "topic": "orders.events",
    "fromBeginning": false,
    "ssl": false,
    "saslMechanism": "none",
    "jsonParse": true,
    "maxMessagesPerSec": 0,
    "reconnect": true
  }
}

🔗 Nodi correlati nella stessa categoria

Pronto a usare Apache Kafka?

Disponibile da subito in tutti i piani FlowForge. Provalo gratis senza carta di credito.

Inizia gratisSfoglia tutti i nodi