From Real-Time Ingestion to Cold Analytics: Designing for Data Usage

Real-time ingestion does not mean real-time consumption. Data temperature and retention should be decided by how data is used, not by how fast it arrives. Illustrated with our tyre usage and wear data lake architecture.

From Real-Time Ingestion to Cold Analytics: Designing for Data Usage

A single tyre telemetry event can end up in very different places. Within minutes, it can warn a fleet manager that a tyre is losing pressure. The next day, it feeds a per-trip aggregate used in a study. Months later, an R&D engineer may pull it again to compare wear across fleets.

Same event, same ingestion flow but the consumption requirements are different.

This is the founding principle of our tyre usage and wear data architecture, and the rest of this article follows its consequences: on temperature, retention, storage and governance. Before choosing storage technologies or processing engines, we need to understand the uses they will support.

Start with consumption requirements

Data temperature describes how actively data needs to be accessed and the trade-offs we can make to serve it economically.

We distinguish three dimensions:

Dimension Key question What it decides
Freshness How far behind the source events can the data be? How the store is fed
Access frequency How often will consumers read it? Whether to keep it always on
Response time How long can a consumer wait for a result? How the store must perform

These dimensions are related, but they do not answer the same question. Access frequency and response time decide how data must be served; freshness decides how the store that serves it is fed.

Historical data queried continuously through an interactive application may still require a hot serving layer. Events arriving in real time may be retained for occasional analysis without requiring an always-on system.

Data age alone does not determine data temperature.

In practice, our consumers sit at very different points on these dimensions. A few operational services need data within seconds; most analysts and data scientists can wait hours or days, but expect complete and consistent history.

Consumers of our tyre usage and wear data lake, positioned on these dimensions

Ingestion velocity: it's all about usage

Ingestion methods vary: some systems ingest batches, others use real-time streaming protocols like MQTT or tools like Kafka. 

The key principle is to decouple ingestion velocity from data temperature. Real-time ingestion doesn't require real-time consumption.

In our tyre usage and wear platform, truck IoT devices send telemetry through MQTT and Kafka. Part of this flow is near real-time, while consolidated datasets are loaded into Snowflake for analytical exploration.

One ingestion, multiple consumption paths

Ingestion & storage of tyre usage and wear data into our Data Lake

MQTT carries telemetry from devices into the ingestion infrastructure. Kafka provides a retained event stream: incoming messages land first in a raw topic, kept as received. This layer buffers incoming events, allowing downstream consumers to catch up after load spikes or temporary failures, provided they recover within the retention window. Kafka Streams consumes the raw topic, decodes each message and writes one cleansed event per measurement channel to that channel's dedicated curated topic. These topics are then loaded into Snowflake through Kafka Connect; analysts query the events there to identify fleet usage trends. Kafka Streams filters the data in flight, while Snowflake serves analytics and history-based processing.

To give a sense of scale: our Kafka topics sustain around 70k messages per second, and the time-series tables in Snowflake now hold more than 100 TB of compressed data. At that volume, the question is no longer "can we serve this in real time?" but "which consumer actually needs us to?".
Today, our near real-time use cases are narrow: trip reconstruction for a subset of boxes, and notifying a fleet when a tyre on one of its vehicles starts losing pressure. Everything else can wait, so we don't pay for real-time serving on 70k events/s.

The same ingestion stream also serves both operational needs and analytical workloads, without duplicating pipelines.

  • It preserves the raw events for retention and replay: if a rule changes, a processing defect is found or an occasional R&D need arises, we can reprocess from the source within the retention window.
  • Over longer time horizons, the same data feeds cold analytics.
  • It enables in-flight controls: in our pipeline, Kafka Streams drops implausible GPS points and messages whose timestamp is too old, typically after a box had no network coverage. Late messages are often incomplete, and at several messages per second per box, losing a frame here and there doesn't change the picture of a trip.

Keep in mind that use cases evolve. What initially serves as a simple archival pipeline may later become a critical source for analytics, regulatory audits, or new business insights.

Hot and cold are serving decisions

We use two logical tiers, and a single rule to assign a dataset to one of them:

A dataset is hot when its consumers need an answer within seconds and read it often enough to justify an always-on serving layer. Otherwise, it is cold.

