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 in nautilus_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 CustomData wrappers (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 up fsspec filesystem (lines 30-45) and an ArrowSerializer (line 51).
  • write uses pq.write_table with filenames from class_to_filename (line 64 in funcs.py).
  • query builds a pyarrow.dataset scan with optional filter expressions via combine_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.RecordBatchStreamWriter internally (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_arrow in nautilus_trader/serialization/arrow/serializer.py.
  • Serialize single objects via ArrowSerializer.serialize or batches via serialize_batch to produce PyArrow tables.
  • Persist batch data using ParquetDataCatalog.write or stream continuously with StreamingFeatherWriter and configurable rotation modes.
  • Query persisted data with ParquetDataCatalog.query supporting time-range filters, or deserialize directly using ArrowSerializer.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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →