# How the RocketMQ Filter Engine Works: A Deep Dive into the Rust Implementation

> Explore the RocketMQ filter engine's Rust implementation. Discover how it compiles subscriptions and evaluates messages using a pluggable SPI and SQL-92 for efficient filtering.

- Repository: [mxsm/rocketmq-rust](https://github.com/mxsm/rocketmq-rust)
- Tags: deep-dive
- Published: 2026-03-07

---

**The RocketMQ filter engine uses a pluggable SPI architecture with a global registry factory and a default SQL-92 implementation to compile consumer subscription expressions and evaluate them against message properties at runtime.**

The `rocketmq-filter` crate in the [mxsm/rocketmq-rust](https://github.com/mxsm/rocketmq-rust) repository provides the broker-side filtering subsystem. This engine determines whether messages match consumer subscription criteria before delivery, reducing unnecessary network traffic and client-side processing.

## Core Architecture of the RocketMQ Filter Engine

The filter engine follows a three-layer design: a trait-based Service Provider Interface (SPI), a singleton factory managing global access, and concrete implementations for specific expression languages.

### The Filter SPI (`Filter` Trait)

All filter implementations must satisfy the `Filter` trait defined in [`rocketmq-filter/src/filter/filter_spi.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-filter/src/filter/filter_spi.rs). This interface establishes the contract between the broker and filtering logic:

```rust
pub trait Filter: Send + Sync + fmt::Debug {
    fn compile(&self, expr: &str) -> Result<Box<dyn Expression>, FilterError>;
    fn of_type(&self) -> &str;
}

```

The **`compile`** method transforms textual expressions (such as `"age > 18 AND region = 'US'"`) into boxed `Expression` trait objects that can be evaluated repeatedly against messages. The **`of_type`** method returns a unique type identifier (e.g., `"SQL92"`) used for factory registration.

### The Global Registry (`FilterFactory`)

`FilterFactory` acts as a singleton registry for all filter implementations, located in [`rocketmq-filter/src/filter/filter_factory.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-filter/src/filter/filter_factory.rs). It maintains a static `DashMap` called `FILTER_REGISTRY` that maps type strings to `Arc<dyn Filter>` instances:

```rust
static FILTER_REGISTRY: LazyLock<DashMap<String, Arc<dyn Filter>>> = LazyLock::new(|| {
    let registry = DashMap::new();
    registry.insert("SQL92".to_string(), Arc::new(SqlFilter::new()));
    registry
});

```

The factory provides **lock-free reads** for filter retrieval and thread-safe registration for custom implementations. Key methods include `instance()` for singleton access, `register()` for adding new filters, and `get()` for lookup by type string.

### Default SQL-92 Implementation (`SqlFilter`)

The `SqlFilter` struct in [`rocketmq-filter/src/filter/filter_sql_filter.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-filter/src/filter/filter_sql_filter.rs) provides the built-in SQL-92 expression support. It implements `Filter` with `"SQL92"` as its type identifier (defined as a constant in [`rocketmq-common/src/common/filter/expression_type.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-common/src/common/filter/expression_type.rs)):

```rust
pub const SQL92: &'static str = "SQL92";

```

Currently, the `compile` method returns `unimplemented!()`, but the architecture expects a full SQL-92 parser producing an expression tree compatible with the `Expression` trait in `rocketmq-filter/src/expression/`.

## Thread-Safety and Concurrency Design

The RocketMQ filter engine achieves zero-cost concurrency through several mechanisms:

* **Send + Sync bounds**: The `Filter` trait requires implementations to be thread-safe, allowing filter instances to move between threads and be shared concurrently.
* **Arc wrapping**: All filters stored in the registry are wrapped in `Arc`, enabling shared ownership without lifetime complications.
* **DashMap storage**: The `FILTER_REGISTRY` uses `DashMap` for lock-free reads and fine-grained locking on writes, preventing contention during high-throughput message filtering.
* **Immutable factory**: The `FilterFactory` struct itself contains no mutable state; all mutability is confined to the interior of the static registry.

## End-to-End Message Filtering Flow

The filter engine integrates into the message consumption pipeline through the following steps:

1. **Consumer subscription**: A consumer subscribes with a selector type (SQL-92 or tag) and expression string.
2. **Filter retrieval**: The broker calls `FilterFactory::instance().get("SQL92")` to obtain the appropriate filter implementation.
3. **Expression compilation**: The filter's `compile` method parses the subscription string into a reusable `Expression` object, which the broker caches for the consumer session.
4. **Message evaluation**: For each incoming message, the broker invokes `expression.evaluate(message)` to check if message properties match the subscription criteria.
5. **Delivery decision**: Messages evaluating to `true` are delivered to the consumer; others are filtered out at the broker level, conserving network bandwidth.

## Working with the Filter Engine

### Retrieving the Built-In SQL-92 Filter

Access the default SQL-92 filter through the factory singleton and compile subscription expressions:

```rust
use rocketmq_filter::filter::{FilterFactory, Filter};
use std::sync::Arc;

fn main() -> Result<(), Box<dyn std::error::Error>> {
    let factory = FilterFactory::instance();

    let sql_filter: Arc<dyn Filter> = factory.get("SQL92")
        .expect("SQL92 filter should be registered");

    let expr = sql_filter.compile("age > 18 AND region = 'US'")?;
    Ok(())
}

```

### Registering a Custom Tag-Based Filter

Implement the `Filter` trait to create custom filtering logic and register it at runtime:

```rust
use rocketmq_filter::filter::{Filter, FilterFactory, FilterError};
use rocketmq_filter::expression::Expression;
use std::sync::Arc;

#[derive(Debug, Default)]
struct TagFilter;

impl Filter for TagFilter {
    fn compile(&self, expr: &str) -> Result<Box<dyn Expression>, FilterError> {
        Ok(Box::new(TagExpression::new(expr.to_string())))
    }

    fn of_type(&self) -> &str {
        "TAG"
    }
}

#[derive(Debug)]
struct TagExpression {
    allowed: Vec<String>,
}

impl TagExpression {
    fn new(tags: String) -> Self {
        let allowed = tags.split(',').map(|s| s.trim().to_string()).collect();
        Self { allowed }
    }
}

impl Expression for TagExpression {
    fn evaluate(&self, message: &rocketmq_common::message::Message) -> bool {
        if let Some(tags) = message.tags() {
            tags.iter().any(|t| self.allowed.contains(t))
        } else {
            false
        }
    }
}

fn main() {
    let factory = FilterFactory::instance();
    factory.register(Arc::new(TagFilter::default()));
}

```

### Listing Registered Filter Types

Introspect the factory to discover available filtering implementations:

```rust
use rocketmq_filter::filter::FilterFactory;

fn main() {
    let factory = FilterFactory::instance();
    for filter_type in factory.registered_types() {
        println!("Registered filter: {}", filter_type);
    }
}

```

## Summary

* The **RocketMQ filter engine** implements a pluggable SPI architecture centered on the `Filter` trait in [`rocketmq-filter/src/filter/filter_spi.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-filter/src/filter/filter_spi.rs).
* **FilterFactory** manages a global, thread-safe registry using `DashMap` and `Arc`, providing lock-free lookup and runtime registration of custom filters.
* The default **SQL-92 filter** (`SqlFilter`) provides the standard expression language for message filtering, identified by the constant `"SQL92"` in `rocketmq-common`.
* All components enforce **Send + Sync** bounds, ensuring safe concurrent access across the broker's thread pool without performance bottlenecks.
* The broker uses this engine to **compile consumer subscriptions once** and **evaluate them against each message**, filtering at the broker level to minimize network overhead.

## Frequently Asked Questions

### How does the RocketMQ filter engine handle concurrent access to filter implementations?

The engine uses **Arc** to share filter instances across threads and **DashMap** for the registry, providing lock-free reads and fine-grained locking only for writes. All filter implementations must implement **Send + Sync**, guaranteeing thread-safety at the trait level.

### Can I implement a custom filter type for the RocketMQ broker?

Yes. Implement the **Filter** trait from [`rocketmq-filter/src/filter/filter_spi.rs`](https://github.com/mxsm/rocketmq-rust/blob/main/rocketmq-filter/src/filter/filter_spi.rs), defining the `compile` method to parse your expression syntax and `of_type` to return a unique identifier. Register your implementation at runtime using `FilterFactory::instance().register(Arc::new(YourFilter))`.

### What is the difference between the filter compilation and evaluation phases?

**Compilation** occurs once when a consumer subscribes, transforming the expression string (e.g., `"age > 18"`) into a boxed **Expression** trait object via `Filter::compile`. **Evaluation** happens for every message, calling `Expression::evaluate` to check if the message matches the compiled criteria, returning a boolean result to the broker.