Market-data systems are often designed as if every collector has a stable link, every venue answers in order, and every packet arrives once. That assumption breaks quickly when a team consumes several exchange feeds over mobile connections, congested international routes, or infrastructure that occasionally changes IP addresses. The difficult problem is not drawing a chart. It is being able to explain, hours later, exactly why one candle has a particular open, high, low, close, and volume.

This is relevant anywhere connectivity varies, including teams operating between Pakistan, Armenia, and global trading venues. A brief outage should not silently rewrite history, and a reconnect burst should not be mistaken for a sudden burst of market activity.

Start with an evidence envelope

I treat each received trade as evidence, not as a finished candle. Before doing aggregation, the collector wraps the source message in a small immutable envelope:

source_venue
source_channel
instrument
trade_id
exchange_sequence
event_time
received_at
payload_hash
collector_id
schema_version

event_time says when the venue believes the trade happened. received_at says when our collector observed it. They answer different questions and should never overwrite one another. The difference between them is useful telemetry: a growing gap can reveal a degraded route even when the WebSocket still appears connected.

The payload_hash is calculated from canonicalized source fields before normalization. It provides a compact way to prove that a later calculation used the same evidence. It is not a substitute for retaining the raw record, and it does not prove that the venue itself was correct. It only protects the boundary between what arrived and what our pipeline later derived.

Make ingestion idempotent

Reconnects commonly replay data. A consumer that simply appends everything will inflate volume and sometimes alter high or low values. A practical deduplication key is venue-specific:

(source_venue, instrument, trade_id)

If a venue does not supply stable trade IDs, the fallback key can include sequence number, event timestamp, side, price, and quantity. That fallback must be labelled as heuristic because two legitimate trades can share those fields.

I keep a durable ingestion offset separately from the business record. A message is acknowledged only after both the immutable event and its offset are committed. On restart, processing resumes from the last committed offset; replay is safe because the write is idempotent.

Detect gaps instead of hiding them

Sequence numbers are valuable even when they are imperfect. If the last accepted message has sequence 81042 and the next one is 81047, the system records a gap for 81043–81046. It should not invent four trades or interpolate prices. The gap becomes an explicit object with a status such as:

  • detected
  • backfill_requested
  • partially_recovered
  • recovered
  • unrecoverable

That state can travel with downstream data. A chart may still render, but an analyst can distinguish a complete interval from one built with missing evidence.

For feeds without sequences, I use heartbeats and expected activity as weaker signals. “No messages” is ambiguous: the market may be quiet, the connection may be stale, or the subscription may have been dropped. Recording the last message time, last heartbeat time, and last successful REST snapshot makes that ambiguity visible.

Build candles in event time

Candles should be derived from event_time, with a deterministic tie-breaker. A stable ordering rule is:

(event_time, exchange_sequence, trade_id)

If a venue lacks one field, its adapter declares the fallback rather than letting database insertion order decide. This matters for the open and close: two trades can have the same timestamp but different sequences.

Every interval moves through explicit states:

  1. provisional while messages can still arrive normally;
  2. final after a configured watermark;
  3. corrected if later evidence changes the result.

The watermark is not a promise that nothing late will arrive. It is a policy based on observed delay. For example, if nearly all events arrive within 20 seconds, a one-minute candle might become final after a 30-second allowance. The threshold should be measured per venue and route, not copied from another deployment.

When a late trade changes a final candle, I create a new candle revision. The earlier value remains addressable. A correction record includes the old hash, new hash, triggering event ID, calculation version, and correction time. Consumers can then choose whether to follow the latest revision or reproduce the view that existed at an earlier point.

Separate the log from the view

One table should not be asked to serve as raw evidence, normalized data, current state, and audit history. A clearer design has four layers:

raw source records
        ↓
normalized immutable trades
        ↓
versioned derived candles
        ↓
API/read models

The raw layer helps diagnose adapter bugs. The normalized layer gives every venue a shared schema. The derived layer can be rebuilt when calculation logic changes. The read model is optimized for clients and can be discarded and regenerated.

This separation also makes schema migration safer. If a decimal-precision rule changes, a team can run the new calculator beside the old one, compare hashes, and promote it only after discrepancies are understood.

Design for constrained links

Bandwidth adaptation should change transport, not meaning. I use several techniques:

  • resume from a durable sequence or cursor rather than download an entire day again;
  • compress batches, but hash canonical uncompressed records;
  • request only missing ranges during backfill;
  • keep reference metadata such as symbol rules in a versioned local cache;
  • send compact integrity manifests before transferring large blocks;
  • apply bounded queues and visible backpressure instead of silently dropping old messages.

A manifest can contain the instrument, time range, record count, first and last sequence, and a Merkle root or batch hash. The receiver compares the manifest with local state and requests only absent blocks. If the connection fails again, completed blocks do not need to be repeated.

The collector should also have a storage-pressure policy. “Keep everything until the disk is full” is not a policy. I reserve space for offsets and gap metadata, rotate raw payloads according to a documented retention period, and raise an alert before pruning anything. If evidence must be dropped, the loss is itself recorded.

Expose quality in the API

Returning a number without its state encourages clients to treat provisional and corrected data as identical. A candle response can expose a compact quality object:

{
  "status": "corrected",
  "revision": 3,
  "source_count": 2,
  "missing_ranges": [],
  "calculation_version": "ohlcv-2.1",
  "finalized_at": "2026-09-23T05:31:30Z"
}

This is especially useful for cached mobile clients. A client can store the candle revision and ask only for intervals whose revisions changed. The server does not need to resend the whole series, and the client does not confuse a corrected close with a brand-new interval.

Test invariants, not only examples

Example-based tests catch familiar cases. Property-based tests are better at producing awkward orders and duplicates. Useful invariants include:

  • ingesting the same batch twice does not change a candle;
  • shuffling arrival order does not change the final result when event ordering keys are unchanged;
  • volume equals the sum of unique accepted quantities;
  • high is never below open, close, or low;
  • a correction increases the revision and preserves the previous revision;
  • replaying from a committed offset produces the same hashes;
  • an unresolved gap can never be labelled complete.

Test fixtures should include clock skew, reconnect bursts, equal timestamps, duplicate IDs, sequence resets, malformed decimals, and schema changes. Synthetic records are preferable here: no customer data or live credentials are required.

Operational questions worth answering

Before calling a pipeline reliable, I want precise answers to these questions:

  • Can we reproduce any published candle from retained evidence and a named code version?
  • Can we tell whether an empty interval was genuinely quiet or merely disconnected?
  • Does a late trade create a visible correction rather than a silent overwrite?
  • Can a new collector resume without duplicating volume?
  • Can a low-bandwidth consumer verify a batch before downloading it fully?
  • Are deletion and retention decisions themselves auditable?

At ARMCP, these questions shape how I think about market-data and blockchain tooling. The broader lesson is that reliability is not the absence of outages. It is the ability to recover without losing provenance, to show uncertainty honestly, and to reproduce every published value.

Disclosure and limits: All identifiers, timestamps, sequences, and examples above are synthetic. This is a technical architecture discussion, not financial, investment, trading, or legal advice. AI assistance was used for drafting and structure; I, Mushegh Manukyan, Founder & CEO of ARMCP, reviewed the technical claims, limits, and final text and take responsibility for it.