Backpressure in Elixir: inviare milioni di messaggi WhatsApp senza crollare
Progettare un motore di invio massivo WhatsApp in Elixir: pipeline Broadway guidate dalla domanda, token bucket per numero, retry con jitter e stati in batch.
Per inviare milioni di messaggi WhatsApp senza crollare, il sistema deve prelevare il lavoro al ritmo che la sua parte più lenta riesce ad assorbire, invece di spingere tutto in una volta. In pratica servono una pipeline guidata dalla domanda che riserva i destinatari da Postgres a piccoli lotti, un token bucket per ogni numero di telefono aziendale che cadenza le chiamate alla Cloud API, errori di throughput che rallentano l’intero numero invece di scatenare una tempesta di retry, e stati di consegna riscritti in batch anziché una riga alla volta. In Turn.io ho costruito il motore di invio massivo ad alto throughput usato per campagne nazionali di salute pubblica che hanno raggiunto 14,7 milioni di persone, con rate limiting, backpressure, tracciamento delle consegne e retry. Questo articolo descrive il design che userei oggi per lo stesso problema, con uno schizzo in Elixir per ogni componente.
In breve
- Tieni lo stato durevole di ogni destinatario in Postgres e tratta gli stage della pipeline in memoria come usa e getta.
- Usa GenStage/Broadway, così è il ritmo dell’API a determinare quanto velocemente leggi dal database, e non il contrario.
- Metti un token bucket per ogni numero di telefono aziendale davanti all’API. Il throughput è un limite per numero, quindi anche il limiter deve esserlo.
- Classifica gli errori: gli errori di throughput mettono in pausa il numero, i limiti per destinatario riprogrammano quel destinatario, gli errori permanenti falliscono subito e quelli ambigui (i timeout) vanno riconciliati, non ritentati alla cieca.
- Fai retry con backoff esponenziale e full jitter.
- I webhook di stato arrivano a un ritmo pari a più volte quello di invio e fuori ordine: applicali in modo monotono e in batch.
- Strumenta tutto con
:telemetryed esporta verso Prometheus. Non puoi ottimizzare ciò che non vedi.
I vincoli con cui devi fare i conti
Tre limiti determinano l’intero design.
Il throughput della Cloud API. Il throughput è per numero di telefono aziendale. Di default la Cloud API consente fino a 80 messaggi al secondo per numero, e i numeri idonei possono essere portati a un limite più alto (fino a 1.000 nel momento in cui scrivo; verifica la documentazione aggiornata di Meta per i tuoi numeri). Se lo superi ricevi errori di throughput (codice di errore 130429). Esiste anche un rate limit per coppia di destinatari (131056) se invii troppi messaggi allo stesso utente in un breve intervallo, e ci sono i messaging limit sul numero di utenti unici che puoi raggiungere con messaggi avviati dall’azienda in una finestra mobile di 24 ore, che dipendono dalla reputazione del tuo account.
Il tuo database. Ogni invio comporta almeno una lettura e una scrittura. Ogni messaggio inviato genera poi dei webhook di stato (sent, delivered, read o failed), quindi le scritture degli stati viaggiano a un multiplo del ritmo di invio. È facile dimensionare tutto sul limite dell’API e dimenticare che il database deve reggere tutto questo.
La capacità condivisa. Probabilmente il tuo numero aziendale gestisce anche conversazioni in corso. Una campagna che usa il 100% del throughput di un numero rende il servizio incapace di rispondere a chi scrive. Lascia margine di proposito.
Architettura
campaign_recipients (Postgres, fonte di verità)
│ riserva N righe (FOR UPDATE SKIP LOCKED), su richiesta
▼
RecipientProducer ──► processors (acquisisce token ► chiama la Cloud API)
│
▼
batcher: scrive i risultati in blocco
Niente viene spinto. I processor chiedono altro lavoro quando sono liberi, il producer riserva solo tante righe quante ne vengono richieste, e il token bucket blocca i processor quando il numero è al limite. Se l’API rallenta, i processor impiegano di più, la domanda cala e il producer smette di leggere dal database. Questa è la backpressure: tutti gli stage rallentano insieme, invece di lasciare crescere le code da qualche parte nel mezzo.
Job Oban o stage in memoria
Un primo design naturale è un job Oban per ogni messaggio. È durevole, ha i retry integrati e funziona fino a un certo punto. Con milioni di messaggi per campagna, però, stai inserendo milioni di righe di job, aggiornando ciascuna più volte mentre passa da uno stato all’altro, e poi facendone pulizia. La tabella dei job diventa la tabella più calda e più gonfia del tuo database, e il controllo del throughput finisce sparso nella configurazione delle code invece di stare in un unico punto.
Funziona meglio dividere le responsabilità:
- Le righe in Postgres sono lo stato durevole: una riga per destinatario, con uno stato, un contatore dei tentativi e un
next_attempt_at. - Oban si occupa dell’orchestrazione: avviare una campagna, programmarla, recuperare periodicamente le righe rimaste bloccate in
sendingdopo un crash, segnare la campagna come completata. Una manciata di job per campagna, non uno per messaggio. - GenStage/Broadway si occupa dell’invio: veloce, in memoria, usa e getta. Se un nodo muore, la pipeline riparte e riprende dalla tabella.
Riservare il lavoro senza contesa
La prenotazione è una singola istruzione. FOR UPDATE SKIP LOCKED permette a più producer (o nodi) di riservare righe in parallelo senza bloccarsi a vicenda e senza prendere la stessa riga:
UPDATE campaign_recipients
SET status = 'sending', claimed_at = now()
WHERE id IN (
SELECT id FROM campaign_recipients
WHERE campaign_id = $1
AND status = 'pending'
AND next_attempt_at <= now()
ORDER BY next_attempt_at
LIMIT $2
FOR UPDATE SKIP LOCKED
)
RETURNING id, phone, template_params, attempt;
Affiancala a un indice parziale, così la prenotazione resta economica anche quando la tabella si riempie di righe completate:
CREATE INDEX campaign_recipients_pending_index
ON campaign_recipients (campaign_id, next_attempt_at)
WHERE status = 'pending';
Un producer guidato dalla domanda
Broadway accetta qualsiasi producer GenStage. Questo tocca il database solo quando gli stage a valle hanno chiesto messaggi, e fa polling quando resta a secco:
defmodule Bulk.RecipientProducer do
use GenStage
@behaviour Broadway.Producer
alias Broadway.Message
@max_batch 500
@poll_interval 1_000
@impl true
def init(opts) do
state = %{campaign_id: Keyword.fetch!(opts, :campaign_id), demand: 0, poll_scheduled?: false}
{:producer, state}
end
@impl true
def handle_demand(incoming, state), do: claim(%{state | demand: state.demand + incoming})
@impl true
def handle_info(:poll, state), do: claim(%{state | poll_scheduled?: false})
defp claim(%{demand: 0} = state), do: {:noreply, [], state}
defp claim(state) do
limit = min(state.demand, @max_batch)
recipients = Bulk.Recipients.claim(state.campaign_id, limit)
messages =
Enum.map(recipients, fn r ->
%Message{data: r, acknowledger: Broadway.NoopAcknowledger.init()}
end)
state = %{state | demand: state.demand - length(messages)}
{:noreply, messages, schedule_poll(state, length(messages) == limit)}
end
# Un lotto pieno significa che probabilmente ci sono altre righe in attesa: nuovo polling subito.
# Altrimenti siamo in pari (o ciò che resta è in backoff): ricontrolliamo più tardi.
defp schedule_poll(%{demand: 0} = state, _full?), do: state
defp schedule_poll(%{poll_scheduled?: true} = state, _full?), do: state
defp schedule_poll(state, full?) do
Process.send_after(self(), :poll, if(full?, do: 0, else: @poll_interval))
%{state | poll_scheduled?: true}
end
end
L’acknowledgement è un no-op perché l’unità di durabilità è la riga del database, non il messaggio. I risultati vengono registrati dal batcher.
Un token bucket per numero di telefono
Broadway ha un’opzione rate_limiting integrata nel producer, e se una pipeline corrisponde a un solo numero di telefono è un buon punto di partenza. In pratica però più campagne e il traffico conversazionale condividono lo stesso numero, quindi voglio un solo limiter per numero, attraverso cui passa tutto ciò che invia da quel numero. Basta un piccolo GenServer:
defmodule Bulk.TokenBucket do
use GenServer
def start_link(opts) do
id = Keyword.fetch!(opts, :phone_number_id)
GenServer.start_link(__MODULE__, opts, name: via(id))
end
@doc "Blocks the caller until this number has capacity for one more send."
def acquire(phone_number_id) do
case GenServer.call(via(phone_number_id), :take) do
:ok ->
:ok
{:wait, ms} ->
Process.sleep(ms)
acquire(phone_number_id)
end
end
@doc "Stops all sends from this number for `ms`, e.g. after a throughput error."
def pause(phone_number_id, ms), do: GenServer.cast(via(phone_number_id), {:pause, ms})
defp via(id), do: {:via, Registry, {Bulk.Registry, {:bucket, id}}}
@impl true
def init(opts) do
rate = Keyword.fetch!(opts, :rate)
{:ok, %{rate: rate, tokens: rate * 1.0, updated_at: now()}}
end
@impl true
def handle_call(:take, _from, state) do
now = now()
state = refill(state, now)
cond do
now < state.updated_at -> {:reply, {:wait, state.updated_at - now}, state}
state.tokens >= 1 -> {:reply, :ok, %{state | tokens: state.tokens - 1}}
true -> {:reply, {:wait, ceil((1 - state.tokens) * 1000 / state.rate)}, state}
end
end
@impl true
def handle_cast({:pause, ms}, state) do
resume_at = max(state.updated_at, now() + ms)
{:noreply, %{state | tokens: 0.0, updated_at: resume_at}}
end
# Durante la pausa updated_at è nel futuro e non si accumulano token.
defp refill(state, now) when now <= state.updated_at, do: state
defp refill(state, now) do
tokens = min(state.rate * 1.0, state.tokens + (now - state.updated_at) * state.rate / 1000)
%{state | tokens: tokens, updated_at: now}
end
defp now, do: System.monotonic_time(:millisecond)
end
Imposta rate sotto il limite reale del numero, per lasciare spazio alle conversazioni in corso. L’attesa avviene nel chiamante, quindi il processo del bucket non si blocca mai. A questi ritmi, una chiamata al GenServer per ogni invio non costa nulla.
Attorno al bucket ci sono due cose da fare bene. La prima: il bucket deve avere un unico proprietario per numero in tutto il cluster. Registry è locale al nodo, quindi o fai girare l’invio di ogni numero su un solo nodo (con un registry globale, o instradando i numeri verso i nodi), oppure sposti il limiter in un punto condiviso. Due nodi, ciascuno con il proprio bucket a 80 al secondo, invieranno allegramente 160 messaggi al secondo. La seconda: la concorrenza deve coprire la latenza. Per la legge di Little, le richieste in volo sono pari al ritmo moltiplicato per la latenza: per sostenere, per esempio, 80 messaggi al secondo con 250 ms di latenza dell’API ti servono almeno 20 processor concorrenti, più un margine per le risposte lente.
La pipeline
defmodule Bulk.Pipeline do
use Broadway
alias Broadway.Message
def start_link(campaign) do
Broadway.start_link(__MODULE__,
name: {:via, Registry, {Bulk.Registry, {:pipeline, campaign.id}}},
producer: [module: {Bulk.RecipientProducer, campaign_id: campaign.id}, concurrency: 1],
processors: [default: [concurrency: 40, max_demand: 5]],
batchers: [default: [concurrency: 1, batch_size: 500, batch_timeout: 1_000]],
context: %{phone_number_id: campaign.phone_number_id}
)
end
@impl true
def process_name({:via, Registry, {registry, key}}, base_name) do
{:via, Registry, {registry, {key, base_name}}}
end
@impl true
def handle_message(_processor, %Message{data: recipient} = message, ctx) do
:ok = Bulk.TokenBucket.acquire(ctx.phone_number_id)
result =
:telemetry.span([:bulk, :send], %{phone_number_id: ctx.phone_number_id}, fn ->
result = WhatsApp.send_template(ctx.phone_number_id, recipient)
{result, %{phone_number_id: ctx.phone_number_id, ok?: match?({:ok, _}, result)}}
end)
Message.put_data(message, {recipient, classify(result, recipient, ctx)})
end
@impl true
def handle_batch(:default, messages, _batch_info, _ctx) do
messages |> Enum.map(& &1.data) |> Bulk.Recipients.record_results()
messages
end
end
Un max_demand basso sui processor è importante. Con una domanda alta, ogni processor accumula molti messaggi che non può ancora inviare perché è in attesa del bucket, e quelle righe restano in sending senza fare nulla.
Errori, retry e backoff con jitter
Non tutti gli errori significano la stessa cosa, e trattarli allo stesso modo è il modo in cui iniziano le tempeste di retry:
# Accettato dall'API: registra l'id del messaggio WhatsApp.
defp classify({:ok, wamid}, _recipient, _ctx), do: {:sent, wamid}
# Throughput raggiunto per questo numero: rallenta l'intero numero, riprova più tardi.
defp classify({:error, %{code: 130429}}, r, ctx) do
Bulk.TokenBucket.pause(ctx.phone_number_id, 1_000)
{:retry, Bulk.Backoff.delay_ms(r.attempt)}
end
# Troppi messaggi a questo singolo utente: riprogramma solo questo destinatario.
defp classify({:error, %{code: 131056}}, r, _ctx), do: {:retry, Bulk.Backoff.delay_ms(r.attempt)}
# Non sappiamo se il messaggio è partito: riconcilia, non reinviare alla cieca.
defp classify({:error, :timeout}, _r, _ctx), do: :unknown
# Tutto il resto è considerato permanente per questo destinatario.
defp classify({:error, error}, _r, _ctx), do: {:failed, error}
E il backoff, con il “full jitter”: un ritardo casuale tra zero e un tetto che cresce in modo esponenziale. Senza jitter, tutti i destinatari falliti nello stesso secondo riprovano nello stesso secondo, e vai di nuovo a sbattere contro il limite tutti insieme.
defmodule Bulk.Backoff do
@base_ms 1_000
@cap_ms 5 * 60_000
def delay_ms(attempt) do
:rand.uniform(min(@cap_ms, @base_ms * Integer.pow(2, attempt)))
end
end
I retry ripassano dalla tabella: riporti la riga a pending, incrementi attempt, imposti next_attempt_at. Dopo un numero massimo di tentativi, la segni come fallita. Siccome i retry sono righe e non processi addormentati, un arretrato di retry non costa nulla in memoria e sopravvive ai deploy.
Idempotenza e il problema del “è partito o no?”
Il caso pericoloso è una richiesta che va in timeout. Meta potrebbe averla accettata, oppure no. Parti dal presupposto che l’API non deduplichi un reinvio al posto tuo, quindi devi decidere cosa è peggio per questa campagna: un messaggio duplicato o uno mancato. Per i promemoria sanitari, di solito il duplicato è il male minore. Per qualsiasi cosa somigli a una richiesta di pagamento, no.
In ogni caso, rendilo riconciliabile. Segna la riga come sending prima di chiamare l’API, e passa il tuo id del destinatario in biz_opaque_callback_data nella richiesta di invio. La Cloud API restituisce quel campo nei webhook di stato, quindi anche se non hai mai salvato l’id del messaggio WhatsApp, uno stato sent o delivered arrivato più tardi ti dice che il messaggio è partito davvero. Un job Oban periodico risolve poi le righe bloccate in sending o unknown: tutto ciò che ha un webhook di stato corrispondente è inviato, tutto ciò che resta muto dopo una finestra generosa riceve la policy scelta per la campagna.
Tracciare gli stati di consegna in batch
I webhook di stato arrivano a un multiplo del ritmo di invio, e non necessariamente in ordine: un read può arrivare prima del suo delivered. Due regole: non far mai tornare indietro un messaggio, e non scrivere mai una riga per ogni webhook.
Assegna un rango agli stati e applica solo gli aggiornamenti che vanno avanti:
CREATE FUNCTION message_status_rank(text) RETURNS int
LANGUAGE sql IMMUTABLE AS $$
SELECT CASE $1
WHEN 'sending' THEN 0 WHEN 'sent' THEN 1 WHEN 'failed' THEN 2
WHEN 'delivered' THEN 3 WHEN 'read' THEN 4
END
$$;
UPDATE campaign_recipients r
SET status = s.status, status_at = s.at
FROM unnest($1::text[], $2::text[], $3::timestamptz[]) AS s(wamid, status, at)
WHERE r.wamid = s.wamid
AND message_status_rank(s.status) > message_status_rank(r.status);
Accumula gli stati in arrivo per al massimo un secondo o qualche centinaio di elementi, riducili in Elixir allo stato di rango più alto per ogni id di messaggio (se lo stesso id compare due volte in un unico UPDATE ... FROM, Postgres applica solo una delle righe corrispondenti, e non puoi scegliere quale), poi esegui un’unica istruzione per l’intero batch. Il receiver dei webhook, da parte sua, deve limitarsi a verificare, salvare e confermare; ho descritto come costruisco quella parte.
Osservabilità
Broadway emette eventi :telemetry per i suoi stage, e il :telemetry.span/3 attorno alla chiamata all’API ti dà eventi di start/stop/exception con le durate. Esportali con telemetry_metrics_prometheus (o PromEx) e costruisci una dashboard con almeno:
- invii al secondo per numero di telefono, confrontati con il ritmo configurato
- tempo passato in attesa in
TokenBucket.acquire(se è sempre alto, il limite è il bucket; se è zero e il throughput è basso, il collo di bottiglia è altrove) - percentili di latenza dell’API e conteggio degli errori per codice
- durata della query di prenotazione e numero di righe
pendingper campagna - tempo tra
sentedelivered, che ti dice qualcosa sul lato dei destinatari, non sul tuo - righe bloccate in
sendingda più di qualche minuto
Imposta alert sugli errori di throughput, non solo sui fallimenti. Un flusso costante di 130429 significa che il ritmo configurato è più alto di quello che Meta ti sta concedendo davvero.
Errori da evitare
- Caricare in memoria l’intero pubblico all’avvio della campagna. Riserva le righe a piccoli lotti.
- Fare retry dentro il processor con
Process.sleep. Tiene in ostaggio un processor e una riga. Rimetti i retry nella tabella. - Lasciare che più nodi condividano un numero senza coordinarsi. I limiti per numero richiedono un proprietario per numero.
- Dimenticare il traffico conversazionale. Quando un milione di persone riceve un messaggio, qualcuno risponde. Il tuo percorso in ingresso ha bisogno di capacità esattamente nel momento di picco della campagna.
Un motore di invio massivo che regge
Un motore di messaggistica massiva è un sistema distribuito, con stato e rate limit, travestito da ciclo su una lista. Se ne stai costruendo uno per WhatsApp, o quello che hai va in difficoltà nei picchi delle campagne, aiuto i team a progettare e costruire sistemi real-time come questo.