How to Serialize Trading Data with Arrow/Parquet in Nautilus Trader: A Complete Guide
Nautilus Trader serializes trading data to Apache Arrow record batches and persists them as Parquet files using the ArrowSerializer class and ParquetDataCatalog, enabling high-performance storage of market data, order book updates, and custom objects.
Nautilus Trader leverages Apache Arrow and Parquet to provide efficient, columnar storage for trading data. This guide explains how to serialize trading data with Arrow/Parquet using the nautilus_trader.serialization.arrow and nautilus_trader.persistence.catalog packages from the nautechsystems/nautilus_trader repository.
Registration: Declaring Arrow-Compatible Types
Before serialization, data classes must be registered with an Arrow schema. The register_arrow function in nautilus_trader/serialization/arrow/serializer.py handles this mapping.
When the package initializes, built-in types like QuoteTick, TradeTick, and Bar are automatically registered:
# From serializer.py (lines 70-84)
for _data_cls in NAUTILUS_ARROW_SCHEMA:
if _data_cls in RUST_SERIALIZERS:
register_arrow(
data_cls=_data_cls,
schema=NAUTILUS_ARROW_SCHEMA[_data_cls],
)
else:
register_arrow(
data_cls=_data_cls,
schema=NAUTILUS_ARROW_SCHEMA[_data_cls],
encoder=make_dict_serializer(NAUTILUS_ARROW_SCHEMA[_data_cls]),
decoder=make_dict_deserializer(_data_cls),
)
Key components:
NAUTILUS_ARROW_SCHEMA: A dictionary mapping classes to PyArrow schemas defined innautilus_trader/serialization/arrow/schema.py.- Python types: Use dictionary-based encoders (
make_dict_serializer) and decoders (make_dict_deserializer). - Rust types: Delegated to the Rust bridge via
RUST_SERIALIZERS(lines 52-60).
Custom instruments are auto-registered by iterating over Instrument.__subclasses__() (lines 86-92). To add support for a new instrument, implement its schema in nautilus_trader/serialization/arrow/implementations/instruments.py.
Conversion: Serializing Objects to Arrow
Single Object Serialization
Use ArrowSerializer.serialize to convert a single object to a PyArrow RecordBatch:
from nautilus_trader.serialization.arrow.serializer import ArrowSerializer
record_batch = ArrowSerializer.serialize(my_trade_tick) # → pyarrow.RecordBatch
The method:
- Unwraps
CustomDatawrappers (lines 13-15). - Looks up encoders in
_ARROW_ENCODERS, falling back to Rust serializers if needed (lines 20-24).
Batch Serialization
For multiple objects, use serialize_batch to create an Arrow Table:
from nautilus_trader.serialization.arrow.serializer import ArrowSerializer
table = ArrowSerializer.serialize_batch(
data=[tick1, tick2, tick3],
data_cls=QuoteTick, # Common type of the batch
) # → pyarrow.Table
For Rust-defined types, this shortcuts to rust_defined_to_record_batch (lines 59-61). For Python types, it assembles RecordBatch objects into a Table (lines 62-64).
Deserialization
Convert Arrow data back to Nautilus objects:
from nautilus_trader.serialization.arrow.serializer import ArrowSerializer
ticks = ArrowSerializer.deserialize(QuoteTick, table) # → List[QuoteTick]
This uses _ARROW_DECODERS. For Rust types, _deserialize_rust creates a wrangler (lines 303-322) to reconstruct native objects from the Arrow table.
Persistence: Writing and Reading Parquet
Using the ParquetDataCatalog
The ParquetDataCatalog in nautilus_trader/persistence/catalog/parquet.py provides high-level read/write operations:
from nautilus_trader.persistence.catalog.parquet import ParquetDataCatalog
from nautilus_trader.model.data import TradeTick
# Initialize catalog
catalog = ParquetDataCatalog(path="./data")
# Write batch
table = ArrowSerializer.serialize_batch([trade1, trade2], TradeTick)
catalog.write(table, TradeTick) # Writes to ./data/trade_tick.parquet
# Query time range
df = catalog.query(
data_type=TradeTick,
start_time="2024-01-01T00:00:00Z",
end_time="2024-01-01T01:00:00Z",
)
Implementation details:
__init__sets upfsspecfilesystem (lines 30-45) and anArrowSerializer(line 51).writeusespq.write_tablewith filenames fromclass_to_filename(line 64 infuncs.py).querybuilds apyarrow.datasetscan with optional filter expressions viacombine_filters.
Streaming with FeatherWriter
For continuous, low-latency persistence, use StreamingFeatherWriter from nautilus_trader/persistence/writer.py:
from nautilus_trader.persistence.writer import StreamingFeatherWriter, RotationMode
from nautilus_trader.cache.cache import Cache
from nautilus_trader.common.component import Clock
writer = StreamingFeatherWriter(
path="./feather_stream",
cache=Cache(),
clock=Clock(),
rotation_mode=RotationMode.SIZE,
max_file_size=500_000_000, # 500 MiB per file
)
# Write any Nautilus object
writer.write(my_bar)
writer.write(my_quote_tick)
writer.flush() # Force buffer write
Key behaviors:
- Detects object class and resolves filename via
class_to_filename(lines 77-90). - Creates per-instrument writers for instrument-specific data (lines 91-96).
- Uses
pyarrow.RecordBatchStreamWriterinternally (line 27).
End-to-End Example: Persisting and Reloading Trade Data
# -------------------------------------------------
# 1. Setup
# -------------------------------------------------
from nautilus_trader.persistence.catalog.parquet import ParquetDataCatalog
from nautilus_trader.serialization.arrow.serializer import ArrowSerializer
from nautilus_trader.model.data import TradeTick
import pandas as pd
catalog = ParquetDataCatalog(path="./data")
# -------------------------------------------------
# 2. Create sample data
# -------------------------------------------------
ticks = [
TradeTick(
instrument_id="BTC-USD-PERP",
price="30000.0",
size="0.005",
side="BUY",
ts_event=pd.Timestamp("2024-02-01T12:00:00Z").value,
ts_init=pd.Timestamp("2024-02-01T12:00:00Z").value,
)
for _ in range(10)
]
# -------------------------------------------------
# 3. Serialize and write
# -------------------------------------------------
table = ArrowSerializer.serialize_batch(ticks, TradeTick)
catalog.write(table, TradeTick) # ./data/trade_tick.parquet
# -------------------------------------------------
# 4. Query back
# -------------------------------------------------
df = catalog.query(
data_type=TradeTick,
start_time="2024-02-01T11:59:00Z",
end_time="2024-02-01T12:01:00Z",
)
print(df.head())
The resulting DataFrame contains the original trade fields. Use ArrowSerializer.deserialize(TradeTick, table) to reconstruct fully-typed TradeTick objects.
Key Source Files
| File | Role | Location |
|---|---|---|
serializer.py |
Central registry, (de)serialization logic, Rust bridge | nautilus_trader/serialization/arrow/serializer.py |
schema.py |
Arrow schema definitions for all built-in types | nautilus_trader/serialization/arrow/schema.py |
parquet.py |
Queryable Parquet catalog, file layout, fsspec integration | nautilus_trader/persistence/catalog/parquet.py |
writer.py |
Rotating Feather writer for streaming use-cases | nautilus_trader/persistence/writer.py |
funcs.py |
Helper utilities (class_to_filename, combine_filters) |
nautilus_trader/persistence/funcs.py |
Summary
- Register every data class with an Arrow schema using
register_arrowinnautilus_trader/serialization/arrow/serializer.py. - Serialize single objects via
ArrowSerializer.serializeor batches viaserialize_batchto produce PyArrow tables. - Persist batch data using
ParquetDataCatalog.writeor stream continuously withStreamingFeatherWriterand configurable rotation modes. - Query persisted data with
ParquetDataCatalog.querysupporting time-range filters, or deserialize directly usingArrowSerializer.deserialize.
Frequently Asked Questions
How do I add Arrow serialization support for a custom data type?
Implement a PyArrow schema for your class in nautilus_trader/serialization/arrow/implementations/ and call register_arrow with your class, schema, and optional encoder/decoder functions. The registration system automatically picks up instrument subclasses, but custom data types require manual registration before serialization.
What is the difference between ParquetDataCatalog and StreamingFeatherWriter?
ParquetDataCatalog in nautilus_trader/persistence/catalog/parquet.py is designed for batch operations, writing complete Arrow tables to Parquet files and supporting time-range queries via pyarrow.dataset. StreamingFeatherWriter in nautilus_trader/persistence/writer.py provides low-latency, continuous persistence with automatic file rotation based on size or time, using the Feather format for efficient streaming writes.
How does Nautilus Trader handle Rust-defined types during serialization?
Rust-defined types (prefixed with nautilus_pyo3) bypass the Python encoder/decoder dictionary and use the RUST_SERIALIZERS registry. The ArrowSerializer detects these types and delegates to rust_defined_to_record_batch for serialization and _deserialize_rust for deserialization, leveraging the Rust bridge for high-performance conversion.
Can I query Parquet files by specific instruments or time ranges?
Yes. The ParquetDataCatalog.query method accepts data_type, start_time, end_time, and optional filter_expr parameters. It constructs a pyarrow.dataset scan and applies filters via combine_filters from nautilus_trader/persistence/funcs.py, allowing efficient retrieval of specific instruments or time windows without loading entire files into memory.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →