How to Use Streaming Writers for Real-Time Data Persistence in Nautilus Trader

Use StreamingFeatherWriter to persist live market data to Feather files with automatic rotation, per-instrument file handling, and configurable flush intervals.

Nautilus Trader provides high-performance streaming writers for real-time data persistence, enabling you to capture live market data to disk with minimal latency. The StreamingFeatherWriter class in nautilus_trader/persistence/writer.py serves as the primary interface for writing ticks, bars, and order-book deltas to Apache Feather format on the fly.

Core Architecture of StreamingFeatherWriter

Writer Initialization and Creation

The StreamingFeatherWriter class, defined at line 55 of nautilus_trader/persistence/writer.py, orchestrates the entire persistence pipeline. The _create_writers method (line 56) pre-creates writers for all built-in Arrow-serializable types, while custom data writers are instantiated lazily as new data types arrive.

Per-Instrument File Handling

When a data object carries an instrument_id attribute—such as QuoteTick, TradeTick, or custom data—the writer stores the file handle under a composite (table, instrument_id) key. This logic appears around lines 90-110, with key construction specifically at lines 112-124. This automatic organization creates separate Feather files for each instrument, preventing data fragmentation and simplifying downstream analysis.

Serialization Pipeline

Data objects are converted to Arrow record batches via ArrowSerializer.serialize_batch (line 41), which transforms incoming objects into Arrow RecordBatch instances before writing to disk. The _extract_obj_metadata method (line 59) attaches instrument-specific metadata—including price and size precision—to the Arrow schema, ensuring that deserialized data retains its original semantic meaning.

Configuring File Rotation and Flushing

Size-Based Rotation

File rotation logic resides in _check_file_rotation (line 72). For size-based rotation, the writer tracks accumulated bytes in _file_sizes and compares against max_file_size (lines 89-92). When a file exceeds the threshold, the writer closes the current stream and opens a new file with an updated timestamp.

Time-Based Rotation

For interval rotation, the writer stores the next rotation timestamp in _next_rotation_times and compares against clock.utc_now() (lines 92-100). Set rotation_mode=RotationMode.INTERVAL and provide a pandas.Timedelta via rotation_interval to enable this behavior.

Scheduled Date Rotation

Align rotation to specific wall-clock times using rotation_time and rotation_timezone (lines 102-112). This mode triggers rotation at a specific hour and minute in the specified timezone, useful for creating daily files that align with market sessions.

Flush Intervals

The check_flush method (line 101) runs on every write operation. If the elapsed time since the last flush exceeds flush_interval_ms, it triggers flush to ensure data durability. The default interval is 1000ms (1 second), but you can reduce this to 100ms for critical data at the cost of increased I/O overhead.

Implementation Example

from nautilus_trader.persistence.writer import StreamingFeatherWriter, RotationMode
from nautilus_trader.cache.cache import Cache
from nautilus_trader.common.clock import Clock
import pandas as pd

# 1️⃣ Initialise core services

cache = Cache()                     # populated elsewhere (e.g., from a data catalog)

clock = Clock()                     # uses system UTC time

# 2️⃣ Create the streaming writer

writer = StreamingFeatherWriter(
    path="/tmp/nautilus_stream",   # directory that will hold the feather files

    cache=cache,
    clock=clock,
    fs_protocol="file",            # local filesystem; use "s3", "gcs", etc. for cloud storage

    flush_interval_ms=500,         # flush twice a second (optional)

    rotation_mode=RotationMode.SIZE,   # rotate when a file reaches 500 MiB

    max_file_size=500 * 1024 * 1024,
    replace=True,                  # delete any existing files in the directory

    include_types=None,            # write *all* supported types

)

# 3️⃣ Feed live objects (ticks, bars, etc.) as they arrive

for tick in live_tick_stream():
    writer.write(tick)             # `tick` can be QuoteTick, TradeTick, Bar, etc.

# 4️⃣ When the stream ends, close the writer to ensure everything is flushed

writer.close()

Key implementation details:

  • The fs_protocol parameter accepts any fsspec protocol ("file", "s3", "gcs", "memory"). The writer creates directories using fs.makedirs (lines 108-110).
  • When rotation triggers, the writer calls either _rotate_identifier_file (per-instrument) or _rotate_regular_file (global tables) to close the existing Arrow stream and open a new file with a fresh timestamp (lines 133-166).
  • After calling writer.close(), the property writer.is_closed returns True (lines 182-190).

Advanced Configuration Options

Parameter Description Example
include_types List of concrete Python types to restrict writing. Only those types will be persisted. include_types=[QuoteTick, TradeTick]
rotation_mode RotationMode.NO_ROTATION, SIZE, INTERVAL, or SCHEDULED_DATES. RotationMode.INTERVAL
rotation_interval pandas.Timedelta for interval-based rotation. pd.Timedelta(hours=6)
rotation_time & rotation_timezone Used with SCHEDULED_DATES to rotate at a specific wall-clock time in a given timezone. rotation_time=dt.time(0,0), rotation_timezone="America/New_York"
max_file_size Byte limit for SIZE rotation. 1024*1024*1024 (1 GiB)
replace If True, deletes any existing files under path before writing. replace=True

Summary

  • Use StreamingFeatherWriter from nautilus_trader/persistence/writer.py to persist live market data to Feather format with minimal overhead.
  • Configure file rotation via size or time intervals to manage disk usage and file granularity.
  • Leverage per-instrument file handling for organized data storage by instrument_id.
  • Set flush intervals to balance durability guarantees with I/O performance.
  • Support for fsspec protocols enables writing to local disk, S3, GCS, or in-memory filesystems.

Frequently Asked Questions

What file format does StreamingFeatherWriter use?

The writer persists data to Apache Feather (version 2) files, which provide high-performance columnar storage compatible with pandas and Arrow ecosystems. According to the Nautilus Trader source code, the writer uses ArrowSerializer.serialize_batch to convert objects into Arrow RecordBatch instances before writing.

How do I prevent memory leaks when writing high-frequency data?

Configure flush_interval_ms to ensure regular buffer flushes (default is 1000ms), and use rotation_mode=RotationMode.SIZE with a reasonable max_file_size to prevent individual files from growing too large. The writer automatically manages file handles through _rotate_identifier_file and _rotate_regular_file, closing old streams before opening new ones.

Can I write to cloud storage like S3?

Yes. Pass fs_protocol="s3" (or "gcs", "az", etc.) and ensure you have the appropriate fsspec implementation installed (e.g., s3fs). The writer uses fs.makedirs (lines 108-110 in writer.py) to create directories remotely, and all file operations go through the fsspec filesystem interface.

What happens to data if my application crashes?

Data written since the last flush may be lost. To minimize data loss, reduce flush_interval_ms (e.g., to 100ms), though this increases I/O overhead. Feather files are written transactionally per batch, so existing files remain valid even if the process terminates unexpectedly. Always call writer.close() during graceful shutdown to ensure final buffers are flushed and the is_closed property (lines 182-190) reflects the correct state.

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 →