Overview
This guide covers the complete technical architecture for processing fleet telematics data - from GPS device to actionable alert. It is intended for engineers building fleet management systems or integrating telematics data into existing platforms.
1. Data Ingestion Layer
1.1 MQTT Broker Setup
GPS devices communicate using MQTT (Message Queuing Telemetry Transport) - a lightweight protocol designed for IoT devices with unreliable connectivity.
Recommended broker: Eclipse Mosquitto (open source) or HiveMQ (enterprise)
Mosquitto configuration (/etc/mosquitto/mosquitto.conf):
listener 8883
cafile /etc/mosquitto/certs/ca.crt
certfile /etc/mosquitto/certs/server.crt
keyfile /etc/mosquitto/certs/server.key
require_certificate true
# Authentication
allow_anonymous false
password_file /etc/mosquitto/passwd
# Persistence
persistence true
persistence_location /var/lib/mosquitto/
# Logging
log_dest file /var/log/mosquitto/mosquitto.log
log_type error warning notice informationTopic structure:
fleet/{company_id}/vehicles/{vehicle_id}/telemetry
fleet/{company_id}/vehicles/{vehicle_id}/alerts
fleet/{company_id}/vehicles/{vehicle_id}/commands1.2 GPS Ping Schema
{
"device_id": "DEV-001234",
"vehicle_id": "VH-MH12AB1234",
"timestamp": "2026-02-05T08:15:30.000Z",
"lat": 18.5204,
"lng": 73.8567,
"altitude": 560.2,
"speed": 45.3,
"heading": 270,
"accuracy": 3.5,
"ignition": true,
"odometer": 45230.5,
"engine_hours": 1205.3,
"fuel_level": 68.5,
"rpm": 2100,
"coolant_temp": 88,
"battery_voltage": 13.8,
"dtc_codes": [],
"signal_strength": -75
}1.3 Message Consumer
// Node.js MQTT consumer
const mqtt = require('mqtt');
const { processTelematics } = require('./processor');
const client = mqtt.connect('mqtts://broker.yourfleet.com:8883', {
clientId: 'telemetry-consumer-01',
cert: fs.readFileSync('./certs/client.crt'),
key: fs.readFileSync('./certs/client.key'),
ca: fs.readFileSync('./certs/ca.crt'),
rejectUnauthorized: true,
reconnectPeriod: 5000,
keepalive: 60
});
client.subscribe('fleet/+/vehicles/+/telemetry', { qos: 1 });
client.on('message', async (topic, payload) => {
try {
const data = JSON.parse(payload.toString());
await processTelematics(data);
} catch (err) {
logger.error('Failed to process telemetry', { topic, err });
// Dead letter queue for failed messages
await deadLetterQueue.push({ topic, payload: payload.toString(), error: err.message });
}
});2. Stream Processing Layer
2.1 Idle Detection Algorithm
// State machine for idle detection
class IdleDetector {
constructor(thresholdMinutes = 5) {
this.threshold = thresholdMinutes * 60 * 1000; // ms
this.idleStart = null;
this.lastPing = null;
}
process(ping) {
const isIdle = ping.ignition === true && ping.speed < 2;
const now = new Date(ping.timestamp);
if (isIdle && !this.idleStart) {
this.idleStart = now;
} else if (!isIdle && this.idleStart) {
const duration = now - this.idleStart;
if (duration >= this.threshold) {
// Emit idle event
return {
type: 'IDLE_END',
vehicle_id: ping.vehicle_id,
start: this.idleStart,
end: now,
duration_minutes: Math.round(duration / 60000),
location: { lat: ping.lat, lng: ping.lng }
};
}
this.idleStart = null;
}
this.lastPing = ping;
return null;
}
}2.2 Harsh Driving Detection
class HarshDrivingDetector {
constructor() {
this.prevPing = null;
}
process(ping) {
if (!this.prevPing) {
this.prevPing = ping;
return null;
}
const timeDelta = (new Date(ping.timestamp) - new Date(this.prevPing.timestamp)) / 1000; // seconds
const speedDelta = ping.speed - this.prevPing.speed; // km/h
const headingDelta = Math.abs(ping.heading - this.prevPing.heading);
const events = [];
// Harsh acceleration: > 10 km/h increase in 3 seconds
if (timeDelta <= 3 && speedDelta > 10) {
events.push({ type: 'HARSH_ACCELERATION', severity: speedDelta > 15 ? 'HIGH' : 'MEDIUM' });
}
// Harsh braking: > 15 km/h decrease in 3 seconds
if (timeDelta <= 3 && speedDelta < -15) {
events.push({ type: 'HARSH_BRAKING', severity: speedDelta < -25 ? 'HIGH' : 'MEDIUM' });
}
// Sharp cornering: > 45 degree heading change at speed > 30 km/h
if (timeDelta <= 3 && headingDelta > 45 && ping.speed > 30) {
events.push({ type: 'SHARP_CORNERING', severity: 'MEDIUM' });
}
// Speeding: > 80 km/h on city roads (configurable per geofence zone)
if (ping.speed > 80) {
events.push({ type: 'SPEEDING', speed: ping.speed, severity: ping.speed > 100 ? 'HIGH' : 'MEDIUM' });
}
this.prevPing = ping;
return events.length > 0 ? events.map(e => ({ ...e, vehicle_id: ping.vehicle_id, timestamp: ping.timestamp, lat: ping.lat, lng: ping.lng })) : null;
}
}3. Geofencing with PostGIS
3.1 Geofence Schema
CREATE EXTENSION IF NOT EXISTS postgis;
CREATE TABLE geofences (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
name VARCHAR(200) NOT NULL,
type VARCHAR(30) CHECK (type IN ('customer_site','depot','restricted','speed_zone','operating_area')),
boundary GEOMETRY(POLYGON, 4326) NOT NULL, -- WGS84 coordinates
speed_limit INTEGER, -- for speed_zone type
company_id UUID NOT NULL,
is_active BOOLEAN DEFAULT true,
created_at TIMESTAMPTZ DEFAULT now()
);
-- Spatial index for fast point-in-polygon queries
CREATE INDEX idx_geofences_boundary ON geofences USING GIST(boundary);3.2 Geofence Entry/Exit Detection
-- Function: check which geofences a point is inside
CREATE OR REPLACE FUNCTION check_geofences(
p_lat FLOAT,
p_lng FLOAT,
p_company_id UUID
) RETURNS TABLE(geofence_id UUID, geofence_name VARCHAR, geofence_type VARCHAR) AS $$
BEGIN
RETURN QUERY
SELECT g.id, g.name, g.type
FROM geofences g
WHERE g.company_id = p_company_id
AND g.is_active = true
AND ST_Contains(
g.boundary,
ST_SetSRID(ST_MakePoint(p_lng, p_lat), 4326)
);
END;
$$ LANGUAGE plpgsql;// Geofence event detection
async function detectGeofenceEvents(ping, prevGeofences) {
const currentGeofences = await db.query(
'SELECT * FROM check_geofences($1, $2, $3)',
[ping.lat, ping.lng, ping.company_id]
);
const currentIds = new Set(currentGeofences.rows.map(g => g.geofence_id));
const prevIds = new Set(prevGeofences.map(g => g.geofence_id));
const events = [];
// Entry events: in current but not in previous
for (const gf of currentGeofences.rows) {
if (!prevIds.has(gf.geofence_id)) {
events.push({ type: 'GEOFENCE_ENTRY', geofence: gf, vehicle_id: ping.vehicle_id, timestamp: ping.timestamp });
}
}
// Exit events: in previous but not in current
for (const gf of prevGeofences) {
if (!currentIds.has(gf.geofence_id)) {
events.push({ type: 'GEOFENCE_EXIT', geofence: gf, vehicle_id: ping.vehicle_id, timestamp: ping.timestamp });
}
}
return { events, currentGeofences: currentGeofences.rows };
}4. Time-Series Storage
4.1 TimescaleDB Schema
-- Enable TimescaleDB
CREATE EXTENSION IF NOT EXISTS timescaledb;
-- Raw telemetry (partitioned by time)
CREATE TABLE telemetry (
time TIMESTAMPTZ NOT NULL,
vehicle_id UUID NOT NULL,
lat FLOAT NOT NULL,
lng FLOAT NOT NULL,
speed FLOAT,
heading SMALLINT,
ignition BOOLEAN,
odometer FLOAT,
fuel_level FLOAT,
rpm INTEGER,
coolant_temp FLOAT
);
SELECT create_hypertable('telemetry', 'time', chunk_time_interval => INTERVAL '1 day');
CREATE INDEX idx_telemetry_vehicle ON telemetry(vehicle_id, time DESC);
-- Retention policy: keep raw data for 90 days
SELECT add_retention_policy('telemetry', INTERVAL '90 days');
-- Continuous aggregate: hourly summaries (kept for 2 years)
CREATE MATERIALIZED VIEW telemetry_hourly
WITH (timescaledb.continuous) AS
SELECT
time_bucket('1 hour', time) AS hour,
vehicle_id,
AVG(speed) AS avg_speed,
MAX(speed) AS max_speed,
SUM(CASE WHEN ignition AND speed < 2 THEN 1 ELSE 0 END) AS idle_pings,
COUNT(*) AS total_pings
FROM telemetry
GROUP BY hour, vehicle_id;5. Alert Engine
5.1 Alert Rules Schema
CREATE TABLE alert_rules (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
company_id UUID NOT NULL,
name VARCHAR(200) NOT NULL,
event_type VARCHAR(50) NOT NULL, -- 'SPEEDING', 'IDLE_END', 'GEOFENCE_EXIT', etc.
conditions JSONB, -- e.g., {"min_duration_minutes": 15, "vehicle_types": ["truck"]}
recipients JSONB, -- e.g., {"roles": ["fleet_manager"], "emails": ["[email protected]"]}
channels VARCHAR[] DEFAULT ARRAY['email','in_app'],
cooldown_minutes INTEGER DEFAULT 60, -- don't re-alert for same vehicle within this period
is_active BOOLEAN DEFAULT true
);5.2 Alert Processing
async function processAlert(event) {
// Find matching alert rules
const rules = await db.query(`
SELECT * FROM alert_rules
WHERE company_id = $1
AND event_type = $2
AND is_active = true
`, [event.company_id, event.type]);
for (const rule of rules.rows) {
// Check cooldown
const recentAlert = await redis.get(`alert:cooldown:${rule.id}:${event.vehicle_id}`);
if (recentAlert) continue;
// Check conditions
if (!matchesConditions(event, rule.conditions)) continue;
// Send alert
await sendAlert(event, rule);
// Set cooldown
await redis.setex(
`alert:cooldown:${rule.id}:${event.vehicle_id}`,
rule.cooldown_minutes * 60,
'1'
);
}
}6. Driver Scoring Algorithm
-- Calculate driver score for a given period
CREATE OR REPLACE FUNCTION calculate_driver_score(
p_driver_id UUID,
p_from TIMESTAMPTZ,
p_to TIMESTAMPTZ
) RETURNS NUMERIC AS $$
DECLARE
v_total_km NUMERIC;
v_harsh_accel INTEGER;
v_harsh_brake INTEGER;
v_cornering INTEGER;
v_speeding INTEGER;
v_score NUMERIC;
BEGIN
SELECT
SUM(distance_km),
SUM(harsh_acceleration_count),
SUM(harsh_braking_count),
SUM(cornering_count),
SUM(speeding_count)
INTO v_total_km, v_harsh_accel, v_harsh_brake, v_cornering, v_speeding
FROM trip_summaries
WHERE driver_id = p_driver_id
AND trip_start BETWEEN p_from AND p_to;
IF v_total_km IS NULL OR v_total_km = 0 THEN
RETURN NULL;
END IF;
-- Score per 100 km (normalised)
v_score := 100
- (v_harsh_accel / v_total_km * 100 * 2)
- (v_harsh_brake / v_total_km * 100 * 3)
- (v_cornering / v_total_km * 100 * 1)
- (v_speeding / v_total_km * 100 * 5);
RETURN GREATEST(0, LEAST(100, ROUND(v_score, 1)));
END;
$$ LANGUAGE plpgsql;7. Infrastructure Sizing
| Fleet Size | GPS Pings/Day | Storage/Month | Recommended Setup |
|---|---|---|---|
| < 50 vehicles | < 720K | ~15 GB | Single PostgreSQL + TimescaleDB server (8 vCPU, 32 GB RAM) |
| 50–200 vehicles | 720K–2.9M | 15–60 GB | Primary + read replica, Redis for caching |
| 200–1000 vehicles | 2.9M–14.4M | 60–300 GB | Kafka for ingestion, TimescaleDB cluster |
| > 1000 vehicles | > 14.4M | > 300 GB | Full streaming architecture (Kafka + Flink) |
*See how IdeaSprout Fleet Management processes telematics data →*