← Resources/GuideFleet Management

Fleet Telematics Data Pipeline: Engineering Guide from GPS Ping to Actionable Alert

5 February 2026·15 min read

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 information

Topic structure:

fleet/{company_id}/vehicles/{vehicle_id}/telemetry
fleet/{company_id}/vehicles/{vehicle_id}/alerts
fleet/{company_id}/vehicles/{vehicle_id}/commands

1.2 GPS Ping Schema

json
{
  "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

javascript
// 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

javascript
// 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

javascript
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

sql
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

sql
-- 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;
javascript
// 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

sql
-- 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

sql
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

javascript
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

sql
-- 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 SizeGPS Pings/DayStorage/MonthRecommended Setup
< 50 vehicles< 720K~15 GBSingle PostgreSQL + TimescaleDB server (8 vCPU, 32 GB RAM)
50–200 vehicles720K–2.9M15–60 GBPrimary + read replica, Redis for caching
200–1000 vehicles2.9M–14.4M60–300 GBKafka for ingestion, TimescaleDB cluster
> 1000 vehicles> 14.4M> 300 GBFull streaming architecture (Kafka + Flink)

*See how IdeaSprout Fleet Management processes telematics data →*

fleet telematicsGPS dataMQTTdata pipelineIoTfleet analyticsengineering