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,
ReadMessagewith 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:
- 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.
- 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.
- 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.
- 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.