NautilusTrader
Concepts
These docs track the unreleased nightly build and may change without notice. Switch to the latest stable docs.

Message Bus

The MessageBus enables communication between system components through message passing. This design creates a loosely coupled architecture where components interact without direct dependencies.

The messaging patterns include:

  • Point-to-Point
  • Publish/Subscribe
  • Request/Response

Messages exchanged via the MessageBus fall into three categories:

  • Data
  • Events
  • Commands

Topic hierarchy

Nautilus keeps market data topics under the data root. Live data publications use the direct data.<kind>... topics, for example data.book.deltas.XCME.ESZ24.

When requested, replayed, or workflow-generated data flows over the message bus as topic-addressable data, the DataEngine publishes it under data.pipeline.<kind>.... Long requests, grouped requests, and aggregation chains can split, transform, and fan data back in before the parent request completes. These messages are still data messages, but they do not claim the same live ordering and timing semantics as normal real-time publications. For example, book deltas on the pipeline path use data.pipeline.book.deltas.XCME.ESZ24.

Correlated request responses are delivered through response handlers keyed by correlation ID. The data.response topic is a capture channel for response publications, not the pipeline data path.

Message integrity

Once a message is created, its fields must not be mutated. This includes container fields such as params maps. Components can read a message and derive local state from it, but they must not rewrite the original.

Immutable messages keep every consumer seeing the same input, preserve what was true at emission time, and remove a class of shared-state races. Replay, debugging, and audit all depend on messages remaining stable after dispatch.

Three ownership rules follow from this:

  • Caller-supplied request options stay on the message.
  • Response metadata returned to the caller stays on the response.
  • Component workflow state (bounded date ranges, grouping state, replay cursors, counters, processing flags) stays in component-owned context keyed by message or request ID.

When a component needs a derived message, it creates a new one with the required values instead of rewriting the original.

Data and signal publishing

While the MessageBus is a lower-level component that users typically interact with indirectly, DataActor and Strategy provide typed methods built on top of it:

def publish_data(self, data_type: DataType, data: CustomData) -> None:
def publish_signal(self, name: str, value, ts_event: int = 0) -> None:

These methods publish custom data and signals without exposing the raw message bus to Python.

Direct access

The Python DataActor and Strategy APIs do not expose self.msgbus. Use custom data or signals for supported Python component messaging. Rust components can use the typed message-bus facade directly.

Messaging styles

NautilusTrader is an event-driven framework where components communicate by sending and receiving messages. Understanding the different messaging styles helps when building trading systems.

This guide explains the three primary messaging patterns available in NautilusTrader:

Messaging stylePurposeBest for
Custom data publish/subscribeStructured trading data exchangeTrading metrics, indicators, data needing persistence
Signal publish/subscribeLightweight notificationsSimple alerts, flags, and status updates
Rust MessageBus publish/subscribeLow‑level, typed topic communicationNative runtime components

Each approach serves different purposes. Use this guide to decide which pattern to use.

Rust MessageBus publish/subscribe to topics

Concept

The MessageBus is the central hub for all messages in NautilusTrader. Rust components can publish typed messages to named topics and subscribe handlers to those topics. This low‑level interface is not part of the Python actor or strategy surface.

Benefits and use cases

Direct message-bus access is for native components that need:

  • Cross-component communication within the system.
  • Flexibility to define typed topics and payloads.
  • Decoupling between publishers and subscribers who don't need to know about each other.
  • Global reach where messages can be received by multiple subscribers.
  • Working with events that do not fit the data actor model.
  • Advanced scenarios requiring full control over messaging.

Considerations

  • You must track topic names manually (typos could result in missed messages).
  • You must define handlers manually.

Custom data publish/subscribe

Concept

Custom data exchanges structured values between data actors and strategies. A CustomData value carries a DataType, payload, event timestamp, and initialization timestamp for routing and event ordering.

Benefits and use cases

The Data publish/subscribe approach works well when you need:

  • Exchange of structured trading data like market data, indicators, custom metrics, or option greeks.
  • Proper event ordering via built-in timestamps (ts_event, ts_init) crucial for backtest accuracy.
  • Data persistence and serialization through registered custom data classes, integrating with NautilusTrader's data catalog system.
  • Standardized trading data exchange between system components.