Freshness is deliberately absent from this rule. It does not decide the tier, it decides how the store is fed. Hot does not mean real-time feeding, and cold does not mean batch. Our platform has both cases, as the storage section below shows.

Characteristic Hot Cold/Historical
Consumer expectation Answer within seconds, predictable response time Answer within seconds to hours
Access pattern Frequent Periodic or infrequent
Main priority Timely access Economical retention and usable history
Serving layer Always on Queried on demand

These are logical tiers, not a direct mapping to a storage product.

Cold data can remain immediately readable. Archived data, on the other hand, may need to be restored before it can be read, which only works if consumers can afford the wait.

On our platform, cold covers two cases: the curated history in Snowflake is queried directly by analysts and data scientists, while the raw copy in blob storage, although readable without restoration, is never queried as such, only read back to backfill Snowflake when a new channel is needed.

Keeping the two tiers separate is what prevents the hot tier from becoming the default destination for everything we ingest. A tier expresses a requirement. Snowflake, Kafka or PostgreSQL is an implementation detail, and it can change. The mapping between them can change without the requirement changing.

Retention is a separate decision

Alongside temperature, retention shapes how we handle a dataset: how long it must be kept. It is often conflated with "hot" and "cold", but it answers a different question.

Retention is driven by two very different things.

  • On the technical side, it is a performance and cost lever: keeping only a rolling window of recent events in the hot tier keeps operational queries fast and predictable, while older data moves to cheaper storage or is dropped.
  • On the regulatory side, it can become a hard constraint. In our fleet telemetry, a position or a driving pattern can be linked to an identifiable driver, which brings GDPR into play.

In our case, GDPR is not theoretical.
A telematics box is attached to a vehicle, the vehicle belongs to a customer fleet, and a position or a driving pattern can be traced back to the driver. Raw telemetry therefore cannot simply accumulate in the cold tier because storage is cheap.

Retention is a property of each dataset, and it is enforced at three points in the pipeline:

  • Kafka keeps only 7 to 10 days of history. Enough to absorb a downstream outage or replay after a processing defect, too short to act as a store. The stream never becomes a hidden archive of personal data.
  • Snowflake holds the curated history. An automatic purge removes records once they exceed the dataset's retention period, which combines the legal limit with the terms of each fleet contract.
  • Blob storage holds the raw events, as received, consolidated per telematics box by a dedicated consumer on the raw topic. They are kept on a rolling window of a few months, long enough to backfill a new channel or reprocess a recent period, then purged automatically. The per-box layout also makes it possible to delete one box's history as soon as its contract expires, without touching anything else.

The raw copy is not a duplicate of the curated data. Kafka Streams extracts the measurement channels we currently know how to use; the raw messages contain more, and newer box models add signals beyond accelerometry, GPS or CAN that no curated pipeline consumes yet. When a study asks for one of those channels, the history already exists: we add the extraction to Kafka Streams for the future and backfill the last few months from the raw copy, instead of waiting months for new data to accumulate.

Storage strategy: temperature drives architecture

Storage decisions should reflect access patterns, retrieval expectations, retention periods and operating costs. In our tyre usage and wear platform, this translates into a range of stores, each chosen for a consumer rather than for the data it holds; the table below shows a representative selection.

Use case Serving tier Store Feeding
Box installation & communication check Hot PostgreSQL Kafka Consumer, seconds
Trip reconstruction (selected box types) Hot Kafka Streams (State Store) Stream, seconds
Per-trip / per-box aggregates for studies Hot PostgreSQL Batch every few hours, from Snowflake
R&D tyre wear studies Cold Snowflake (native tables) Micro-batches every few seconds via Kafka Connect
Raw event retention & backfill Cold Blob storage (raw, per box) Continuous, dedicated consumer on the raw topic

The smallest hot use case is the application technicians use to check that a newly installed box communicates correctly. It needs the last few messages of one box, within seconds, and nothing else. A lightweight Kafka consumer keeps only those recent messages in PostgreSQL, so the application never touches the stream itself. Serving the full 70,000 messages per second to a real-time query engine, so that this one use case could exist, would cost a great deal for no additional value.

The aggregates that studies consume through our API are served from PostgreSQL with sub-second response time, yet they are computed by a job that queries Snowflake and runs every few hours. The consumers read frequently and expect fast answers, which makes the serving layer hot; they tolerate data that is hours old, so the pipeline that feeds it can be batch, derived from the historical store. Serving temperature and feeding latency are two independent decisions, and this is the clearest example we have of it.

