From MQTT to BigQuery: an IoT telemetry pipeline that survives production
Most IoT tutorials stop at "the message arrived." Production starts after that: the same message arrives twice, the device clock is wrong, somebody deploys a firmware that sends Celsius where Fahrenheit was agreed, and the bill for streaming inserts quietly triples.
This is the pipeline we build for client fleets: MQTT at the edge, a subscriber that treats validation as a first-class job, and BigQuery as the analytical home. By the end you will have the full path — and the failure handling that keeps it alive.
The architecture
devices → MQTT broker → subscriber (batches, validates) → BigQuery
↓
dead-letter topic (bad payloads)
Three rules shape everything:
- The broker is not a database. MQTT retains at most one message per topic — history has exactly one reliable destination: the warehouse.
- The subscriber owns semantics. Devices are expensive to update; the subscriber is where contracts are enforced and old firmware versions are tolerated.
- At-least-once is the only honest delivery guarantee. Design the table so duplicates are survivable, then dedupe at query time.
1. Fix the topic and payload contract
Topics carry routing metadata; the payload carries measurements:
telemetry/{orgId}/{deviceId}/readings
{
"recordedAt": "2026-09-20T08:14:03Z",
"temperatureC": 4.2,
"humidityPct": 61.8,
"seq": 18211
}
Two things are worth more than any clever encoding: recordedAt is the device's recording time (arrival time
lies under retry), and seq is a monotonic per-device counter that makes duplicates detectable.
2. Batch in the subscriber
Streaming one row per message to BigQuery is the fastest way to spend your budget. Batch by size and time —
whichever flushes first:
import mqtt from 'mqtt';
import { BigQuery } from '@google-cloud/bigquery';
const bq = new BigQuery();
const table = bq.dataset('telemetry').table('readings');
const BATCH_SIZE = 500;
const FLUSH_MS = 5_000;
let buffer: ReadingRow[] = [];
let flushTimer: NodeJS.Timeout | null = null;
function enqueue(row: ReadingRow) {
buffer.push(row);
if (buffer.length >= BATCH_SIZE) {
void flush();
} else if (flushTimer === null) {
flushTimer = setTimeout(() => void flush(), FLUSH_MS);
}
}
async function flush() {
if (flushTimer) clearTimeout(flushTimer);
flushTimer = null;
if (buffer.length === 0) return;
const rows = buffer;
buffer = [];
try {
await table.insert(rows, { raw: false });
} catch (err) {
// Partial failures carry row indexes — re-queue only those.
const partial = (err as any)?.errors;
if (Array.isArray(partial)) {
const failed = partial.map((e: any) => rows[e.index]);
buffer.unshift(...failed);
} else {
buffer.unshift(...rows); // whole batch failed; retry wholesale
}
}
}
const client = mqtt.connect(process.env.MQTT_URL!, { reconnectPeriod: 2_000 });
client.on('message', (topic, payload) => {
const row = parseAndValidate(topic, payload); // returns null on garbage
if (row) enqueue(row);
});
client.subscribe('telemetry/+/+/readings', { qos: 1 });
The parseAndValidate gate is where bad firmware gets contained. Send rejected payloads to a dead-letter topic
instead of discarding them — silent data loss is worse than a noisy queue.
3. Shape the table for the queries you will run
Every telemetry query is time-bounded, so partition by day and cluster by the fields you filter on:
CREATE TABLE telemetry.readings (
org_id STRING,
device_id STRING,
recorded_at TIMESTAMP,
received_at TIMESTAMP,
temperature_c FLOAT64,
humidity_pct FLOAT64,
seq INT64
)
PARTITION BY DATE(recorded_at)
CLUSTER BY org_id, device_id;
Partitioning is what makes retention a one-liner — PARTITION BY DATE(recorded_at) plus a table-level expiry means
raw readings age out on schedule:
ALTER TABLE telemetry.readings
SET OPTIONS (partition_expiration_days = 90);
4. Query with duplicates in mind
At-least-once delivery means the same (device_id, seq) can appear more than once. Dedupe at read time with a
window function, and make it a view so nobody has to remember:
CREATE VIEW telemetry.readings_deduped AS
SELECT * EXCEPT (rn) FROM (
SELECT
*,
ROW_NUMBER() OVER (PARTITION BY device_id, seq ORDER BY received_at) AS rn
FROM telemetry.readings
)
WHERE rn = 1;
A cold-chain report on top of that is plain SQL — hourly averages with excursion flags:
SELECT
device_id,
TIMESTAMP_TRUNC(recorded_at, HOUR) AS hour,
AVG(temperature_c) AS avg_temp,
MAX(temperature_c) AS max_temp,
COUNTIF(temperature_c > 8.0) AS excursion_samples
FROM telemetry.readings_deduped
WHERE recorded_at >= TIMESTAMP_SUB(CURRENT_TIMESTAMP(), INTERVAL 24 HOUR)
GROUP BY device_id, hour
ORDER BY hour DESC;
5. The failure checklist
Before this pipeline carries anything that matters:
- Duplicate burst: replay a day of messages twice. Row counts must double,
readings_dedupedmust not. - Clock skew: inject one device with a date three days off. It should land in its own partition and be visible to nothing else — queries are time-bounded, so skewed rows stop contaminating dashboards.
- Contract break: publish a payload with a missing field. The subscriber must dead-letter it and keep processing the topic.
- Subscriber restart: kill it under load. QoS 1 redelivers; the size+time flush means nothing waits forever.
Why this shape
The pipeline is boring on purpose. MQTT does what it is good at (cheap pub/sub at the edge), the subscriber does what it is good at (policy, batching, validation), and BigQuery does what it is good at (cheap storage, fast time-bounded scans). Every interesting failure — duplicates, skew, bad firmware — is handled where it can be handled cheaply, not in the device and not in the dashboard.