Considerations

  • The payload must expose ts_event and ts_init.
  • Persistence requires registering a serializable custom data class.

Quick overview code

from dataclasses import dataclass

from nautilus_trader.model import CustomData
from nautilus_trader.model import DataType


@dataclass
class GreeksData:
    delta: float
    gamma: float
    ts_event: int
    ts_init: int


data_type = DataType("GreeksData")
data = CustomData(
    data_type,
    GreeksData(
        delta=0.75,
        gamma=0.1,
        ts_event=1_630_000_000_000_000_000,
        ts_init=1_630_000_000_000_000_000,
    ),
)
self.publish_data(data_type, data)

self.subscribe_data(data_type)


def on_data(self, data: CustomData) -> None:
    if data.data_type == data_type:
        greeks = data.data
        self.log.info(f"Delta: {greeks.delta}, Gamma: {greeks.gamma}")

See Custom data for registration and persistence.

Signal publish/subscribe

Concept

Signals are a lightweight way to publish and subscribe to simple notifications within the actor framework. This is the simplest messaging approach, requiring no custom class definitions.

Benefits and use cases

The Signal messaging approach works well when you need:

  • Simple, lightweight notifications/alerts like "RiskThresholdExceeded" or "TrendUp".
  • Quick, on-the-fly messaging without defining custom classes.
  • Broadcasting alerts or flags as simple primitive values.
  • Easy API integration with straightforward methods (publish_signal, subscribe_signal).
  • Multiple subscriber communication where all subscribers receive signals when published.
  • Minimal setup overhead with no class definitions required.

Considerations

  • Each signal carries a single value. The value is converted to a string on publish, so handlers always receive signal.value as a str; complex data structures are not preserved.
  • Differentiate between signals in the on_signal handler with signal.name.

Quick overview code

# Define signal constants for better organization (optional but recommended)
import types

from nautilus_trader.common import LogColor
from nautilus_trader.core.datetime import unix_nanos_to_dt

signals = types.SimpleNamespace()
signals.NEW_HIGHEST_PRICE = "NewHighestPriceReached"
signals.NEW_LOWEST_PRICE = "NewLowestPriceReached"

# Subscribe from a DataActor or Strategy
self.subscribe_signal(signals.NEW_HIGHEST_PRICE)
self.subscribe_signal(signals.NEW_LOWEST_PRICE)

# Publish from a DataActor or Strategy
self.publish_signal(
    name=signals.NEW_HIGHEST_PRICE,
    value=signals.NEW_HIGHEST_PRICE,  # value can be the same as name for simplicity
    ts_event=bar.ts_event,  # timestamp from triggering event
)


# Handler (fixed callback name)
def on_signal(self, signal):
    match signal.name:
        case signals.NEW_HIGHEST_PRICE:
            self.log.info(
                f"New highest price was reached. | "
                f"Signal value: {signal.value} | "
                f"Signal time: {unix_nanos_to_dt(signal.ts_event)}",
                color=LogColor.GREEN,
            )
        case signals.NEW_LOWEST_PRICE:
            self.log.info(
                f"New lowest price was reached. | "
                f"Signal value: {signal.value} | "
                f"Signal time: {unix_nanos_to_dt(signal.ts_event)}",
                color=LogColor.RED,
            )

Summary and decision guide

Decision guide: Which style to choose?

Use caseRecommended approachSetup required
Native system‑level communicationRust MessageBus publish/subscribeTyped topic and handler
Structured Python component dataDataActor custom data methodsDataType, CustomData, and on_data()
Simple Python alerts and notificationsDataActor signal methodsSignal name and on_signal()

External egress and ingress

The MessageBus can write serialized messages to external streams. This section describes the external egress and ingress sides of the external bus. Rust-native live nodes use injected MessageBusExternalEgress and MessageBusExternalIngress surfaces, so the core node does not depend on Redis, a broker, shared-memory implementation, or socket protocol.

Redis is the built‑in external backing for serializable messages. The minimum supported Redis version is 6.2, required for the MINID stream trimming used by autotrim.

