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_protocolparameter accepts any fsspec protocol ("file","s3","gcs","memory"). The writer creates directories usingfs.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 propertywriter.is_closedreturnsTrue(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
StreamingFeatherWriterfromnautilus_trader/persistence/writer.pyto 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →