How the RocketMQ Filter Engine Works: A Deep Dive into the Rust Implementation
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 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. This interface establishes the contract between the broker and filtering logic:
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. It maintains a static DashMap called FILTER_REGISTRY that maps type strings to Arc<dyn Filter> instances:
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 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):
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
Filtertrait 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_REGISTRYusesDashMapfor lock-free reads and fine-grained locking on writes, preventing contention during high-throughput message filtering. - Immutable factory: The
FilterFactorystruct 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:
- Consumer subscription: A consumer subscribes with a selector type (SQL-92 or tag) and expression string.
- Filter retrieval: The broker calls
FilterFactory::instance().get("SQL92")to obtain the appropriate filter implementation. - Expression compilation: The filter's
compilemethod parses the subscription string into a reusableExpressionobject, which the broker caches for the consumer session. - Message evaluation: For each incoming message, the broker invokes
expression.evaluate(message)to check if message properties match the subscription criteria. - Delivery decision: Messages evaluating to
trueare 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:
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:
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:
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
Filtertrait inrocketmq-filter/src/filter/filter_spi.rs. - FilterFactory manages a global, thread-safe registry using
DashMapandArc, 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"inrocketmq-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, 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.
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 →