When external egress is configured, outgoing publish messages are first dispatched to in-process subscribers, then serialized into the existing BusMessage wire record:

  • topic: the exact message bus topic used by the internal publish call, for example data.quotes.BINANCE.BTCUSDT or events.order.S-001.
  • type: the canonical payload type name, for example QuoteTick or OrderEventAny.
  • encoding: the payload encoding selected from the message bus encoding policy.
  • payload: serialized bytes encoded with the selected encoding.

An external producer that writes directly to a Redis stream must include topic, type, and payload. The topic must be a valid publish topic and cannot contain * or ?. The encoding field is optional and defaults to JSON when omitted. The receiving node skips entries without type because it cannot select a payload decoder.

External egress receives that record as publish(BusMessage). This outbound call must not block the node's bus thread. Bounded egress implementations drop on a full queue instead of applying back-pressure to the trading loop. Closing the message bus closes the configured egress.

Inbound external streams are exposed through the separate Rust MessageBusExternalIngress trait. Ingress yields the same BusMessage { topic, payload_type, encoding, payload } shape. republish_external_message decodes supported inbound messages and republishes them internally without forwarding the message back out. The inbound payload type must first be registered for streaming on the receiving message bus; unregistered types are skipped without decoding.

For custom data, egress writes and ingress expects an envelope in the Redis payload field, not the bare custom object. The canonical JSON envelope is:

{
  "type": "MyData",
  "data_type": {
    "type_name": "MyData",
    "metadata": {
      "source": "external"
    },
    "identifier": "optional-storage-key"
  },
  "payload": {
    "value": 42,
    "ts_event": 0,
    "ts_init": 0
  }
}

The envelope requires type and payload; data_type is optional on ingress and defaults to the message type with no metadata or identifier. In the canonical emitted form, data_type.type_name uses the same custom type name, metadata is an object that may be empty, and identifier is present only when assigned. The envelope type must match the Redis stream type. The envelope payload is the bare object passed to the registered class's from_json(...) method. MessagePack uses the same map fields encoded as MessagePack bytes.

For Python custom data, register the class before starting the node:

from nautilus_trader.model import register_custom_data_class

register_custom_data_class(MyData)

The external‑client subscription registers the payload type for streaming, while register_custom_data_class(...) installs the process‑wide JSON decoder. Both registrations are required. See Custom data for the class requirements.

For Redis, messages are transmitted via a Multiple-Producer Single-Consumer (MPSC) channel to a separate Rust task. That task writes the message to Redis streams.

Offloading I/O to a separate task keeps the publishing thread unblocked.

With MessagePack or JSON, Rust-native external egress forwards serializable typed publications. This includes instruments, quotes, trades, bars, book deltas, depth-10 snapshots, mark/index/funding updates, option greeks (OptionGreeks), account state, portfolio snapshots, order events, position events, and custom data. With the defi feature this also includes DeFi blocks, pools, liquidity updates, fee collects, and flash events. Full order book snapshots, GreeksData records, option chain slices, and DeFi pool swaps are not forwarded because those types do not implement Serde serialization.

With SBE or Cap'n Proto, Rust-native external egress forwards the built-in market data payloads with schema codecs: quotes, trades, bars, book deltas, depth-10 snapshots, mark price updates, index price updates, funding rate updates, and option greeks. Other payload types are dropped with a debug log when those schema encodings are selected.

Configuration

The message bus external backing technology uses a behavior config plus a technology-owned backing config. MessageBusConfig controls message bus behavior. RedisMessageBusConfig owns Redis connection settings and implements MessageBusBackingFactory.

use nautilus_common::{
    enums::SerializationEncoding,
    msgbus::{MessageBusBackingFactory, MessageBusConfig},
};
use nautilus_infrastructure::redis::msgbus::RedisMessageBusConfig;

let config = MessageBusConfig {
    encoding: SerializationEncoding::Json,
    encoding_market_data: Some(SerializationEncoding::Sbe),
    timestamps_as_iso8601: true,
    buffer_interval_ms: Some(100),
    autotrim_mins: Some(30),
    use_trader_prefix: true,
    use_trader_id: true,
    use_instance_id: false,
    streams_prefix: "streams".to_string(),
    types_filter: Some(vec!["QuoteTick".to_string(), "TradeTick".to_string()]),
    ..Default::default()
};