Which boxes are enrolled in the hot path is decided by need and by capability: a study has to ask for them, and the box model has to emit what trip reconstruction requires. Nothing enters the hot tier by default.

Snowflake is not inherently a cold storage system. In our example, it serves the historical analytics path because that is the role we assign to it, and it is also the source from which part of our hot serving layer, the PostgreSQL aggregates, is derived.

Example of our tyre usage and wear data storage decision tree

Optimise the full cost of access

A lower storage price does not automatically produce a lower total cost.

In our tyre usage and wear data lake, most Snowflake tables holding curated telemetry are time series: one row per telematics box and per emission timestamp. We cluster these tables on two keys, the telematics box identifier and the emission date of the message. This choice is driven directly by how our data scientists query the data. A typical study targets a single box, or a fleet of boxes, over a bounded period: a few hours or days. With these two clustering keys, a query filtering on box_id and a date range only scans the micro-partitions that actually contain the relevant boxes and days, instead of the whole history.

This has two effects:

  • Queries return faster for the consumer.
  • The warehouse consumes less compute, because far less data is read per query.

What we optimised is the cost of each access, which is where most of the spend on historical data actually sits.

On a 3.8 To table, a study on one fleet over a few days touches a tiny fraction of the micro-partitions. Without clustering on box_id and emission date, the same query would have to scan a large share of the table to find a handful of boxes and days, making it both slow and expensive to run.

Snowflake query detail: 3.8To table, 2.5% partitions scanned

Clustering is not free: automatic clustering is billed in credits and rewrites micro-partitions. In our case, that cost scales with the ingestion rate rather than with the size of the table. Events arrive ordered by emission time, so only freshly ingested partitions need to be reorganised by box_id; once clustered, historical partitions are almost never rewritten.

Schema evolution and cold data governance

Cold data changes shape even when nobody queries it.

Since this platform went live, new telematics box models have added channels, some units have been harmonised and fields renamed. Every change had to land without breaking studies that started months earlier on the same tables. That is why the way consumers access the data has to tolerate evolution over years.

Treat consumer interfaces as contracts

Our consumers never query the physical tables. They query curated views that act as a contract: stable names, stable types, stable units. Behind that interface, the tables are free to evolve.

Most changes can be absorbed by the view itself. When a field is renamed, the view keeps exposing the old name. When a unit is harmonised, the view converts. When a new box model adds channels, the view exposes them only once we decide to. The table changes; the queries written against the view keep working.

Some changes cannot be hidden this way. If distance, once computed from GPS points, is now read from the vehicle's CAN odometer, the column keeps the same name and unit, but the values are no longer comparable. No view can reconcile them, and an analysis mixing both would be wrong.

For those, we publish a new version of the view alongside the previous one. Both coexist while consumers migrate, and the older version is retired once the last study relying on it has closed.

From cold data to actionable insight

Historical data becomes useful when consumers can understand it, trust it, and query it at an acceptable cost.

When an R&D data scientist analyses tyre wear evolution over several months, the retained history in Snowflake is only the starting point: the analysis also needs consistent units, reliable timestamps, relevant vehicle context such as truck type or number of tyres per axle, and a clear treatment of incomplete observations.

Curated datasets and views provide an entry point. Reusable aggregations can reduce repeated processing, while detailed records remain available for investigations.

We also turn recurring needs into shared assets. Some enrichments are valuable to many consumers, not just one analysis. In our case, ambient temperature, weather conditions, and road type are relevant to almost anyone studying a trip. Rather than letting each team reconstruct that context on its own, we consolidate it once, at the trip level, and expose it alongside the data so every consumer benefits from the same enriched view.

Design for consumption

None of this is really about technologies like Kafka, Snowflake or PostgreSQL. Those choices will change. What we try to keep is the habit of asking one question before picking a store: who reads this data, how often, and how long can they wait?

Arrival speed and age are properties of the data. Temperature is a decision about its consumers. When we mix the two up, we either overpay for real-time systems nobody queries, or we bury useful history in places where nobody can reach it.

Starting from consumption doesn't just save money. It also gives the data we keep a better chance of still being useful for new use cases in a few years.