Project Argos: Exploring GPS Ingestion with Kafka and Go

Exploring micro-batching, consumer behavior, and EMA smoothing in a local GPS pipeline, including the delivery and ordering gaps exposed by its design.

Scope
Local experiment · Simulated GPS data
What I built
A Node.js gateway, Kafka, a Go processor applying EMA smoothing, and MongoDB latest-state storage.
Engineering question
What happens to delivery and ordering across a multi-stage pipeline?

View the complete source code and Docker setup on GitHub

Argos is a local learning project exploring how WebSocket ingestion, a broker, concurrent processing, and a latest-state database interact under simulated GPS traffic.

The questions were practical: where should batching happen, what happens when a worker stops, and how does simple smoothing change noisy coordinates?

The repository demonstrates these mechanisms. It does not establish production capacity, positioning accuracy, ETA improvements, or a no-loss delivery guarantee.

Implementation review — September 12, 2026: Corrected earlier reliability, benchmark, and smoothing claims. The separate implementation remains unchanged by this editorial update.

Why separate ingestion and processing?

The pipeline uses a Node.js gateway, Kafka, a Go processor applying exponential moving average (EMA) smoothing, and MongoDB upserts. Splitting these concerns lets me experiment with ingestion and processing independently, at the cost of more operational complexity.

Fig 1. A recorded worker-interruption experiment. Continued gateway activity does not by itself prove end-to-end delivery.

Batching in the Node.js gateway

I chose Node.js for WebSocket ingestion because its asynchronous I/O model fits this gateway. That choice alone does not prove a particular connection capacity or superiority over other runtimes.

The gateway buffers messages in memory and sends when the buffer reaches 5,000 entries. This is threshold-based batching, not immediate forwarding. The following abbreviated excerpt omits connection setup and error handling; consult the repository for the implementation.

// Node.js / Ingestion Gateway (Micro-Batching Implementation)
let messageBuffer = [];
const BATCH_SIZE = 5000;

wss.on('connection', (ws) => {
  ws.on('message', async (message) => {
    const rawData = JSON.parse(message);
    messageBuffer.push({
      key: rawData.truck_id,
      value: JSON.stringify(rawData),
    });

    // Await the send promise after clearing the local buffer; not fire-and-forget
    if (messageBuffer.length >= BATCH_SIZE) {
      const batchToSend = [...messageBuffer];
      messageBuffer = [];

      await producer.send({
        topic: 'raw-telematics-data',
        messages: batchToSend,
      });

      console.log(
        `[Ingestion] 📦 Dispatched batch of ${batchToSend.length} spatial coordinates to Kafka.`
      );
    }
  });
});

The reviewed gateway has no periodic timed flush. Low traffic can leave a partial batch waiting. A crash loses unsent memory, and a failed send after clearing the buffer is logged without an application-level requeue. Loss can include an entire in-flight batch, not only the 4,999 messages below the threshold.

Retaining messages in Kafka

Kafka separates ingestion from worker consumption and can retain broker-acknowledged records while workers are unavailable. End-to-end durability still depends on producer acknowledgements, broker replication and retention, offset commits, and successful database writes. A broker is not an automatic zero-loss guarantee.

Processing coordinates in Go

The implementation uses EMA, not a Kalman Filter. For each coordinate, the update combines the new observation with the previous smoothed value using a fixed alpha. It does not implement a motion model or covariance-based state estimation.

The Go consumer launches processing goroutines and upserts by truck ID. The reviewed code does not bound in-flight work with a worker pool. Consumer-group scale-out is limited by the available partitions, and the small EMA calculation does not establish that CPU is the bottleneck.

// Simplified processing excerpt; consumer and offset lifecycle are omitted
func processTelemetry(ctx context.Context, payload []byte, coll *mongo.Collection, alpha float64) {
    var rawData TelematicsData
    if err := json.Unmarshal(payload, &rawData); err != nil {
        log.Printf("[Worker] Invalid payload: %v", err)
        return
    }

    // EMA smoothing, with per-truck state managed by the helper
    cleanData := applySpatialSmoothing(rawData, alpha)

    // Real-Time State: Upsert instead of Append-Only
    filter := bson.M{"truck_id": cleanData.TruckID}
    update := bson.M{"$set": cleanData}
    opts := options.Update().SetUpsert(true)

    // Count only successful writes; this does not establish message delivery
    if _, err := coll.UpdateOne(ctx, filter, update, opts); err != nil {
        log.Printf("[Worker] Database write failed: %v", err)
        return
    }

    // Thread-safe atomic counter for logging
    current := atomic.AddUint64(&processedCount, 1)
    if current%500 == 0 {
        log.Printf("[Worker] Completed %d upserts", current)
    }
}