let redis_config = RedisMessageBusConfig::default();
let backing = redis_config.create(trader_id, instance_id, config.clone())?;

Existing Rust callers can continue using RedisMessageBusFactory::new(redis_config), which delegates to the config implementation.

Backing config

A RedisMessageBusConfig is required when using the built-in Redis backing. For a default Redis setup on the local loopback you can pass RedisMessageBusConfig::default().

Redis selection is explicit in the Rust type. The config does not use a user-facing selector such as type = "redis" or backing_type = "redis".

Rust-native callers that inject MessageBusExternalEgress with LiveNodeBuilder::with_external_msgbus_egress pass concrete connection details when they construct that egress surface. The core message bus does not require a RedisMessageBusConfig for injected egress.

The Rust live runtime accepts external_streams in MessageBusConfig, and consumes inbound BusMessages when callers inject a MessageBusExternalIngress with LiveNodeBuilder::with_external_ingress. The config names the external stream keys; the injected ingress is the concrete runtime source. Rust callers can install RedisMessageBusConfig with LiveNodeBuilder::with_external_msgbus_factory. Building fails when a factory is combined with separately injected egress or ingress. A factory always installs egress and creates ingress only when external_streams is non‑empty.

Python exposes the same builder method for built‑in backing configs, currently RedisMessageBusConfig. The existing RedisMessageBusFactory wrapper remains supported. Python does not accept arbitrary factory classes.

The built-in Redis ingress starts each configured stream at the current timestamp, so entries that already exist when the node starts are not replayed. After startup it advances the last-seen ID for each stream and preserves those IDs across connection retries. Use cache recovery or the event store when durable pre-start replay is required; external_streams provides live forwarding, not a consumer-group backlog.

Encoding

Rust-native external message bus egress supports these encoding names:

  • JSON (json)
  • MessagePack (msgpack)
  • Cap'n Proto (capnp, with the Rust capnp feature)
  • SBE (sbe, with the Rust sbe feature)

Use the encoding config option to control the message writing encoding. Use encoding_market_data to override the encoding for market data payloads backed by the external bus binary codecs. Use encoding_builtin to override account state, portfolio snapshot, order event, and position event payloads. Custom and unmapped payload types always use encoding.

MessageBusConfig::validate requires the default encoding to support custom payloads, so it must be JSON or MessagePack. Category overrides must be supported by every published payload type in that category. SBE and Cap'n Proto can currently be used only for encoding_market_data, and only when the matching Rust feature is enabled. encoding_builtin = "sbe" and encoding_builtin = "capnp" fail validation until those schema codecs cover the built-in event category.

The Redis cache payload path supports MessagePack and JSON only. SBE and Cap'n Proto are schema payload encodings for Rust-native external message bus egress, not Redis cache encodings, and selecting either for a Redis cache payload is an error.

The json encoding is used by default for human readability and interoperability. Use msgpack when payload size and serialization performance are a primary concern.

Timestamp formatting

By default timestamps are formatted as UNIX epoch nanosecond integers. Alternatively you can configure ISO 8601 string formatting by setting timestamps_as_iso8601 to true.

Message stream keys

Message stream keys identify individual trader nodes and organize messages within streams. The trader- prefix, trader ID, and instance ID segments are optional and controlled by the options below; the streams prefix is always included. With every segment enabled, a trader key has the following structure:

trader-{trader_id}:{instance_id}:{streams_prefix}

With the default options (use_trader_prefix and use_trader_id enabled, use_instance_id disabled) the base stream key is trader-{trader_id}:{streams_prefix}.

These options control Redis stream keys. They do not rewrite the topic passed to an injected MessageBusExternalEgress; that topic remains the internal message bus publish topic. When stream_per_topic is True, Redis egress appends the topic to the stream key. When it is False, Redis stores all messages on the base stream key and keeps the topic as a message field.

The following options are available for configuring message stream keys:

Trader prefix

If the key should begin with the trader- prefix.

Trader ID

If the key should include the trader ID for the node.

Instance ID

Each trader node is assigned a unique instance ID, which is a UUIDv4. This instance ID helps distinguish individual traders when messages are distributed across multiple streams. You can include the instance ID in the trader key by setting the use_instance_id configuration option to True. This is particularly useful when you need to track and identify traders across various streams in a multi-node trading system.

