Data systems · E24 · Implementation

A data callback is not a durable dataset

A scheduled market-data recorder separates callbacks, buffers and file writes. Each boundary has a different failure mode.

Feed API callbacksPython buffers / CSVScheduler
Data must travel beyond the callback into a durable, auditable record before it can support a dataset claim.
Figure 1. A callback is not a dataset. Data must travel beyond the callback into a durable, auditable record before it can support a dataset claim. Ingestion schematic. Original vector illustration.

Follow the information

From input to outcome

Callbacks append observations to buffers; a timer checks when to persist them. A successful callback or write is not a completeness guarantee for the eventual dataset.

Callbacks append observations to buffers; a timer checks when to persist them. A successful callback or write is not a completeness guarantee for the eventual dataset.
Figure 2. Information flow. Solid arrows carry observations, tensors or artifacts; other routes are explicitly labelled. Signal shapes, matrices and network icons are schematic, not measured samples or literal neuron counts. Open full-size SVG ↗ On narrow screens, scroll the diagram horizontally.

Read this alongside Figure 1: Data must travel beyond the callback into a durable, auditable record before it can support a dataset claim. The module map and layer-level figures below expand the operations in this route.

A data callback is not a durable dataset: architectureFeed callbacks: Bars / quotes / trades → Per-symbol buffers: Separate record types → Periodic save check: Elapsed host time → CSV writers: Append + flush → Offline dataset: Needs coverage audit. A high-level module map; comparison branches and training details are explained in the article.DATA SYSTEMS / E24 / MODULE MAP01 INPUTFeed callbacksBars / quotes / trades02 MODULEPer-symbol buffersSeparate record types03 MODULEPeriodic save checkElapsed host time04 MODULECSV writersAppend + flush05 OUTPUTOffline datasetNeeds coverage audit
Source-grounded module map. Boxes summarize operations, not individual neurons; comparison arms and training paths are detailed below. On a small screen, scroll the diagram horizontally.
Feed callbacks — Bars / quotes / trades

The architecture in context

The system we are building

A recorder should keep its hot callback small and move storage work into a controlled path. This script collects several feed record types into per-symbol buffers and periodically writes them to corresponding CSV streams. A scheduler determines when the recorder is active. The same architecture applies to laboratory sensors and online ML telemetry.

Who does what in the stack

Feed API callbacks
Deliver asynchronous source observations.
Python buffers / CSV
Separate ingestion from persistence.
Scheduler
Controls recording windows without defining event time.

The project-specific work maps callback fields into structured rows, maintains per-symbol state and manages scheduled recording windows. The upstream broker API supplies events; it does not certify that the resulting files are complete, correctly ordered or suitable for a later learning target.

Framework responsibility map. Each row maps a library or custom component to its job; rows are not a sequential inference graph.
Framework responsibility map. Each row maps a library or custom component to its job; rows are not a sequential inference graph. Open full-size SVG ↗

Open up the implementation

Design the ingestion clock before the model clock

A concrete operation-level view of this implementation; no unobserved neural architecture is implied.
A concrete operation-level view of this implementation; no unobserved neural architecture is implied. Open full-size SVG ↗

Event time identifies what happened in the market or sensor. Arrival time records when the system could know it. Flush time records persistence. Conflating those clocks makes both latency claims and causal feature reconstruction unreliable. A callback should minimize work before handing data to a controlled processing stage.

The mathematical contract

tevent≠tarrival≠tflusht_{\rm event}\neq t_{\rm arrival}\neq t_{\rm flush}

Larger buffers reduce write frequency but increase loss exposure on interruption. Smaller flushes can dominate I/O. Provider reconnects and duplicate events require explicit semantics; a continuously running scheduler is not proof of continuous observation.

Implementation and resource card

Capacity / budget
Ingestion component, not neural training. Capacity is buffer growth, flush interval and loss accounting.
Execution evidence
This revision inspects and explains the archived implementation. It does not rerun the original workload. No unrecorded convergence time, throughput or accelerator result is supplied.
Current reproduction context
Current workstation, supplied by the author: Apple M4, 128 GB unified RAM, 40 GPU cores and 16 CPU cores. This is context for prospective reproduction, not attribution of every archived run. Python and framework versions are not fully locked for these historical sources; declarations, when available, are identified separately.

From explanation to a reproducible check

Use a mock callback stream with a duplicate, a reconnect gap and a delayed event. Verify persisted order and coverage metadata. This article does not connect to a broker or initiate live capture.

Preserve input identities, configuration and failure records with the result. A successful numerical check only establishes the operation it exercises: it does not certify an entire dataset, model or deployed system. Reproduce the interface on a small deterministic input before optimizing throughput or increasing workload size.

A closer look at the implementation

The code that carries the idea

The excerpt writes every buffered row, flushes the file handle and clears the buffer. That ordering avoids clearing data before the write call, but a Python flush is not a complete crash-durability protocol. Failure during a partial write can also complicate replay or recovery.

Python · file · lines 233–254
    def check_and_save(self):
        current_time = time.time()
        if current_time - self.last_save_time >= self.save_interval:
            self.save_buffers_to_csv()
            self.last_save_time = current_time

    def save_buffers_to_csv(self):
        for ticker in self.tickers:
            # Save Real-time bars
            if self.realtime_bars_buffer[ticker]:
                for row in self.realtime_bars_buffer[ticker]:
                    self.csv_writers[ticker]['bars'].writerow(row)
                self.csv_files[ticker]['bars'].flush()
                self.realtime_bars_buffer[ticker].clear()
            
            # Save NBBO
            if self.nbbo_buffer[ticker]:
                for row in self.nbbo_buffer[ticker]:
                    self.csv_writers[ticker]['nbbo'].writerow(row)
                self.csv_files[ticker]['nbbo'].flush()
                self.nbbo_buffer[ticker].clear()

Verbatim archive excerpt from scheduled_futures_l1_stream.py. Context-dependent historical code, not a standalone runnable program. Comments retain their original wording; the article distinguishes implemented behavior from stale or overbroad comments.

The boundary that matters

Host receipt time, exchange event time and bar-completion time are different clocks. A file’s presence is not proof of uninterrupted coverage. No connections, subscriptions or recordings were started during this editorial work.

Keep building

Other posts of interest