Data systems · E25 · Implementation

When does an online bar actually become observable?

A timestamp-bucket aggregator closes a bar when the next bucket arrives—not merely when a clock crosses its nominal boundary.

Python state machineTimestamp arithmeticdeepcopy
The aggregate’s timestamp and the event that makes it available are different points on the timeline.
Figure 1. When is a bar complete?. The aggregate’s timestamp and the event that makes it available are different points on the timeline. Illustrative event-time diagram. Original vector illustration.

Follow the information

From input to outcome

The aggregate is updated while its bucket is current. A later-bucket arrival triggers emission of the copied old aggregate; its chart timestamp is earlier than the time the consumer receives it.

The aggregate is updated while its bucket is current. A later-bucket arrival triggers emission of the copied old aggregate; its chart timestamp is earlier than the time the consumer receives it.
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: The aggregate’s timestamp and the event that makes it available are different points on the timeline. The module map and layer-level figures below expand the operations in this route.

When does an online bar actually become observable?: architectureIncoming source bar: Timestamp + OHLCV → Bucket assignment: floor(t / resolution) → Current aggregate: Mutable OHLCV state → Later bucket arrives: Emit copied old bar → Downstream consumer: Receives completed event. A high-level module map; comparison branches and training details are explained in the article.DATA SYSTEMS / E25 / MODULE MAP01 INPUTIncoming source barTimestamp + OHLCV02 MODULEBucket assignmentfloor(t / resolution)03 MODULECurrent aggregateMutable OHLCV state04 MODULELater bucket arrivesEmit copied old bar05 OUTPUTDownstream consumerReceives completed event
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.
Incoming source bar — Timestamp + OHLCV

The architecture in context

The system we are building

The aggregator maps source timestamps into fixed-width buckets. It holds one mutable aggregate, updates its extrema and volume, and returns a completed copy when a later bucket appears. This is a compact online state machine: the consumer receives an event only when the input stream provides evidence that the current bucket has ended.

Who does what in the stack

Python state machine
Owns bucket transitions and aggregation.
Timestamp arithmetic
Defines bucket identity independently of callback timing.
deepcopy
Detaches emitted bars from mutable current state.

The custom implementation controls the precise emission boundary and preserves a copied completed bar before initializing the next one. It is not a neural layer, but its semantics determine which information every downstream model may use at a given instant.

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

A bar label is not its completion time

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 ↗

The state machine emits the previous aggregate when it observes a newer bucket. If no later event arrives, the last bucket remains un-emitted. A gap does not automatically create intermediate bars. A timestamp printed at the bucket start is therefore not the time at which all its values became available.

The mathematical contract

b(t)=⌊t/Δ⌋Δb(t)=\lfloor t/\Delta\rfloor\Delta

A watermark-based implementation could trade delay against tolerance for late arrivals, but that is a different completion policy. The source does not justify treating a late older event as part of the current bucket. Explicit shutdown behavior is needed for finite recordings.

Implementation and resource card

Capacity / budget
Constant-size current aggregate; no learned weights. Completion latency depends on later arrivals, not only interval length.
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

Test boundary-minus-one, exact boundary, a multi-bucket jump and end-of-stream. Keep chart label, last contributing event and emission timestamp in the test ledger.

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’s strict greater-than comparison is the transition. If no new input arrives, no completion event is emitted. If the next input jumps several buckets, the implementation does not automatically create observations for the missing intervals. Deep-copying the completed aggregate prevents subsequent mutation from rewriting a previously emitted bar.

Python · file · lines 39–61
        # Bar 10:00:45 -> falls into 10:00:45 bucket.
        this_bar_bucket = (bar_ts // self.target_res) * self.target_res
        
        # 3. Initialization check (First run)
        if self.current_bar is None:
            self._init_new_bar(bar, this_bar_bucket)
            return None
            
        # 4. Check if we moved to a NEW bucket
        # If the incoming bar belongs to a later bucket than the one we are building...
        if this_bar_bucket > self.current_bucket_start:
            
            # A. Close the old bar
            completed_bar = deepcopy(self.current_bar)
            
            # B. Start the new bar (using the incoming data)
            self._init_new_bar(bar, this_bar_bucket)
            
            # C. Return completed bar
            # This triggers the Strategy Logic to run
            return completed_bar
            
        # 5. Same bucket -> Merge Data

Verbatim archive excerpt from aggregator.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

A late bar with an earlier bucket needs an explicit policy; treating every non-new bucket as the current one can corrupt state. The final bucket also needs a declared shutdown or watermark rule rather than being assumed complete.

Keep building

Other posts of interest