Streams prefix

The streams_prefix string enables you to group all streams for a single trader instance or organize messages for multiple instances. Configure this by passing a string to the streams_prefix configuration option, ensuring other prefixes are set to false.

Stream per topic

Indicates whether the producer will write a separate stream for each topic. This is particularly useful for Redis backings, which do not support wildcard topics when listening to streams. If set to False, all messages will be written to the same stream.

Redis does not support wildcard stream topics. For better compatibility with Redis, it is recommended to set this option to False.

Types filtering

When messages are published on the message bus, they are serialized and written to a stream if a backing for the message bus is configured and enabled. To prevent flooding the stream with data like high-frequency quotes, you may filter out certain types of messages from external publication.

To enable this filtering mechanism, pass a list of payload type names to the types_filter parameter in the message bus configuration. Listed types are excluded from external publication.

from nautilus_trader.config import MessageBusConfig

# Create a MessageBusConfig instance with types filtering
message_bus = MessageBusConfig(types_filter=["QuoteTick", "TradeTick"])

Stream auto-trimming

Use autotrim_mins to set a lookback window in minutes and autotrim_maxlen to set an approximate maximum number of entries for each Redis stream. You can configure either policy or both. When both are set, the message bus removes entries that exceed either the time window or the entry-count threshold.

Redis applies autotrim_maxlen with approximate trimming for better write performance, so a stream may contain slightly more entries than the configured threshold.

The Redis implementation trims each stream at most once per minute, so entries can remain up to roughly a minute longer than the autotrim_mins window.

External streams

The message bus within a LiveNode (node) is referred to as the "internal message bus". A producer node is one which publishes messages onto an external stream (see external egress and ingress). The consumer node listens to external streams to receive and publish deserialized message payloads on its internal message bus.

Set the LiveDataEngineConfig.external_clients with the list of client_ids intended to represent the external streaming clients. The DataEngine will filter out subscription commands for these clients, ensuring that the external streaming provides the necessary data for any subscriptions to these clients. When the Rust DataEngine skips an external-client subscription, it registers the corresponding streaming payload type for inbound republishing on the message bus.

Example configuration

The following example details a streaming setup where a producer node publishes Binance data externally, and a downstream consumer node publishes these data messages onto its internal message bus.

Producer node

We configure the MessageBus of the producer node to publish to a "binance" stream. The settings use_trader_id, use_trader_prefix, and use_instance_id are all set to false to ensure a simple and predictable stream key that the consumer nodes can register for.

let message_bus = MessageBusConfig {
    use_trader_id: false,
    use_trader_prefix: false,
    use_instance_id: false,
    streams_prefix: "binance".to_string(), // <---
    stream_per_topic: false,
    autotrim_mins: Some(30),
    ..Default::default()
};

let redis_config = RedisMessageBusConfig {
    connection_timeout: 2,
    response_timeout: 2,
    ..Default::default()
};

let mut node = LiveNode::builder(trader_id, Environment::Live)?
    .with_msgbus_config(message_bus)
    .with_external_msgbus_factory(Box::new(redis_config))
    .build()?;
node.run().await?;

Consumer node

We configure the MessageBus of the consumer node to receive messages from the same "binance" stream. A RedisMessageBusConfig creates ingress from external_streams, and LiveNode::run publishes the received messages onto the node's internal message bus. We declare the client ID "BINANCE_EXT" as an external client so the DataEngine does not attempt to send data commands to this client ID.

let data_engine = LiveDataEngineConfig {
    external_clients: Some(vec![ClientId::from("BINANCE_EXT")]),
    ..Default::default()
};

let message_bus = MessageBusConfig {
    external_streams: Some(vec!["binance".to_string()]), // <---
    ..Default::default()
};

let redis_config = RedisMessageBusConfig {
    connection_timeout: 2,
    response_timeout: 2,
    ..Default::default()
};

let mut node = LiveNode::builder(trader_id, Environment::Live)?
    .with_data_engine_config(data_engine)
    .with_msgbus_config(message_bus)
    .with_external_msgbus_factory(Box::new(redis_config))
    .build()?;
node.run().await?;
  • Actors - Actors use the message bus for event handling.
  • Architecture - Message bus role in system architecture.

On this page