
Ingestion patterns per analytics realtime
Dimensionare batch insert, flush e partizionamento per tempo su ClickHouse, rendere l'ingestione idempotente con chiavi deterministiche e monitorare backpressure e costo query.
Cosa imparerai
- Dimensionare batch insert, flush e partizionamento per tempo sul carico reale
- Rendere l'ingestione idempotente con chiavi deterministiche e retry con backoff
- Monitorare backpressure, parti in attesa e costo query prima di scalare il cluster
Collegamenti
Ingestion patterns per analytics realtime
Il binario di questa lezione è ml-tabellare e il tema è il punto di ingresso: come far entrare milioni di eventi al minuto in ClickHouse senza saturare cluster, storage e reperibilità. La scelta non è tra veloce e lento, ma tra dimensione dei batch, frequenza di flush, partizioni, deduplicazione e backpressure.
Una piattaforma deve ingerire milioni di eventi al minuto e renderli interrogabili senza saturare cluster, storage e reperibilità. La scelta non è tra veloce e lento, ma tra dimensione dei batch, frequenza di flush, partizioni, deduplicazione e backpressure.
Il team subisce spike dopo ogni campagna e vede query lente proprio durante le review operative. La soluzione non è aumentare il cluster. Serve cambiare batch size, schema di partizionamento, politica di retry e percorso di ingestione, così da proteggere latenza e costo al tempo stesso.
Il concetto in una frase
Gli ingestion pattern decidono come batch, buffer, partizioni e deduplicazione portano milioni di eventi al minuto dentro ClickHouse senza saturare cluster e costi.
Il criterio guida e i cinque passi
Il criterio dell’ingestion pattern non si legge dalla scala dei dati ma dalla catena misurata a monte: throughput di picco, cardinalità delle chiavi e latenza accettabile per decisione. Ogni scelta successiva discende da quelle tre misure.
- Misura throughput di picco, cardinalità e latenza. Prima di scegliere qualcosa, campiona il carico reale:
SELECT count() AS eventi, sum(time_sec) AS sec_total
FROM web_events
WHERE ingestion_time > now() - interval 1 day
ORDER BY ingestion_time
SAMPLE 1000000
Annota il throughput e la cardinalità dei principali filtri per fissare max_insert_block_size e index_granularity. Il risultato fissa batch size, cardinalità e latenza ammissibile per ogni tipologia di query.
- Raggruppa gli eventi in batch insert con flush programmati. Ogni insert scrive poi la parte temporaneamente in logs/merges/: il
flush_on_async_insertdecide quando trasformarla in partizione stabile. Conflush_on_async_insert = 0(default) la scrittura va in una coda in memoria e scarica il batch solo quando dimensione o intervallo di tempo lo giustificano:
[merge_tree]
<merge_tree>
<flush_on_async_insert>0</flush_on_async_insert>
<flush_after_inserts>0</flush_after_inserts>
</merge_tree>
Il flush_after_inserts scatta un flush al raggiungimento di un numero fisso di insert (da 1 a 10): conviene usarlo quando le query storiche si aspettano un evento preciso.
- Partiziona per tempo e ordina per chiavi di filtro dominanti. La partizione per timestamp riduce le query storiche alla lettura di poche partizioni; l’ordinamento primario e secondario accelera la scansione. Con alte cardinalità un ordinamento debole espande i merge e posticipa il compaction:
| Ottimizza | Dove la usa |
|---|---|
| Partizione (tempo) | Va a fondo su ingestion_time >= today |
| Ordinamento (chiave filtro) | Index-scan veloce quando il filtro è un’uguaglianza |
-
Rendi ogni scrittura idempotente con chiavi deterministe e retry. L’ingestione in produzione fallisce e ritenta, quindi i duplicati sono inevitabili: le tabelle
MergeTreecon deduplicazione viadeduplicate_keyeliminano le righe equivalenti usando una firma canonicasha1(...). Il percorso di ingestione gestisce quindi retry e backpressure, non l’ultima richiesta vinta: retry con ritardo crescente e backoff esponenziale da 1 a 60 secondi. -
Monitora backpressure, parti in attesa e costo query. Quando l’ingestione eccede la capacità, misura
active_insertsequeue_sizeprima di decidere: aumentare il batch, allentare il compaction o degradare per priorità (tabelle calde e fredde):
| Segnale | Azione |
|---|---|
active_inserts troppo alto | Aumenta flush_after_inserts e max_insert_block_size |
Il problema in produzione
Il collo di bottiglia reale non è la velocità del singolo insert, ma l’accumulo di parti: ogni scrittura asincrona crea una piccola partizione che, senza compaction, moltiplica i file sul disco e rallenta le query di scansione. I guasti tipici sono sempre gli stessi: un parts explosion dopo un picco (migliaia di partizioni in attesa sotto parts_to_merge) che blocca i merge e gonfia lo storage, un merge backlog dove le code di fusione non tengono il ritmo dell’ingestione, un out-of-memory su join perché il join senza indice legge e ordina milioni di righe in RAM, oppure un piano index-only lento perché l’indice coperto non copre le colonne dell’aggregazione.
Le leve di partizionamento non allineate alle query peggiorano tutto: una partizione con chiave troppo larga rallenta ogni merge, una partizione per tempo assente rende ogni query una scansione full-table, e una partizione mal allineata ai filtri domina il costo. Inoltre una finestra di compaction troppo stretta non compatta, una troppo larga congela il disco, e il replica lag su ON CLUSTER lascia le query sulle repliche indietro di minuti rispetto al leader, allungando la latenza percepita. Il verdetto è semplice: le scritture vanno bufferizzate e scritte in batch partizionato per tempo, con deduplicazione deterministica, mentre lo streaming riga per riga resta solo dove il momento di decisione non tollera attese.
ClickHouse non è un database tradizionale: è un OLAP colonnare progettato per ingestione ad alta velocità e letture aggregate, non per transazioni interattive. Ogni colonna è memorizzata in verticale, quindi il motore esegue query per colonna con esecuzione vettoriale e comprime i blocchi interi. Il rovescio è che un join contro milioni di righe senza indice può esaurire la RAM del server. Il contrasto decisivo è con un row-store OLTP: un database transazionale ottimizza insert riga per riga, ma soffre davanti a set aggregati di milioni di righe. La confusione tra i due è l’origine di molte query lente: un JOIN analitico sulla tabella transazionale invece che su una copia compattata porta a out-of-memory e lock.
La trappola più frequente
L’errore più tipico è scambiare familiarità con comprensione, e trattare il framework come risposta invece che come strumento. Se la formalizzazione non lascia spazio a ipotesi, eccezioni e limiti, stai costruendo un rituale invece di una pratica analitica. Il controllo resta operativo: dimensione dei batch, cardinalità, retry, deduplicazione e costo query.
Il secondo errore comune è confondere il costo del compaction con l’esigenza di compattarlo subito. Un OPTIMIZE TABLE ... FINAL dopo ogni ingestione su milioni di righe costa quanto gli insert appena recuperati, perché rilegge e riscrive tutto; la compaction di routine resta in background, e l’intervento finale serve solo per query che soffrono.
Analogamente, la parte indicizzata deve essere quella delle query storiche dominanti, non ogni colonna mai interrogata: altrimenti sprechi storage e raddoppi i tempi di merge. E la partizione per tempo non è un’opzione, ma il fulcro: senza di essa ogni ingestione in batch va a fondo e ogni query soffoca.
Infine il backpressure si ignora: misurare active_inserts e queue_size prima di aggiungere capacità evita di costruire cluster più grandi per un problema di ingestione. La scelta tra materialized view e mutation spesso viene sbagliata: la materialized view accoppia automaticamente la scrittura alla partizione, ma non aggiorna i key-value senza ricostruire lo stato corrente; la mutation fa l’update in tempo, ma scarta i log vecchi e costa I/O proporzionale ai dati interessati.
Checklist operativa
- Fai sempre il background merge di
MergeTree: non compattare mai le tabelle a mano. OPTIMIZEcosta I/O e CPU: non eseguireOPTIMIZE TABLE ... FINALdopo ogni insert.- I rapporti di compaction sono indicativi, in base ai tuoi blocchi e codec.
- I codec zstd, snappy e gzip hanno tradeoff diversi di velocità e rapporto di compressione.
- Backup dei soli dati su
ON CLUSTER: metadati e settings rimangono critici. - Non allentare il replica lag e tieni le repliche aggiornate.
- Evita join su milioni di righe senza
primary keyo indice: è la via più sicura verso l’out-of-memory. - I sorgenti di join devono avere materialized columns.
L’esempio da cui tutto nasce: Yandex.Metrica
Yandex.Metrica, tra i servizi di web analytics più usati al mondo, genera decine di miliardi di eventi al giorno: è per interrogarli che ClickHouse è stato costruito e pubblicato open source da Yandex nel 2016. La documentazione ufficiale del progetto cita volumi oltre i 20 miliardi di righe al giorno su singoli cluster, con tempi di risposta interattivi. Ogni pattern della lezione discende da quella scala: bufferizzare prima di scrivere con batch insert, partizionare per tempo, deduplicare con chiavi deterministiche e monitorare il backpressure invece di comprare cluster più grandi.
Verdetto: le scritture vanno bufferizzate e scritte in batch partizionato per tempo, con deduplicazione deterministica, mentre lo streaming riga per riga resta solo dove il momento di decisione non tollera attese: misura throughput di picco, cardinalità e latenza prima di scegliere qualunque parametro.
Domande per ripassare
- Quando un batch insert protegge latenza e costi rispetto a scritte singole?
- Come scegli parte e ordinamento per le tue query dominanti?
- Come rendi l’ingestione idempotente dinanzi a retry e duplicati?
- Quale segnale di backpressure ti fa degradare per priorità?
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.