
Kafka Connect: integrazione senza codice
Come usare Kafka Connect per integrare database, file system e servizi esterni senza scrivere consumer/producer.
Cosa imparerai
- Configurare source e sink connector con converter, transform e topic di destinazione
- Attivare la dead letter queue con tolleranza agli errori su ogni sink
- Calibrare tasks.max e monitorare lag e tasso di errore con soglie di escalation
Collegamenti
Kafka Connect: integrazione senza codice
Questa lezione, nel binario ml-tabellare del modulo, ti mostra come portare dati dentro e fuori Kafka senza scrivere un solo producer o consumer: la disciplina resta la stessa, ma il codice diventa configurazione dichiarativa.
L’idea in una frase
Kafka Connect è un runtime dichiarativo che sposta l’integrazione da codice duplicato a connector configurati, osservabili e recuperabili.
La procedura in cinque passi
- Scegli il connector source o sink adatto alla sorgente e alla destinazione.
- Configura converter, transform leggere e topic di destinazione.
- Attiva la dead letter queue con tolleranza agli errori su ogni sink.
- Calibra
tasks.maxsul throughput reale e sulle partizioni disponibili. - Monitora lag e tasso di errore con soglie di escalation esplicite.
Il problema che Connect risolve
Prima o poi un team deve portare in Kafka dati che vivono in un database, in un SaaS o in un sistema operativo, e la tentazione è scrivere un microservizio diverso per ogni fonte. Kafka Connect evita questo con un’integrazione standardizzata e configurabile. La promessa “senza codice” va però presa con cautela: il valore reale dipende da come gestisci configurazione, offset, errori e osservabilità. Questa lezione mostra dove finisce il no code e dove comincia l’ingegneria vera.
Senza Connect, ogni nuova sorgente richiede un consumer o un producer scritto a mano, con retry, gestione errori e tracciamento degli offset. Moltiplica per dieci sistemi e ottieni dieci pezzi di codice quasi identici da mantenere, ognuno con i suoi bug. Connect sposta quel lavoro in una configurazione dichiarativa eseguita da un runtime già pronto.
Leggi un connector come checklist di integrazione: sorgente, destinazione, converter, transform, dead letter queue, retry e lag. Un connector è affidabile quando fallisce in modo osservabile e recuperabile, non quando nasconde la complessità dietro poche righe di configurazione. La domanda da tenere a mente è dove finiscono i dati quando qualcosa va storto.
Source e sink connector
Connect distingue due ruoli. Un source connector legge da un sistema esterno e scrive su Kafka: è il caso di Debezium per il change data capture da PostgreSQL, o di un connector S3 che legge file. Un sink connector fa il percorso inverso, legge da Kafka e scrive su un sistema esterno come ClickHouse, S3, Elasticsearch o una destinazione JDBC.
La maggior parte delle pipeline di analytics combina i due: cattura un cambiamento da un database operativo con un source connector e lo riversa nel data warehouse con un sink. Il punto delicato è quasi sempre la trasformazione e la gestione degli errori nel mezzo, non il trasporto in sé.
Esempio: sink connector ClickHouse
{
"name": "clickhouse-sink",
"config": {
"connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
"topics": "user_events",
"clickhouse.url": "jdbc:clickhouse://localhost:8123",
"clickhouse.table": "user_events",
"tasks.max": "4",
"errors.tolerance": "all",
"errors.deadletterqueue.topic.name": "dlq_user_events"
}
}
Con questa configurazione, Connect consuma automaticamente dal topic user_events, scrive su ClickHouse e dirotta gli eventi problematici sul topic dlq_user_events per ispezione. Il parametro tasks.max controlla quante istanze parallele lavorano, mentre errors.tolerance su all evita che un singolo messaggio malformato fermi l’intero connector.
Dead letter queue
La dead letter queue gestisce gli errori senza bloccare la pipeline. Quando un messaggio non può essere scritto, per esempio per formato non valido o vincolo violato, invece di fermare il connector lo si reindirizza a un topic DLQ. Da lì un processo separato, o un analista, può ispezionarlo e correggere il problema mentre il flusso principale continua a scorrere.
È una scelta che cambia il comportamento sotto stress. Senza DLQ, un evento corrotto può fermare tutta l’ingestione; con la DLQ, l’errore resta circoscritto e visibile. È esattamente quel “fallire in modo osservabile e recuperabile” che separa un’integrazione robusta da una fragile.
Single message transform
Connect supporta trasformazioni leggere sui messaggi senza uno stream processor esterno:
"transforms": "addPrefix",
"transforms.addPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.addPrefix.regex": ".*",
"transforms.addPrefix.replacement": "clickhouse.$0"
Le single message transform vanno bene per operazioni semplici come rinominare un topic o aggiungere un campo. Per logica più articolata, che incrocia più messaggi o mantiene stato, conviene passare a Kafka Streams o a Flink invece di forzare le SMT oltre il loro scopo.
Come impostare la scelta di integrazione
La prima domanda non è “quale connector uso” ma “quale decisione dovrà migliorare grazie a questa integrazione”. Una pipeline ha valore se riduce l’incertezza di una scelta concreta, non se aggiunge un’altra freccia tra due riquadri in un diagramma. Definisci in modo esplicito l’unità di lavoro (il topic, l’evento, lo schema, il connector stesso), il segnale osservato (lag, throughput, tasso di errore, compatibilità dello schema) e la baseline di lettura.
Prima di adottare un connector esistente controlla mapping degli schemi, gestione degli errori, throughput atteso, gestione delle credenziali e possibilità di replay. Un’integrazione che sembra semplice può diventare il punto meno governato dell’intera piattaforma proprio perché nessuno la tratta come codice da presidiare.
Errori comuni
Il primo errore è trattare il “no code” come “no engineering” e non presidiare un connector come faresti con il codice. Il secondo è non controllare la qualità del dato in ingresso: schemi che cambiano, eventi duplicati e timezone incoerenti producono conclusioni false anche a valle di un’integrazione tecnicamente perfetta. Il terzo è confondere correlazione e causalità quando interpreti le metriche di prodotto alimentate dalla pipeline.
Per ridurre questi rischi mantieni tre controlli minimi: una definizione esplicita della metrica di salute, un confronto per fonte e una verifica contro un periodo precedente. E soprattutto una dead letter queue su ogni sink, perché senza un posto dove far cadere gli errori, prima o poi sarà l’intera pipeline a cadere.
Verdetto: usa i connector per source e sink standard con dead letter queue sempre attiva, limita le transform ai ritocchi senza stato e sposta ogni logica con stato su Kafka Streams o Flink.
Riferimenti:
- Confluent. (2024). “Kafka Connect Documentation.” docs.confluent.io.
- ClickHouse. (2024). “Kafka Connect Sink.” clickhouse.com/docs.
L’esempio che fa da riferimento
Debezium è diventato il source connector di riferimento per portare il change data capture da PostgreSQL e MySQL dentro Kafka senza producer scritti a mano. Sul lato opposto, il sink verso ClickHouse documentato da ClickHouse mostra il pattern speculare: consumo continuo dal topic e scrittura analitica con dead letter queue per i record invalidi. Confluent, fondata nel 2014 dagli autori di Kafka, ha fatto di questo runtime dichiarativo uno strato standard della piattaforma. Il filo comune è la lezione del modulo: l’integrazione è affidabile quando fallisce in modo visibile e recuperabile.
Domande per verificare la lezione
- Quale sistema esterno colleghi e usi un source o un sink connector?
- Dove finiscono i record invalidi e chi li ispeziona?
- Quale trasformazione fai dentro Connect e perché non richiede stato?
- Quale metrica di lag o errore ti dice che il connector sta degradando?
Bloccato su questo argomento o vuoi applicarlo al tuo caso? Prenota una call di 15 minuti con un analista esperto.
Percorso collegato
Lezioni da leggere insieme
Questi collegamenti portano la lezione dentro il resto del corso: basi da riprendere, passaggi successivi e connessioni tematiche tra moduli.