Skip to content
Reliable Data Engineering
Practice problem hard iotmqtttime-seriesstreaminglakehouse
Practise with timer, notes and rubric

Design a Connected-Vehicle Telemetry Platform

Problem

A car manufacturer has 5 million connected vehicles. Each vehicle sends sensor signals (speed, battery state, temperatures, GPS, error codes). The company wants:

  1. Near-real-time fault detection (e.g. battery overheating) → alert the driver/service centre within a minute.
  2. Fleet analytics for engineering (battery degradation by model/climate, error-code frequency after software updates).
  3. ML training data (predictive maintenance).
  4. Compliance with GDPR (location is personal data) and regional data residency.

Clarifying questions

QuestionAssumed answer
Signals & frequency?~200 signals; most sampled at 1 Hz on the car, uploaded in batches every 10–60 s; error codes event-driven
Connected at once?~2 M vehicles at peak
Connectivity?Intermittent (tunnels, garages), so data can arrive hours/days late
Real-time needs?Only for a defined set of critical signals
Consent?Location/driving behaviour only with driver consent; varies per market

1. Estimates

Naive: 2M connected cars × 200 signals × 1 Hz = 400M values/s, far too much to ship raw.
With edge filtering (send-on-change, priorities): ~50 values/s per car.
Per upload every 30 s: 30 s × 50 values × ~4 B (delta-encoded, compressed) ≈ 6 KB
2M cars / 30 s ≈ 67k uploads/s × 6 KB ≈ 400 MB/s compressed ingest
  → large: Kafka 512+ partitions, split into regional clusters
Storage: 400 MB/s × 86,400 ≈ 35 TB/day compressed at peak rates → retain raw high-res 30–90 days, downsample for long-term

The numbers force edge processing (send changes/deltas, not every sample) and tiered retention. Say this out loud.

2. Architecture

flowchart LR
    subgraph Vehicle
        ECU[ECUs / CAN bus] --> TCU[Telematics unit<br/>buffer, compress, edge rules]
    end
    TCU -->|MQTT over TLS<br/>mutual cert auth| BRK[MQTT broker fleet<br/>regional]
    BRK --> BRIDGE[Bridge / connector]
    BRIDGE --> K[[Kafka regional<br/>topics by signal group<br/>key = vehicle_id]]
    K --> DEC[Decoder: binary → signals<br/>via signal catalog DBC]
    DEC --> RT[Stream: critical rules<br/>CEP, thresholds]
    RT --> ALERT[Alert service →<br/>app push, service centre]
    DEC --> BR[(Bronze: decoded signals<br/>Delta, cluster by vehicle_id, ts)]
    BR --> SIL[(Silver: cleaned, unit-normalised,<br/>resampled, trips)]
    SIL --> GOLD[(Gold: fleet KPIs, battery health,<br/>DTC frequency, features)]
    GOLD --> ENG[Engineering analytics / BI]
    GOLD --> ML[Predictive maintenance ML]
    CAT[(Signal catalog:<br/>ids, units, scaling, model/SW version)] --> DEC
    CONS[(Consent service)] --> SIL

3. Deep dives

3.1 Ingestion protocol

3.2 Decoding with a signal catalog

Raw payloads are binary frames (CAN signal ids, scaling factors). A versioned signal catalog (per model and software version) maps ids → name, unit, scale, offset. The decoder joins against the catalog (broadcast/cached). Wrong catalog version = garbage data, so catalog version is stored with each row.

3.3 Time-series data layout

Long/narrow vs wide format:

FormatExampleProsCons
Narrow (EAV)vehicle_id, ts, signal_name, valueAny signal, schema-stableHuge row counts, pivot for analysis
Widevehicle_id, ts, speed, soc, batt_temp, ...Easy analytics, compresses wellSparse, schema changes with new signals
HybridNarrow in bronze, wide per signal group in silver (resampled to 1 s / 10 s)Best of bothExtra step

Clustering by (vehicle_id, ts) or (signal_group, date) depending on dominant queries (single vehicle diagnostics vs fleet-wide analysis); often keep both layouts in silver for the two access patterns.

Downsampling for long retention: 1 Hz raw for 30–90 days → 1-min aggregates (min/max/avg/last) forever.

3.4 Real-time fault detection

flowchart LR
    S[Decoded signals] --> KEY[keyBy vehicle_id]
    KEY --> R1["Rule: batt_temp > 60 °C<br/>for 3 consecutive readings"]
    KEY --> R2["Rule: DTC P0A80 appears<br/>after SW update X"]
    KEY --> R3["Model: anomaly score on<br/>cell voltage spread"]
    R1 --> DEDUP[Alert dedup / cooldown<br/>per vehicle × rule]
    R2 --> DEDUP
    R3 --> DEDUP
    DEDUP --> OUT[Alert service]

3.5 Privacy and residency

4. Trade-offs

DecisionChoiceAlternative
Device protocolMQTT → Kafka bridgeHTTPS batches (simpler, less efficient), AWS IoT Core / Azure IoT Hub (managed)
StorageDelta lakehouse + downsamplingTime-series DB (InfluxDB, TimescaleDB): great for recent per-device queries, not for petabyte fleet analytics; often used alongside for the live vehicle view
Real-time scopeOnly critical signalsEverything real-time (cost explodes, no business need)
FormatNarrow bronze + wide silverSingle format for all

5. Failure modes

FailureHandling
Fleet-wide reconnect after outage (thundering herd)Randomised backoff in TCU firmware, broker autoscaling, Kafka buffering
Wrong signal catalog deployedDecoded value range checks per signal → quarantine; reprocess from bronze raw payloads
Clock drift on vehicleServer receive time recorded; correct via GPS time; flag implausible timestamps
Firmware bug floods error codesPer-vehicle rate limiting, anomaly alert on DTC volume per SW version

6. What separates a senior answer

7. Follow-up questions

Engineering asks: "did error code X increase after software update 24.3?" How do you support that?

Gold table of DTC events joined with an SCD2 table of vehicle software versions (effective-dated from OTA update logs): rate of DTC X per 1,000 vehicles per day, before vs after the update, controlling for model/region/age. The SCD2 join (point-in-time software version) is the key modelling element.

How do you build training data for predictive maintenance?

Labels from service/warranty records (component replaced on date D). Features from telemetry windows before D (e.g. 30 days of downsampled signals), point-in-time correct. Negative examples from vehicles without failure at comparable mileage/age. Class imbalance handled; strict time-based train/test split.


Self-assessment rubric