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

> Persist live market data in real-time using Nautilus Trader streaming writers. Easily save data to Feather files with automatic rotation and per-instrument handling.

- Repository: [Nautech Systems/nautilus_trader](https://github.com/nautechsystems/nautilus_trader)
- Tags: how-to-guide
- Published: 2026-02-16

---

**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`](https://github.com/nautechsystems/nautilus_trader/blob/main/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`](https://github.com/nautechsystems/nautilus_trader/blob/main/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

```python
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`](https://github.com/nautechsystems/nautilus_trader/blob/main/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`](https://github.com/nautechsystems/nautilus_trader/blob/main/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.