Aggregators¶
Aggregate raw trades into bars using different policies. All functions accept numpy arrays and release the GIL.
Common Parameters¶
All aggregation functions share the same input signature:
| Parameter | Type | Description |
|---|---|---|
timestamps |
int64[] |
Trade timestamps in nanoseconds |
prices |
float64[] |
Trade prices |
quantities |
float64[] |
Trade quantities |
is_buy |
uint8[] |
Buy/sell flag (non-zero = buy, 0 = sell) |
is_buy is the inverse of the side field on PyTrade from DataReader, where 0 = buy and 1 = sell. Convert with (trades['side'] == 0).astype(np.uint8) — passing side straight through inverts every bar's buy volume.
Returns: numpy.ndarray with PyExtBar structured dtype. Only fully
closed bars are returned -- a trailing bar still open when the trade array
ends is dropped, so the bar count depends only on the input, not on where
the array happens to stop. Node, QuickJS, and Codon return the same count on
identical input. To see a bar the moment it closes during a live run, use
the bus-fed BarAggregator / Strategy.lastClosedBar path instead of this
batch function.
PyExtBar Dtype¶
| Field | Type | Description |
|---|---|---|
start_time_ns |
int64 |
Bar open timestamp (unix ns) |
end_time_ns |
int64 |
Bar close timestamp (unix ns) |
open_raw |
int64 |
Open price * 10^8 |
high_raw |
int64 |
High price * 10^8 |
low_raw |
int64 |
Low price * 10^8 |
close_raw |
int64 |
Close price * 10^8 |
volume_raw |
int64 |
Total volume * 10^8 |
buy_volume_raw |
int64 |
Buy volume * 10^8 |
trade_count |
int64 |
Number of trades in bar |
close_reason |
uint8 |
Why the bar closed: 0 threshold, 2 forced, 3 warmup. Every bar this batch path returns closed on its own threshold, so 0 here. |
Functions¶
aggregate_time_bars(..., interval_seconds)¶
Aggregate trades into fixed time intervals.
| Parameter | Type | Description |
|---|---|---|
interval_seconds |
float |
Bar duration in seconds |
aggregate_tick_bars(..., tick_count)¶
Aggregate trades into bars with a fixed number of trades.
| Parameter | Type | Description |
|---|---|---|
tick_count |
int |
Trades per bar |
aggregate_volume_bars(..., volume_threshold)¶
Aggregate trades into bars when cumulative volume exceeds a threshold.
| Parameter | Type | Description |
|---|---|---|
volume_threshold |
float |
Volume per bar |
aggregate_range_bars(..., range_size)¶
Aggregate trades into bars when price range exceeds a threshold.
| Parameter | Type | Description |
|---|---|---|
range_size |
float |
Maximum price range per bar |
aggregate_renko_bars(..., brick_size)¶
Aggregate trades into Renko bars with a fixed brick size. Unlike the other aggregators here, Renko can return more bars than there were input trades: a trade that gaps past more than one brick width closes the brick that was forming and synthesizes the bricks in between, walking the price from there to the trade in clean brick-size steps (see bar types).
| Parameter | Type | Description |
|---|---|---|
brick_size |
float |
Renko brick size |
aggregate_heikin_ashi_bars(..., interval_seconds)¶
Aggregate trades into Heikin-Ashi bars (smoothed candlesticks).
bars = lrvx.aggregate_heikin_ashi_bars(timestamps, prices, quantities, is_buy,
interval_seconds=60.0)
| Parameter | Type | Description |
|---|---|---|
interval_seconds |
float |
Bar duration in seconds |
Example¶
import numpy as np
import lrvx
# Load trade data
reader = lrvx.DataReader("./data")
trades = reader.read_trades()
ts = trades['exchange_ts_ns']
px = trades['price_raw'] / 1e8
qty = trades['qty_raw'] / 1e8
is_buy = (trades['side'] == 0).astype(np.uint8) # PyTrade.side: 0 = buy
# Create 1-minute time bars
bars = lrvx.aggregate_time_bars(ts, px, qty, is_buy, interval_seconds=60.0)
print(f"Generated {len(bars)} time bars")
# Create 500-tick bars
tick_bars = lrvx.aggregate_tick_bars(ts, px, qty, is_buy, tick_count=500)
print(f"Generated {len(tick_bars)} tick bars")
# Access bar data
opens = bars['open_raw'] / 1e8
highs = bars['high_raw'] / 1e8
lows = bars['low_raw'] / 1e8
closes = bars['close_raw'] / 1e8