Choosing latest-state storage

SetUpsert(true) updates state by truck_id rather than preserving every event as a history record. That is a storage-model choice, not evidence of query speed.

An atomic update does not guarantee the newest event wins: concurrent goroutines can complete out of order, and the reviewed update has no event-time or sequence guard. A unique truck-ID index, per-key ordering, and a policy for stale updates and state recovery need explicit validation. Query latency has not been measured here.


Local Experiments and Evidence Limits

The recordings below illustrate local runs. They are not a reproducible end-to-end benchmark or proof of production resilience.

The original demo setup used a MacBook Pro with an M1 Pro and Docker Desktop. A repeatable result would also need the exact commit, dependency versions, payload distribution, resource limits, duration, and raw measurements.

1. Simulated Load

Fig 2. A simulator run targeting 10,000 telemetry messages per second; the target is not a verified sustained database-write rate.

  • Purpose: Explore threshold-based micro-batching under simulated traffic.
  • Scope: The simulator's configured emission target differs from accepted, broker-acknowledged, and successfully persisted throughput. Messages per second also differs from concurrent WebSocket connections or unique trucks.
  • Evidence gap: This article does not provide reconciled event counts, repeated runs, p95/p99 end-to-end latency, consumer lag, or raw CPU results. I therefore do not claim a sustained 10,000-message-per-second capacity or a measured batching improvement.

2. Worker Interruption

Fig 3. Intentional worker outage simulated via Chaos Testing.

  • Purpose: Observe gateway and consumer behavior while stopping and restarting the worker.
  • What the recording illustrates: Gateway activity continues during a worker interruption, and worker processing resumes afterward.
  • What it does not establish: Every accepted event was persisted exactly once. In the reviewed worker, ReadMessage with a consumer group commits offsets before the spawned processing goroutine completes; a later write failure or crash can leave a committed event unpersisted. See kafka-go's ReadMessage contract.

3. Smoothing: Intended Use, Not Validated Business Impact

EMA can reduce abrupt changes between observations but also introduces lag. A visually smoother path is not proof of greater positioning accuracy. No ground-truth route, ETA baseline, or false-alarm evaluation is included here.

Question

Intended Use

Evidence Still Needed

Coordinate smoothing

Reduce visible coordinate jitter using EMA

Ground-truth position error and smoothing lag

Logistics planning

Explore whether processed coordinates help downstream planning

ETA and false-alarm comparisons; no improvement established


Engineering trade-offs

These are implementation limits, not assumptions that the risks are acceptable for a real fleet:

  1. Batching and lag: Threshold batching and consumer backlog affect freshness. There is no measured sub-second latency guarantee, and low-volume traffic needs a timed-flush policy.
  2. Local infrastructure: The reviewed Compose setup runs one Kafka broker, ZooKeeper, and MongoDB; gateway and processor run separately on the host. This is not a fully containerized application or a high-availability cluster. Persistence and restart behavior need explicit configuration and testing.
  3. Delivery lifecycle: Unsent and in-flight batches can be lost. Offset handling must coordinate successful persistence with commits, including earlier offsets still in flight; merely adding explicit commits is not enough with unordered processing.
  4. Ordering and recovery: Bound concurrency, preserve processing order per truck, reject stale updates, and define how EMA state survives restart or partition reassignment before claiming reliable latest-state behavior.

The next experiment

Publish the test commit and configuration, workload and duration, emitted/accepted/acknowledged/persisted counts, errors, lag, and end-to-end latency percentiles. Reconcile event IDs or offsets using an audit sink: a latest-state collection cannot establish delivery history from its final document count. Test failure boundaries separately, including producer failure, worker crash after reading, database errors, and restart.


Source and related work

The gateway, processor, Compose configuration, and simulator are available for review. This project is useful as an inspectable architecture experiment; stronger operational claims need additional implementation work and reproducible evidence.

Explore the Argos data pipeline repository on GitHub

Career Pipeline gives this kind of experiment a practical place in my learning workflow: choose one question, investigate it, and attach the result as evidence. For a smaller personal data workflow with different constraints, see Pocket CFO.