How to Use the Cache and Message Bus for Inter-Component Communication in Nautilus Trader

Use the Cache in nautilus_trader/cache/cache.pyx as a shared in-memory database for state, and the MessageBus in nautilus_trader/common/component.pyx as a pub/sub broker for asynchronous commands and events between strategies, execution engines, and risk managers.

Nautilus Trader isolates core functionalities into independent components that must share state and exchange information without tight coupling. The nautechsystems/nautilus_trader repository provides two central services for this purpose: the Cache, which acts as an in-memory database, and the Message Bus, which implements a publish-subscribe pattern with request-reply capabilities. Understanding how to leverage these services enables you to build decoupled, scalable trading systems where strategies, execution engines, and custom services communicate efficiently.

Understanding the Cache and Message Bus Architecture

The Cache as Shared State Store

The Cache class, defined at line 81 in nautilus_trader/cache/cache.pyx, serves as the single source of truth for all transient system state. It holds market data snapshots, order books, positions, accounts, and custom objects that any component can query or update. All components receive a reference to the same Cache instance, guaranteeing a consistent view of the system state without direct object references between components.

The Message Bus as Event Broker

The MessageBus class, defined at line 2178 in nautilus_trader/common/component.pyx, implements a hierarchical pub/sub system for asynchronous communication. Components publish messages on topics such as commands.execution.submit_order or events.trades, and subscribers receive messages based on topic patterns. This decouples producers from consumers, allowing you to add or replace components without modifying existing code.

Setting Up the Cache and Message Bus

When initializing a TradingNode or custom component, you typically receive pre-configured instances of both services. However, for standalone testing or custom infrastructure, you can instantiate them directly:

from nautilus_trader.common.component import MessageBus
from nautilus_trader.cache.cache import Cache
from nautilus_trader.common.config import CacheConfig, MessageBusConfig
from nautilus_trader.core.rust.common import UUID4
from nautilus_trader.core.rust.time import Clock

# Initialize core services

clock = Clock()
msgbus = MessageBus(
    trader_id=UUID4(),
    clock=clock,
    config=MessageBusConfig(),
)
cache = Cache(config=CacheConfig())

The MessageBus requires a trader_id and clock for message timestamping and routing, while the Cache accepts configuration for database persistence options.

Publishing Commands via the Message Bus

Strategies and other components publish commands to trigger actions in remote engines. For example, a strategy submitting an order publishes to the execution topic:

class MyStrategy(Strategy):
    def on_bar(self, bar):
        # Construct the command

        submit_order = SubmitOrder(
            trader_id=self.trader_id,
            instrument_id=self.instrument_id,
            order_side=OrderSide.BUY,
            quantity=Quantity(100),
            price=Price(10_000.0),
            order_type=OrderType.LIMIT,
            time_in_force=TimeInForce.GTC,
        )
        # Publish to execution engine

        self._msgbus.publish("commands.execution.submit_order", submit_order)

The ExecutionEngine subscribes to "commands.execution.*" and processes the SubmitOrder command asynchronously.

Subscribing to Events and Data

Components register handlers to receive messages on specific topics. Use the subscribe method to attach a callback:

def _on_trade(self, trade: TradeTick):
    self._log.info(f"Received trade: {trade}")

# Register during component initialization

msgbus.subscribe(
    topic="events.execution.trade",
    handler=self._on_trade,
)

The MessageBus stores subscriptions in an internal mapping and invokes handlers synchronously when messages are published to matching topics.

Reading and Writing Shared State with the Cache

The Cache provides type-safe accessors for common trading objects and generic key-value storage for custom data:


# Store a custom indicator or calculator

cache.set("my_ma", MovingAverage(window=20))

# Retrieve from any component

ma = cache.get("my_ma")
latest_price = self.cache.quote_tick(self.instrument_id).price
ma.update(latest_price)

# Query trading state

position = self.cache.position(position_id)
order_book = self.cache.order_book(instrument_id)

The Cache class implements methods like set, get, position, and order_book as thin wrappers around internal dictionaries, providing O(1) access to shared state.

End-to-End Communication Flow

A complete interaction between strategy, execution engine, and cache follows this pattern:

Strategy (on_bar) ──► MessageBus.publish("commands.execution.submit_order")
ExecutionEngine (subscriber) ──► processes order
ExecutionEngine updates Cache ──► cache.orders[order_id] = order
ExecutionEngine publishes ──► MessageBus.publish("events.execution.order_updated")
Strategy (subscriber) ──► reads updated state from Cache
    order = self.cache.order(order_id)

This flow demonstrates how the Message Bus handles command routing and event notification, while the Cache maintains the authoritative state that components query after receiving notifications.

Advanced Configuration and Wildcard Patterns

Hierarchical Topic Matching

The MessageBus supports wildcard subscriptions for flexible event handling. According to the source code in component.pyx (lines 85-99), * matches any number of characters and ? matches a single character:


# Subscribe to all execution events

msgbus.subscribe("events.execution.*", handler)

# Subscribe to specific command patterns

msgbus.subscribe("commands.*.submit_order", handler)

Thread Safety and Performance

The MessageBus is not thread-safe and requires all components to run on the same event-loop thread, as documented in the class docstring (lines 22-25 of component.pyx). Publishing from external threads requires marshaling messages through the event loop.

Persistence Options

For production deployments requiring durability, configure the Cache with a CacheDatabaseFacade or the MessageBus with a RedisMessageBusDatabase. These adapters synchronize state to external storage, allowing recovery after process restarts without losing positions or pending orders.

Summary

  • The Cache in nautilus_trader/cache/cache.pyx provides a shared, in-memory database for market data, positions, and custom objects accessible to all components.
  • The Message Bus in nautilus_trader/common/component.pyx implements pub/sub and request-reply patterns for asynchronous command and event exchange.
  • Components communicate by publishing messages to hierarchical topics (e.g., commands.execution.submit_order) and subscribing to relevant event streams.
  • Use wildcard patterns (*, ?) for flexible subscription matching across topic hierarchies.
  • The architecture enforces single-threaded access to the message bus; all components must share the same event loop.
  • Optional persistence layers allow the cache and message bus to survive process restarts for long-running live trading.

Frequently Asked Questions

How do I share custom data between a strategy and an indicator?

Store the custom object in the Cache using cache.set("key", value) from your strategy or indicator initialization. Any component can retrieve it later with cache.get("key"). The Cache maintains strong references to objects until explicitly removed or overwritten.

What happens if I publish a message to a topic with no subscribers?

The Message Bus silently drops messages published to topics without registered listeners. This is by design—components should not assume synchronous delivery or that consumers exist. If you require guaranteed processing, implement a request-reply pattern or check subscription counts before publishing.

Can I use the Message Bus from multiple threads?

No. The Message Bus is not thread-safe and must only be used from the single event-loop thread that owns it, as enforced by the architecture in nautilus_trader/common/component.pyx. If you need to publish from external threads (e.g., a data feed callback), marshal the message through the event loop using loop.call_soon_threadsafe() or similar mechanisms provided by the platform's networking layer.

How do I persist cache state across process restarts?

Configure the Cache with a CacheDatabaseFacade pointing to a backing store (e.g., Redis or PostgreSQL) during initialization. When CacheConfig includes database settings, the cache synchronizes writes to external storage automatically. For the Message Bus, use a RedisMessageBusDatabase to log published messages, enabling replay or recovery of event streams after a restart.

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 →