apache/rocketmq
60.5
Adequate · 24 September 2026
222.5k
lines of production code
Java
primary language
4
measurements over time
What this system is
This system is a distributed message broker platform that manages asynchronous message passing, consumption, and transactional integrity across clustered nodes. It provides a comprehensive security framework with ACL 2.0 for authentication and authorization, alongside a proxy layer that supports gRPC and remoting protocols for flexible client access. The architecture emphasizes high availability through controller-managed master-slave replication, RocksDB-backed configuration storage, and advanced consumption models like Lite Topics and Pop consumption.
How it got here
2013–2017 — Broker architecture and build system overhaul
79 changes.
This period focused on restructuring the RocketMQ codebase, introducing a Maven multi-module layout and migrating the build infrastructure to Bazel. Significant architectural changes were made to the broker, including the implementation of a plugin-based processor model, RocksDB-backed configuration storage, and enhanced client connection management. The work also established foundational support for new messaging features such as SQL92 filtering, OpenMessaging compatibility, and various consumption modes like Lite and LMQ.
2018–2022 — Proxy v2 and Controller architecture
88 changes.
This period focused on introducing a standalone Controller module with jRaft support for broker management and implementing a comprehensive gRPC v2 Proxy with multi-protocol remoting. Significant work also included refactoring the broker's transactional message handling, adding OpenTelemetry metrics, and establishing Bazel build support.
2023–2026 — ACL 2.0 and infrastructure modernization
41 changes.
This period focused on implementing the comprehensive ACL 2.0 authentication and authorization framework, including migration tools and persistent storage backends. Concurrently, the project modernized core infrastructure by migrating broker configurations and state to RocksDB, introducing a JRaft-based controller, and adding tiered storage capabilities.
Features
Add BitsArray, BloomFilter, and BloomFilterData utility classes
The filter module now includes new utility classes to support efficient bit manipulation and probabilistic filtering. BitsArray provides a wrapper for byte arrays to allow easy single-bit operations (set, get, and, or, xor, not). BloomFilter implements a simple Bloom filter using Guava's MurmurHash3 to calculate bit positions for strings, enabling fast membership checks. BloomFilterData holds the bit positions and total bit count generated by the filter, facilitating data exchange and validation between filter instances.
filter/src/main/java/org/apache/rocketmq/filter/util · high confidence
Add FileWatchService for monitoring configuration file changes
The srvutil module now includes a new FileWatchService that monitors specified files for content changes using MD5 hashing. When a watched file is modified, the service triggers a listener callback, enabling dynamic reloading of configuration or certificates without requiring a restart. The service ignores file deletion events to ensure stability during file replacement operations.
rocketmq-srvutil · high confidence
Add OpenMessaging (OMS) example code
Added three new example classes—SimpleProducer, SimplePullConsumer, and SimplePushConsumer—in the openmessaging package to demonstrate how to use the OpenMessaging API with RocketMQ. These examples show how to create a producer for sending messages (sync, async, and oneway) and how to create pull and push consumers for receiving messages, using the OMS access point URL format.
example/src/main/java/org/apache/rocketmq/example/openmessaging · high confidence
Add OpenMessaging Promise implementation for RocketMQ
The \openmessaging\ module now includes a new \DefaultPromise\ class and \FutureState\ enum in the \io.openmessaging.rocketmq.promise\ package. This implementation provides the core asynchronous promise/future behavior required to support the OpenMessaging specification within RocketMQ, handling state transitions (doing, done, cancelled), result retrieval with timeouts, and listener notifications.
openmessaging/src/main/java/io/openmessaging/rocketmq/promise · high confidence
Add OpenMessaging domain model implementations for RocketMQ
This change introduces the core domain objects required to implement the OpenMessaging specification within the RocketMQ connector. It adds \BytesMessageImpl\ to handle binary message bodies and headers, \SendResultImpl\ to expose message IDs and properties from send operations, and \ConsumeRequest\ to encapsulate the context (message, queue, process queue) for consumption tasks. Additionally, it defines \NonStandardKeys\ and \RocketMQConstants\ to manage RocketMQ-specific metadata (such as consumer/producer groups and delivery timing) that maps to the OpenMessaging abstraction.
openmessaging/src/main/java/io/openmessaging/rocketmq/domain · high confidence
Add OpenMessaging utility classes for bean population and message conversion
Introduces BeanUtils and OMSUtil helper classes in the OpenMessaging RocketMQ integration. BeanUtils provides reflection-based population of JavaBeans from Properties and KeyValue objects, supporting primitive type conversion. OMSUtil handles bidirectional conversion between OpenMessaging BytesMessage and RocketMQ Message objects, mapping system and user headers, and converting send results.
openmessaging/src/main/java/io/openmessaging/rocketmq/utils · high confidence
Add batch message sending examples
New example classes, SimpleBatchProducer and SplitBatchProducer, are added to demonstrate how to send messages in batches. SimpleBatchProducer shows sending a small batch of messages, while SplitBatchProducer illustrates handling large batches by splitting them into smaller chunks to stay within size limits.
example/src/main/java/org/apache/rocketmq/example/batch · high confidence
Add default configuration files for broker, proxy, and tools
The distribution package now includes default configuration files for the broker (broker.conf), the proxy (rmq-proxy.json), and the tools module (tools.yml). These files provide out-of-the-box settings for cluster names, broker roles, and default access credentials, simplifying initial setup and deployment.
distribution/conf · high confidence
Add gRPC v2 EndTransaction activity
The proxy now exposes a new gRPC v2 endpoint for ending transactions. This change introduces the \EndTransactionActivity\ class, which validates the transaction ID and topic, maps the requested resolution (commit or rollback) to the internal transaction status, and delegates the operation to the messaging processor. Users can now complete distributed transactions via the v2 gRPC interface.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/transaction · high confidence
Add transactional message example
Added example code for using RocketMQ transactional messages, including a producer that sends messages with a custom transaction listener implementation to handle local transaction execution and status checking.
example/src/main/java/org/apache/rocketmq/example/transaction · high confidence
Added ACL 1.0 to 2.0 migration tooling
The auth module now includes an \AuthMigrator\ component that automatically converts legacy ACL 1.0 configurations (defined in YAML files under the \conf/acl\ directory) into the new ACL 2.0 format. This migration reads existing \PlainAccessConfig\ data—such as access keys, secret keys, admin status, and topic/group permissions—and creates corresponding users and authorization policies in the new system. The migration is controlled by the \authConfig.isMigrateAuthFromV1Enabled\ flag and runs asynchronously to ensure existing access control settings are preserved during the transition.
auth/src/main/java/org/apache/rocketmq/auth/migration · high confidence
Added MessageRequestModeManager for pop consuming mode configuration
The broker now includes a new MessageRequestModeManager component to manage and persist the message request mode settings for topics and consumer groups. This manager stores configuration data in a concurrent map structure and handles serialization/deserialization to support the Pop Consuming feature, allowing the broker to track and apply specific consumption modes.
broker/src/main/java/org/apache/rocketmq/broker/loadbalance · high confidence
Added OpenTracing and standard trace message examples
The tracemessage example package now includes new demonstration programs for integrating with OpenTracing (Jaeger). OpenTracingProducer, OpenTracingPushConsumer, and OpenTracingTransactionProducer show how to register OpenTracing hooks to trace message sending, consumption, and transactional operations. Additionally, TraceProducer and TracePushConsumer provide standard examples for basic message tracing functionality.
example/src/main/java/org/apache/rocketmq/example/tracemessage · high confidence
Added PR merge utility script
A new Python script, dev/merge\_rocketmq\_pr.py, has been added to the repository. This tool automates the process of merging GitHub pull requests into the Apache RocketMQ project by handling git operations such as fetching branches, resolving conflicts, creating squashed commits with proper attribution, and pushing the result to the upstream repository.
dev · high confidence
Added SQL and Tag filter example code
New example files have been added to demonstrate message filtering capabilities: SqlFilterProducer and SqlFilterConsumer show how to use SQL expressions to filter messages based on tags and properties, while TagFilterProducer and TagFilterConsumer illustrate filtering by specific tags using logical operators.
example/src/main/java/org/apache/rocketmq/example/filter · high confidence
Added Windows and Linux deployment scripts for controller quick-start and independent modes
New shell (.sh) and Windows batch (.cmd) scripts have been added to the distribution/bin/controller directory to simplify local testing and deployment. Users can now quickly start a full stack (NameServer and Brokers) using the \fast-try\ scripts, or deploy a standalone 3-node controller cluster independently using \fast-try-independent-deployment\. Additionally, scripts for running NameServers in a controller plugin configuration (\fast-try-namesrv-plugin\) are provided for both platforms, enabling easier validation of controller modes without manual command-line invocation.
distribution/bin/controller · high confidence
Added broadcast push consumer example with default NameServer address
A new PushConsumer example has been added to the broadcast package, demonstrating how to configure a consumer for broadcasting message delivery. The example sets the message model to BROADCASTING, subscribes to TopicTest with specific tag filters, and includes a commented-out line to set the default NameServer address (127.0.0.1:9876), addressing previous issues where the example lacked a default configuration.
example/src/main/java/org/apache/rocketmq/example/broadcast · high confidence
Added configuration files for 2-master 2-slave broker deployments
New property files have been added to the distribution configuration for both asynchronous and synchronous 2-master 2-slave broker topologies. These files define the cluster settings for broker-a and broker-b, specifying their roles (ASYNC\_MASTER, SYNC\_MASTER, or SLAVE) and ensuring all brokers use ASYNC\_FLUSH for disk persistence, allowing users to deploy these specific high-availability architectures.
distribution/conf/2m-2s-async, distribution/conf/2m-2s-sync · high confidence
Added configuration files for 2-master no-slave broker deployment
New configuration files have been added to the distribution package to support a two-master, no-slave broker topology. This includes specific property files for broker-a, broker-b, and a trace-enabled broker (broker-trace), all configured as ASYNC\_MASTER nodes within the DefaultCluster with ASYNC\_FLUSH disk settings.
distribution/conf/2m-noslave · high confidence
Added configuration files for controller mode deployments
New configuration files are now included in the distribution to support running the controller in various modes. This includes standalone and 3-node cluster configurations for the controller itself, as well as setups where the controller runs as a plugin within the NameServer. Additionally, quick-start configuration files for brokers and NameServers are provided to facilitate easy setup of controller-mode environments.
distribution/conf/controller · high confidence
Added core classes for SQL92 message filtering support
Introduced the \UnaryType\ enum and \PropertyExpression\ class within the filter module to support message filtering based on SQL92. The \UnaryType\ enum defines unary operators such as NEGATE, IN, NOT, BOOLEANCAST, and LIKE, while \PropertyExpression\ provides the logic to evaluate property names against the message context, forming the foundational building blocks for the new filtering capability.
filter/src/main/java/org/apache/rocketmq/filter/constant, rocketmq-filter · high confidence
Added dledger cluster management script
A new shell script, fast-try.sh, has been added to the distribution/bin/dledger directory to facilitate the quick start and stop of a dledger cluster. This script automates the launching and termination of the NameServer and three Broker instances (n0, n1, n2) by executing the corresponding mqbroker and mqnamesrv binaries with specific Java memory configurations, ensuring the required configuration files are present before starting services.
distribution/bin/dledger · high confidence
Added order message producer and consumer examples
New example files (Consumer.java and Producer.java) were added to the ordermessage package to demonstrate how to send and consume ordered messages using RocketMQ's DefaultMQProducer and DefaultMQPushConsumer. The producer example shows how to use a MessageQueueSelector to ensure messages with the same order ID are sent to the same queue, while the consumer example illustrates orderly consumption with auto-commit and suspension logic.
example/src/main/java/org/apache/rocketmq/example/ordermessage · high confidence
Added quickstart examples for sending and consuming messages
New example files, Consumer.java and Producer.java, have been added to the quickstart package to demonstrate how to use DefaultMQPushConsumer and DefaultMQProducer. These examples show users how to configure a consumer group, subscribe to a topic, and register a message listener, as well as how to instantiate a producer, send messages synchronously, and handle the resulting SendResult.
example/src/main/java/org/apache/rocketmq/example/quickstart · high confidence
Added request-response RPC example code
New example files (AsyncRequestProducer, RequestProducer, ResponseConsumer) were added to demonstrate the request-response pattern in RocketMQ. These examples show how to send synchronous and asynchronous requests using DefaultMQProducer.request() and how to handle incoming messages and send replies using MessageUtil.createReplyMessage() and a dedicated reply producer.
example/src/main/java/org/apache/rocketmq/example/rpc · high confidence
Added sample configuration for multi-broker container deployment
The distribution now includes a sample configuration set in the \2container-2m-2s\ directory to support deploying multiple brokers within containerized environments. This set provides configuration files for two NameServers and two BrokerContainers, each hosting a master-slave pair (broker-a and broker-b). These files demonstrate how to configure broker roles, storage paths, and container-specific settings like \brokerConfigPaths\ to enable high-availability setups using the BrokerContainer feature.
distribution/conf/container · high confidence
Benchmark tools now support ACL authentication and configurable message compression
The benchmark examples (Producer, Consumer, BatchProducer, TransactionProducer) now include built-in support for RocketMQ Access Control List (ACL) authentication via the new AclClient helper, allowing users to enable ACL by passing the -a flag and providing access/secret keys. Additionally, the producer benchmarks now support configurable message compression (ZLIB, LZ4, SNAPPY) with adjustable levels and size thresholds via command-line options, enabling users to benchmark performance under compressed message scenarios.
example/src/main/java/org/apache/rocketmq/example/benchmark · high confidence
Broker configuration storage migrates to RocksDB with LMQ support
The broker now stores topic, consumer offset, and subscription group configurations in RocksDB instead of JSON files, introducing new managers (RocksDBTopicConfigManager, RocksDBConsumerOffsetManager, RocksDBSubscriptionGroupManager) that handle loading, persisting, and merging data with version checks. This change adds support for LMQ (Local Message Queue) topics and subscription groups within the RocksDB-backed config system, allowing these entities to be managed alongside standard configurations without requiring a migration to RocksDB for the consume queue store.
broker/src/main/java/org/apache/rocketmq/broker/config/v1 · high confidence
Container now tracks client channel lifecycle events
The BrokerContainer now monitors client connection states (connect, close, exception, idle, and active) via the new ContainerClientHouseKeepingService. This ensures that channel statistics are correctly updated for both master and slave brokers within the container, providing better visibility into client connectivity.
rocketmq-container · high confidence
Controller metrics support via OpenTelemetry
The controller now exposes operational metrics using OpenTelemetry, including gauges for disk usage and active broker count, counters for request and DLedger operations, and histograms for request and operation latency. Metrics are exported via OTLP gRPC, Prometheus HTTP server, or JSON logging, with configuration for cardinality limits and aggregation strategies.
controller/src/main/java/org/apache/rocketmq/controller/metrics · high confidence
Initial Bazel build support for RocketMQ
This change introduces the foundational Bazel build configuration files (\bazel/BUILD.bazel\, \bazel/GenTestRules.bzl\, and \bazel/rocketmq\_apis.BUILD\) to enable building and testing RocketMQ using the Bazel build system. It includes a helper macro for generating Java test rules and a minimal build file to expose the \rocketmq\_apis\ submodule's protocol buffer definitions as an external repository, allowing the project to be built with Bazel alongside existing build methods.
bazel · high confidence
Initial SQL92 message filtering support
This change introduces the core expression evaluation engine for SQL92-based message filtering in RocketMQ. It adds the \FilterFactory\ and \FilterSpi\ interface to manage filter registration, along with the \SqlFilter\ implementation that delegates to the new parser. The \expression\ package provides the full set of expression classes (binary, unary, logic, comparison, and constant) required to parse and evaluate SQL92 filter conditions against message properties.
filter/src/main/java/org/apache/rocketmq/filter/expression · high confidence
Initial implementation of BrokerOuterAPI for cluster communication
The broker now includes a new \BrokerOuterAPI\ class that handles all outbound network communication with NameServers and Controllers. This component manages broker registration, heartbeat reporting, route table retrieval, and controller-specific operations such as master election and sync state set updates, replacing previous ad-hoc or distributed implementations with a centralized, configurable client.
broker/src/main/java/org/apache/rocketmq/broker/out · high confidence
Initial implementation of OpenMessaging 0.1.0-alpha push and pull consumers
This change introduces the core consumer implementations for the OpenMessaging (OMS) bridge, adding \PushConsumerImpl\ and \PullConsumerImpl\ along with a new \ClientConfig\ class to manage RocketMQ-specific settings (such as consumer groups, thread pools, and redelivery limits). The consumers are configured to identify themselves with the OMS language code and support attaching/detaching queues, with push consumers registering message listeners and pull consumers utilizing a scheduled pull service with a local message cache. A specific behavioral detail is that the nameserver address is only set directly from access points when the \OMS\_RMQ\_DIRECT\_NAME\_SRV\ environment variable is enabled.
openmessaging/src/main/java/io/openmessaging/rocketmq/consumer · high confidence
Initial implementation of the OpenMessaging 0.3.0 specification
This change introduces the core implementation for the OpenMessaging (OMS) specification version 0.3.0 within the RocketMQ connector. It adds the \MessagingAccessPointImpl\ class, which serves as the entry point for creating OMS-compliant producers, push consumers, and pull consumers, while explicitly noting that the \StreamingConsumer\ and \ResourceManager\ features are not yet supported. Additionally, it includes the \LocalMessageCache\ component to handle message consumption logic, including offset management, message polling, and acknowledgment processing against the underlying RocketMQ broker.
rocketmq-openmessaging · high confidence
Initial project scaffolding and build configuration
The repository has been initialized with essential project scaffolding files, including the Apache License 2.0, NOTICE, and BUILDING documentation. A Bazel build system has been introduced (via WORKSPACE, MODULE.bazel, and .bazelrc) to manage dependencies and build targets, alongside a .gitignore and .licenserc.yaml for development hygiene. The README has been updated to provide a comprehensive quick-start guide for running RocketMQ locally, in Docker, and on Kubernetes, along with links to the broader ecosystem.
(repo-wide) · high confidence
Introduce ACL 2.0 authorization context and evaluation framework
This change adds the core authorization context model and evaluation engine for the new ACL 2.0 system. It introduces \AuthorizationContext\ and \DefaultAuthorizationContext\ to represent the subject, resource, and actions for an access request, along with model classes like \Acl\, \Policy\, and \PolicyEntry\ to define access rules. The \AuthorizationEvaluator\ now uses these contexts to enforce policies via \AuthorizationStrategy\, while \AuthorizationCompatibility\ handles legacy protocol matching for specific request types like heartbeats and unregisters. This provides the foundational logic for evaluating access decisions in the new authorization model.
auth/src/main/java/org/apache/rocketmq/auth/authorization · high confidence
Introduce Bazel build support for the tools module
The \tools\ module now includes a \BUILD.bazel\ file, enabling builds via the Bazel build system. This configuration defines the \tools\ and \tests\ Java libraries, specifying their source files and dependencies on internal modules (such as \remoting\, \client\, \common\, and \srvutil\) and external Maven artifacts (including Netty, Commons CLI, RocksDB, and Fastjson2).
tools · high confidence
Introduce ControllerRequestProcessor for centralized request handling
The controller module now uses a dedicated ControllerRequestProcessor to handle incoming network requests. This processor implements NettyRequestProcessor and routes commands such as broker registration, heartbeat, master election, and configuration updates to the ControllerManager. It also adds OpenTelemetry-based metrics collection for request latency and success/failure status, providing better observability for controller operations.
controller/src/main/java/org/apache/rocketmq/controller/processor · high confidence
Introduce LMQ subscription group support
The broker now includes a dedicated LmqSubscriptionGroupManager that extends the standard SubscriptionGroupManager to handle Light Message Queue (LMQ) groups. This change allows the broker to recognize, create, and manage subscription groups specifically identified as LMQs, ensuring they are handled separately from regular consumer groups during lookup and update operations.
broker/src/main/java/org/apache/rocketmq/broker/subscription · high confidence
Introduce Netty event tracking and request code distribution metrics
The remoting layer now exposes new observability capabilities for the Netty transport. A new NettyEvent class and NettyEventType.ACTIVE event allow external components to subscribe to and react to specific Netty channel lifecycle events via a dedicated executor. Additionally, a RemotingCodeDistributionHandler Netty channel handler has been added to track inbound and outbound request/response code distributions using LongAdders, providing a snapshot mechanism for metrics collection.
rocketmq-remoting · high confidence
Introduce ReplicasManager for controller-mode broker coordination
The broker now includes a new ReplicasManager component that handles registration with the controller, synchronizes metadata, and manages role transitions (master/slave) in controller mode. This component ensures reliable broker registration, handles controller address discovery, and coordinates high-availability state changes, replacing previous ad-hoc logic with a dedicated manager for replica coordination.
broker/src/main/java/org/apache/rocketmq/broker/controller · high confidence
Introduce RocketMQ BrokerContainer for hosting multiple brokers
Adds a new BrokerContainer module that allows a single process to host and manage multiple broker instances (masters, slaves, and DLedger brokers) via a shared remoting server and outer API. This introduces a new startup entry point (BrokerContainerStartup) and configuration model (BrokerContainerConfig) for deploying multiple brokers in a single JVM, along with lifecycle hooks (BrokerBootHook) for pre- and post-start customization.
container/src/main · high confidence
Introduce RocksDB-backed storage and locking for Pop consumption
The broker's Pop consumption subsystem now uses RocksDB to persist consumer state, replacing previous in-memory or alternative storage mechanisms. This change introduces a new PopConsumerContext to track pop operations, a PopConsumerLockService for distributed locking, and a PopConsumerRocksdbStore implementation that manages consumer records with configurable block cache and write buffer sizes. Users benefit from improved durability and performance of pop consumption operations, with the ability to tune RocksDB parameters through MessageStoreConfig.
broker/src/main/java/org/apache/rocketmq/broker/pop · high confidence
Introduce Tiered Storage module for offloading message data
This change introduces the Tiered Storage module, a technical preview feature that allows RocketMQ brokers to offload message data from local disk to cheaper, larger storage mediums (such as POSIX file systems or object stores like S3). By enabling this via the \messageStorePlugIn\ configuration, users can extend message retention times at a lower cost and flexibly specify different TTLs per topic. The module includes a new \TieredMessageStore\ plugin, configuration options for backend providers and read-ahead caching, and specific metrics for monitoring dispatch and storage operations.
tieredstore · high confidence
Introduce centralized JSON-based proxy configuration system
The proxy now uses a new configuration framework centered on \ConfigurationManager\ and \Configuration\ classes to load settings from a JSON file (defaulting to \rmq-proxy.json\). This system supports locating the config file via the \RMQ\_PROXY\_HOME\ environment variable or the \com.rocketmq.proxy.configPath\ system property, and it unifies the management of both \ProxyConfig\ and \AuthConfig\. Additionally, a new \MetricCollectorMode\ enum is introduced to allow users to explicitly configure whether metrics are collected by the proxy itself, sent to an external address, or disabled entirely.
proxy/src/main/java/org/apache/rocketmq/proxy/config · high confidence
Introduce dedicated long-polling services for Lite and LMQ consumption modes
The broker now includes new long-polling infrastructure classes to support Lite consumption and Local Message Queue (LMQ) scenarios. PopLiteLongPollingService handles long-polling for Lite consumers using a clientId-based key and a lightweight in-memory map, while LmqPullRequestHoldService extends the standard PullRequestHoldService to manage hold requests specifically for LMQ topics, including logic to clean up empty polling entries. Supporting classes such as PollingHeader, PollingResult, NotificationRequest, ManyPullRequest, and PullRequest provide the necessary data structures and result enums for these new polling flows.
broker/src/main/java/org/apache/rocketmq/broker/longpolling · high confidence
Introduce expression-based message filtering with bloom filter optimization
The broker now supports advanced message filtering using SQL92 expressions (in addition to existing tag filtering). This change introduces a new filtering architecture in the broker, including \ConsumerFilterManager\ to manage filter data, \ConsumerFilterData\ to store compiled expressions and metadata, and \ExpressionMessageFilter\ to evaluate messages against these expressions. A key performance enhancement is the integration of Bloom filters to quickly determine if a message might match a consumer's filter before full evaluation, reducing unnecessary processing. The system also handles retry topics correctly by resolving the original topic for filter evaluation.
broker/src/main/java/org/apache/rocketmq/broker/filter · high confidence
Introduce jRaft-based controller implementation
The controller module now includes a new implementation of the controller logic using the jRaft consensus library (JRaftControllerStateMachine and RaftReplicasInfoManager). This adds an alternative to the existing DLedger-based controller, enabling users to run the controller with jRaft for leader election and state replication. The change introduces jRaft-specific state machine handling, broker heartbeat and liveness tracking, and request processing for controller operations such as broker registration, master election, and sync state set management.
rocketmq-controller · high confidence
Introduce local RocksDB-backed authentication and authorization metadata providers
The auth module now includes new local metadata providers that persist user and ACL data to RocksDB, replacing the previous reliance on external or default providers. This change adds a \LocalAuthenticationMetadataProvider\ and \LocalAuthorizationMetadataProvider\ that store user profiles and access control lists in RocksDB column families, backed by Caffeine in-memory caches for performance. It also introduces the corresponding manager implementations (\AuthenticationMetadataManagerImpl\ and \AuthorizationMetadataManagerImpl\) and authentication/authorization chain handlers (\DefaultAuthenticationHandler\, \AclAuthorizationHandler\, \UserAuthorizationHandler\) that use these providers to validate user status, check signatures, and enforce policy decisions. This provides a self-contained, persistent storage backend for the new ACL 2.0 authentication and authorization system.
rocketmq-auth · high confidence
Introduce new proxy startup, configuration, and admin gRPC interfaces
The proxy now features a new startup entry point (ProxyStartup) that accepts command-line arguments for configuration paths and proxy mode (local or cluster), and initializes both a data-plane gRPC server and a dedicated admin gRPC server for management operations. This includes new model converters for admin responses and a remoting protocol handler interface, enabling users to start the proxy with explicit configuration and interact with its new administrative capabilities.
proxy/src/main/java/org/apache/rocketmq/proxy · high confidence
Introduce proxy common infrastructure for message renewal and context management
The proxy now includes a new set of common classes to support message receipt handling and renewal strategies. This includes a \ProxyContext\ for managing request-scoped data (such as client ID, namespace, and action), a \MessageReceiptHandle\ to track message consumption state and renewal attempts, and a \RenewStrategyPolicy\ that defines the back-off intervals for message visibility renewal. Additionally, supporting classes like \Address\ for host/port parsing, \ContextVariable\ constants, and \RenewEvent\/\ReceiptHandleGroupKey\ for event tracking have been added to facilitate these operations.
proxy/src/main/java/org/apache/rocketmq/proxy/common · high confidence
Introduce srvutil module with server utilities and file watching
The new srvutil module provides shared server utilities, including command-line argument parsing (help and namesrvAddr options) via ServerUtil, a standardized shutdown hook mechanism via ShutdownHookThread, and a FileWatchService for monitoring file changes. The module is built with Bazel and includes unit tests for these components.
srvutil · high confidence
Introduce standalone Controller module for broker management
Adds a new standalone Controller component that manages broker registration, heartbeat monitoring, and master election. The module introduces a \ControllerManager\ to initialize and run the service, a \Controller\ interface defining the broker management API (register, elect, sync state), and a \BrokerHeartbeatManager\ to track broker liveness. It supports two backend implementations, DLedger and jRaft, and includes a startup class (\ControllerStartup\) for launching the server and an \ElectPolicy\ interface to allow configurable master election strategies.
controller/src/main/java/org/apache/rocketmq/controller · high confidence
Introduce transactional message metrics tracking and flushing
The broker now tracks and persists metrics for transactional messages. A new \TransactionMetrics\ class stores per-topic counts, and a \TransactionMetricsFlushService\ background thread periodically persists these metrics to disk, ensuring that transactional message activity is observable and survives broker restarts.
broker/src/main/java/org/apache/rocketmq/broker/transaction · high confidence
Introduces AbortProcessException to enforce immediate hook termination
The \common\ module now includes the \AbortProcessException\ class, a new runtime exception designed specifically for broker hooks (SendMessageHook, ConsumeMessageHook, and RPCHook). When thrown by a hook implementation, this exception forces the broker to immediately return an error response to the client with a specific response code, bypassing normal processing flow. This provides a standardized mechanism for hooks to enforce strict validation or policy decisions that require immediate request rejection.
common/src/main · high confidence
Introduces EscapeBridge for remote message escaping in Slave Acting Master mode
The new EscapeBridge component in the broker's failover package enables the asynchronous escaping of messages to a remote broker when the local broker is acting as a master for a slave. This change supports the Slave Acting Master feature by providing mechanisms to send messages (including transactional half-messages) to remote brokers via an async executor, ensuring message durability and availability during failover scenarios.
broker/src/main/java/org/apache/rocketmq/broker/failover · high confidence
Introduces RocksDB-backed configuration storage and new cold data flow control
The broker now supports a new configuration management layer (ConfigManagerV2) that persists topic, subscription group, and consumer offset metadata to RocksDB, providing a high-performance alternative to the previous file-based storage. Additionally, a new cold data flow control mechanism is introduced, featuring a service that tracks consumer group read activity and applies adaptive or simple throttling strategies to prevent cold storage bottlenecks, along with a hold service to manage pull request suspension for cold data reads.
rocketmq-broker · high confidence
Introduces broker-side infrastructure for the new Lite Topic message model
This change adds the core broker components required to support the new Lite Topic message model. It introduces lifecycle managers (AbstractLiteLifecycleManager, LiteLifecycleManager, and RocksDBLiteLifecycleManager) to handle topic TTL, subscription validity, and offset tracking for both file-based and RocksDB-backed message stores. A new LmqPrefixIndex (backed by a PatriciaTrie) accelerates wildcard dispatch by indexing LMQ names. Exclusive subscription eviction is now enforced via server-side tombstones managed by ExclusiveEvictionTombstones, and metadata utilities (LiteMetadataUtil) provide configuration checks for Lite-specific group and topic behaviors.
broker/src/main/java/org/apache/rocketmq/broker/lite · high confidence
Introduces default broker election policy for controller mode
Adds the DefaultElectPolicy implementation, which defines the logic for selecting a new master broker when the current one fails. The policy prioritizes the existing master if it remains valid, respects a preferred broker ID if specified, and otherwise selects the best candidate from the sync state set (or all replicas if necessary) based on a ranking of epoch, max offset, and election priority.
controller/src/main/java/org/apache/rocketmq/controller/elect/impl · high confidence
Introduces gRPC v2 common utilities for message conversion and validation
The proxy now includes a new set of common classes in the \proxy.grpc.v2.common\ package to support the gRPC v2 protocol. \GrpcConverter\ handles the translation of internal message structures to gRPC format, including system properties and user attributes. \GrpcValidator\ enforces rules for topics, consumer groups, tags, and lite topics, rejecting invalid inputs with specific gRPC error codes. \ResponseBuilder\ and \ResponseWriter\ manage the mapping of internal exceptions to gRPC status codes and the safe writing of responses to clients, ensuring proper error handling and cancellation checks.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common · high confidence
Introduces new ACL 2.0 authentication evaluation and context components
This change adds the core authentication evaluation logic and context structures for the new ACL 2.0 system. It introduces the \AuthenticationEvaluator\ to coordinate authentication checks, along with \AuthenticationContext\ and \DefaultAuthenticationContext\ to hold request details like credentials and signatures. The \DefaultAuthenticationContextBuilder\ is added to parse these credentials from both gRPC metadata and Remoting commands. Additionally, the \AuthenticationFactory\ is introduced to manage the lifecycle of providers and strategies, while \StatelessAuthenticationStrategy\ and \StatefulAuthenticationStrategy\ provide the actual evaluation logic, with the latter supporting caching via Caffeine.
auth/src/main/java/org/apache/rocketmq/auth/authentication · high confidence
Introduces new client-side rebalance strategies and admin API
The client now supports two new message queue allocation strategies: a consistent hash-based strategy (AllocateMessageQueueConsistentHash) and a machine-room proximity strategy (AllocateMachineRoomNearby), allowing consumers to prefer local brokers or use deterministic hashing for rebalancing. Additionally, a new MqClientAdmin interface and implementation are provided, exposing asynchronous, future-based administrative operations such as querying message stats, topic statistics, and managing topics and subscription groups directly from the client.
rocketmq-client · high confidence
Introduction of ACL 2.0 configuration and broker-side conversion logic
This change introduces the foundational configuration and data-conversion components for the new ACL 2.0 authentication and authorization system. The \AuthConfig\ class in the auth module now provides a structured configuration model for enabling authentication and authorization, specifying providers, strategies, and cache settings, while also supporting migration from the previous V1 ACL format. In the broker module, new \AclConverter\ and \UserConverter\ classes handle the translation between the internal ACL 2.0 domain models (Subjects, Policies, Resources) and the legacy \AclInfo\/\UserInfo\ protocol structures, ensuring that the broker can process and expose the new security model via existing remoting protocols.
auth/src/main/java/org/apache/rocketmq/auth/config, broker/src/main/java/org/apache/rocketmq/broker/auth · high confidence
Introduction of Broker Plugin and Processor Architecture
The broker now supports an extensible plugin architecture and a refactored request processing model. A new \BrokerAttachedPlugin\ interface allows external modules to hook into the broker lifecycle (load, start, shutdown, metadata sync, and status changes). Additionally, core request handling has been reorganized into dedicated processor classes (such as \AckMessageProcessor\, \PeekMessageProcessor\, and \AbstractSendMessageProcessor\) and a \PullMessageResultHandler\ interface, replacing the previous monolithic processing logic to improve modularity and maintainability.
broker/src/main/java/org/apache/rocketmq/broker/processor · high confidence
Introduction of RebalanceLockManager for message queue locking
The broker now includes a new RebalanceLockManager component in the rebalance package to manage distributed locks for message queues. This manager handles lock acquisition, expiration checks, and batch locking operations for consumer groups, ensuring that message queues are correctly assigned and locked during the rebalancing process.
broker/src/main/java/org/apache/rocketmq/broker/client/rebalance · high confidence
Introduction of RocketMQ ACL 2.0 client support
The client module now includes the foundational classes for the new ACL 2.0 security model. This change introduces the \AclSigner\ for calculating request signatures, \AclUtils\ for handling network address validation (including IPv6), and \SessionCredentials\ for managing access keys and secret keys. It also adds the \AccessChannel\ configuration to distinguish between local and cloud deployments, and provides the \MQAdmin\ and \MqClientAdmin\ interfaces to expose management capabilities such as topic creation, offset querying, and cluster statistics.
client/src/main/java · high confidence
Introduction of static topic queue mapping and LMQ support
The broker now supports static topic queue mapping and Light Message Queues (LMQ). A new TopicQueueMappingManager and TopicQueueMappingCleanService manage static topic queue mappings, allowing for persistent, versioned queue assignments that are cleaned up based on broker load and epoch changes. Additionally, a new LmqTopicConfigManager extends TopicConfigManager to handle LMQ-specific topic configurations, ensuring that LMQ topics are treated as existing and selected appropriately without requiring full configuration updates.
broker/src/main/java/org/apache/rocketmq/broker/topic · high confidence
Multi-protocol remoting server supports HTTP/2 alongside Remoting
The proxy now includes a new MultiProtocolRemotingServer that allows both the standard RocketMQ Remoting protocol and HTTP/2 to be served on the same port. This is enabled by a new MultiProtocolTlsHelper that configures SSL/TLS with ALPN to negotiate HTTP/2, and a ProtocolNegotiationHandler that routes connections to the appropriate protocol handler. A new RemotingProxyOutClient interface is also introduced to support client-side invocation.
proxy/src/main/java/org/apache/rocketmq/proxy/remoting · high confidence
Namespace v2 example code added
The namespace example directory now includes three new Java files demonstrating how to use the namespace v2 feature: ProducerWithNamespace, PullConsumerWithNamespace, and PushConsumerWithNamespace. These examples show how to configure DefaultMQProducer, DefaultMQPullConsumer, and DefaultMQPushConsumer with the setNamespaceV2 method to operate within a specific namespace instance.
example/src/main/java/org/apache/rocketmq/example/namespace · high confidence
Namesrv introduces KVConfigManager for namespace-scoped key-value configuration
The NameServer now supports storing and retrieving key-value configuration data organized by namespace. A new KVConfigManager handles loading, persisting, and managing these configurations in memory and on disk, while KVConfigSerializeWrapper provides the serialization structure for the config table. This enables users to manage specific configuration items within isolated namespaces via the NameServer.
namesrv/src/main/java/org/apache/rocketmq/namesrv/kvconfig · high confidence
Nearby routing support via zone-based filtering in NameServer
The NameServer now supports a 'nearby routing' feature that filters topic route data based on the consumer's availability zone. A new ZoneRouteRPCHook intercepts GET\_ROUTEINFO\_BY\_TOPIC responses when zone mode is enabled, removing broker data and associated queue data that do not match the requested zone name, while preserving masters and slaves in the same zone. This allows clients to receive route information containing only brokers within their specific zone, optimizing network latency for zone-aware deployments.
namesrv/src/main/java/org/apache/rocketmq/namesrv/route · high confidence
New AdminService interface for proxy management operations
Added the AdminService interface in the proxy module, defining an asynchronous API for cluster administration tasks such as retrieving broker cluster info, topic routes, consumer connections, and message statistics, as well as mutating topic and subscription group configurations. This interface serves as the backend contract for the proxy's admin capabilities, ensuring all operations are delegated through the proxy's managed broker client rather than opening direct broker links.
proxy/src/main/java/org/apache/rocketmq/proxy/service/admin · high confidence
New Broker2Client component for managing client communication
The broker now includes a new \Broker2Client\ class in the \org.apache.rocketmq.broker.client.net\ package to handle interactions with connected clients. This component centralizes logic for notifying consumers of subscription changes, checking transaction states for producers, and resetting consumer offsets based on time or force flags. It also supports specific serialization formats for C++ clients during offset resets, providing a structured way for the broker to send commands to client channels.
broker/src/main/java/org/apache/rocketmq/broker/client/net · high confidence
New BrokerFastFailure component for proactive queue cleanup
A new BrokerFastFailure class has been added to the broker latency package. It introduces a scheduled background task that proactively cleans expired requests from various thread-pool queues (send, pull, lite-pull, heartbeat, transaction, ack, and admin broker). When a request exceeds its configured wait time in a queue, it is removed and the client receives a SYSTEM\_BUSY response, helping to prevent thread starvation and improve broker responsiveness under load.
broker/src/main/java/org/apache/rocketmq/broker/latency · high confidence
New BrokerHousekeepingService for NameServer channel lifecycle management
The NameServer now includes a dedicated BrokerHousekeepingService that implements the ChannelEventListener interface to manage broker connection states. This service ensures that when a broker channel is closed, encounters an exception, or becomes idle, the NameServer's RouteInfoManager is notified to destroy the associated channel data, thereby keeping the routing table consistent with active broker connections.
rocketmq-namesrv · high confidence
New DLedger role change handler for broker state transitions
Added DLedgerRoleChangeHandler to manage broker role transitions (leader, follower, candidate) within the DLedger replication subsystem. This component ensures that when a broker's role changes, it correctly updates its internal state, registers with the cluster, and manages slave synchronization tasks, thereby supporting high-availability scenarios such as slave acting as master.
broker/src/main/java/org/apache/rocketmq/broker/dledger · high confidence
New JRaft-based controller implementation
The controller module now includes a new implementation of the Controller interface backed by the JRaft consensus library (JRaftController). This change introduces the core controller logic, a state machine, and a closure mechanism for handling asynchronous RPC responses, along with a set of request and response task classes (such as BrokerCloseChannel, GetBrokerLiveInfo, and RaftBrokerHeartBeatEvent) required to manage broker lifecycle and heartbeat events via the new Raft-based architecture.
controller/src/main/java/org/apache/rocketmq/controller/impl, controller/src/main/java/org/apache/rocketmq/controller/impl/closure, controller/src/main/java/org/apache/rocketmq/controller/impl/task · high confidence
New LMQ example applications demonstrating multi-dispatch and consumption patterns
Added four new example programs in the LMQ (Local Message Queue) package to demonstrate how to use RocketMQ's multi-dispatch feature. LMQProducer shows how to send messages to a parent topic while dispatching copies to multiple LMQ sub-topics using the INNER\_MULTI\_DISPATCH property. LMQPushConsumer, LMQPullConsumer, and LMQPushPopConsumer provide examples of consuming these dispatched messages using push, pull, and pop-consume modes respectively, including the necessary client-side route compensation logic required for LMQ topics.
example/src/main/java/org/apache/rocketmq/example/lmq · high confidence
New OpenMessaging producer implementation for RocketMQ
Added AbstractOMSProducer and ProducerImpl classes to implement the OpenMessaging specification for RocketMQ producers. The implementation supports synchronous, asynchronous, and oneway message sending, with specific handling for BytesMessage types. Configuration includes conditional name server address resolution via the OMS\_RMQ\_DIRECT\_NAME\_SRV environment variable, and client language is tagged as 'OMS'.
openmessaging/src/main/java/io/openmessaging/rocketmq/producer · high confidence
New OpenTelemetry-based broker metrics framework
The broker now exposes a comprehensive set of observability metrics using the OpenTelemetry standard, replacing or augmenting previous internal metrics. This change introduces a new metrics framework in the broker module, including constants, managers, and calculators for tracking broker health (topic/consumer group counts, permissions), message flow (throughput, message sizes, commit/rollback counts), client connections, and detailed consumer lag and latency. It also adds specific metrics for the Pop consumption model (buffer sizes, revive lag/latency) and remoting stats. The implementation supports multiple export backends, including Prometheus, OTLP gRPC, and JSON logging, and includes optimizations like pre-built AttributeKey instances to reduce allocation overhead.
broker/src/main/java/org/apache/rocketmq/broker/metrics · high confidence
New Proxy Admin gRPC management interface with ACL 2.0 enforcement
The proxy now exposes a dedicated gRPC admin service (ProxyAdminGrpcService) for cluster management operations such as retrieving runtime stats, topic routes, consumer connections, and subscription details, as well as executing administrative actions like resetting group offsets or sending messages. This service is secured by a new ProxyAdminAuthInterceptor that enforces ACL 2.0 policies using specific proxy.admin.\* resources, ensuring that administrative actions are properly authenticated and authorized. The implementation includes a ProxyAdminForwarder to route requests to the correct proxy in a cluster environment and supports TLS certificate reloading for the admin server.
rocketmq-proxy · high confidence
New ReceiptHandleManager interface and DefaultReceiptHandleManager implementation
The proxy now introduces a dedicated \ReceiptHandleManager\ interface and its \DefaultReceiptHandleManager\ implementation to centralize the management of message receipt handles. This component maintains a concurrent map of receipt handle groups keyed by channel and consumer group, providing methods to add, remove, and query unacked message counts. It includes a scheduled task that scans for messages nearing their visibility timeout and submits renewal tasks to a dedicated worker thread pool, while also listening for client unregistration events to automatically clear stale handle data.
proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt · high confidence
New SlaveSynchronize component for master-slave data replication
A new SlaveSynchronize class has been added to the broker module to handle the synchronization of configuration and state data from the master broker to the slave. This component implements specific sync methods for topic configurations, consumer offsets, delay offsets, subscription group configurations, and message request modes, ensuring that the slave broker maintains consistency with the master's state. It also includes logic to sync timer metrics when the timer wheel is enabled, providing a centralized mechanism for maintaining high availability through data replication.
broker/src/main/java/org/apache/rocketmq/broker/slave · high confidence
New benchmark scripts for batch and transaction producers
The distribution/benchmark directory now includes dedicated startup scripts for batch and transaction producers (batchproducer.sh and tproducer.sh), alongside a centralized shutdown.sh script that supports stopping producer, consumer, transaction producer, and batch producer instances. These scripts leverage an updated runclass.sh utility that handles JVM configuration, GC logging (including RAM disk support on macOS and /dev/shm on Linux), and classpath setup, providing a more structured way to run and manage benchmark workloads.
distribution/benchmark · high confidence
New broker logging configuration and transaction metadata schema
The broker now ships with a dedicated \rmq.broker.logback.xml\ configuration file that defines structured, asynchronous logging for specific subsystems (default, broker, protection, watermark, RocksDB, and store) with configurable log directories and file rotation policies. Additionally, a new \transaction.sql\ schema file is included to define the \t\_transaction\ table structure for storing transaction metadata.
broker/src/main/resources · high confidence
New broker utility classes for message validation and atomic counting
Added two new utility classes to the broker module: HookUtils, which centralizes pre-send message validation logic (including checks for topic length, batch consistency, timer message handling, and LMQ quota enforcement), and PositiveAtomicCounter, a thread-safe counter that ensures values remain positive by masking the sign bit. These utilities support the broker's message processing pipeline by providing reusable, standardized checks and counters.
broker/src/main/java/org/apache/rocketmq/broker/util · high confidence
New broker-side message tracing context and hook interfaces
The broker now exposes dedicated context objects and hook interfaces for tracing message operations. New classes \ConsumeMessageContext\ and \SendMessageContext\ capture detailed metadata for consumption and sending events respectively, including fields for namespace, account statistics, commercial usage, and message identifiers. Corresponding interfaces \ConsumeMessageHook\ and \SendMessageHook\ allow external components to intercept and observe these events before and after message processing, enabling enhanced monitoring and auditing capabilities for message flow within the broker.
broker/src/main/java/org/apache/rocketmq/broker/mqtrace · high confidence
New channel abstraction and serialization interfaces
The proxy module introduces a new set of interfaces and utilities in the \proxy.processor.channel\ package to standardize channel handling. This includes the \ChannelProtocolType\ enum to identify protocol variants (gRPC v1/v2, remoting), and interfaces \ChannelExtendAttributeGetter\ and \RemoteChannelConverter\ to abstract channel attributes and conversion logic. Additionally, \RemoteChannelSerializer\ provides JSON serialization and deserialization for \RemoteChannel\ objects using Fastjson2, enabling consistent handling of remote channel metadata across the proxy.
proxy/src/main/java/org/apache/rocketmq/proxy/processor/channel · high confidence
New common infrastructure for lifecycle management, typed attributes, and logging configuration
The rocketmq-common module introduces several foundational classes to support new message models and operational observability. It adds a LifecycleAwareServiceThread base class to provide reliable thread-start synchronization, and a ThreadFactoryImpl that supports broker-container identification and logs uncaught thread exceptions. A new attribute framework (BooleanAttribute, EnumAttribute, LongRangeAttribute, StringAttribute) is included to validate topic and message properties. Additionally, a DefaultJoranConfiguratorExt is added to manage Logback configuration discovery for various components, and a ThreadPoolQueueSizeMonitor is provided to track thread-pool queue saturation.
rocketmq-common · high confidence
New controller event model and serialization infrastructure
The controller module now uses a new internal event system to manage broker state changes. This introduces a set of event types (AlterSyncStateSet, ApplyBrokerId, ElectMaster, CleanBrokerData, UpdateBrokerAddress) and corresponding message classes, along with serializers to handle their serialization and deserialization. This change provides the underlying mechanism for the controller to track and propagate broker lifecycle events, such as address updates and data cleanup, within the new jRaft-based controller implementation.
controller/src/main/java/org/apache/rocketmq/controller/impl/event · high confidence
New distribution/bin scripts for RocketMQ components and utilities
This change introduces a comprehensive set of new shell and batch scripts in the distribution/bin directory to manage RocketMQ services. It adds dedicated launchers for the Broker (mqbroker), NameServer (mqnamesrv), Controller (mqcontroller), Proxy (mqproxy), and BrokerContainer (mqbrokercontainer), along with their Windows counterparts. The Broker script now supports an --enable-proxy flag to switch between standard broker and proxy modes. New utility scripts include mqadmin for administrative tasks, mqshutdown for stopping services, and os.sh for system tuning. Additionally, helper scripts like play.sh/play.cmd for quick starts, cachedog.sh/cleancache.sh for memory management, and export.sh for configuration export are provided.
distribution/bin · high confidence
New examples for scheduled and timer message delivery
Added four new example classes in the schedule package to demonstrate message scheduling capabilities: ScheduledMessageProducer and ScheduledMessageConsumer show how to send and receive messages with a fixed delay using delay time levels, while TimerMessageProducer and TimerMessageConsumer illustrate timer-based delivery for RocketMQ 5.0+, supporting arbitrary time delays via specific delivery timestamps or millisecond offsets.
example/src/main/java/org/apache/rocketmq/example/schedule · high confidence
New gRPC pipeline and authentication metadata providers in the proxy
The proxy now includes new infrastructure for handling gRPC requests and authentication metadata. Specifically, it adds \ProxyAuthenticationMetadataProvider\ and \ProxyAuthorizationMetadataProvider\ to bridge proxy-level metadata with the new RocketMQ ACL 2.0 authentication and authorization models. Additionally, a new \ContextInitPipeline\ and \RequestPipeline\ interface are introduced to initialize gRPC context (such as client ID, protocol type, and namespace) from incoming request headers, enabling more structured request processing within the proxy.
proxy/src/main/java/org/apache/rocketmq/proxy/auth, proxy/src/main/java/org/apache/rocketmq/proxy/grpc/pipeline · high confidence
New gRPC v2 client activity handler for heartbeat and termination
The proxy now includes a new \ClientActivity\ class in the gRPC v2 client package to handle client lifecycle events. This component implements the \heartbeat\ method, which registers producers and consumers (including new Lite Push and Simple consumer types) based on client settings, and the \notifyClientTermination\ method, which unregisters these clients and cleans up channel resources when a client disconnects. This change introduces the core logic for managing client connections and subscriptions in the v2 gRPC protocol implementation.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client · high confidence
New gRPC v2 route handling for topic and assignment queries
The proxy now implements the gRPC v2 route activity, introducing new endpoints for querying topic routes and consumer group assignments. This change enables clients to retrieve message queue locations and broker assignments via the v2 protocol, including support for FIFO and lite topic modes in assignment logic, and ensures that permission checks are correctly mapped to the gRPC response format.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route · high confidence
New heartbeat management implementations for the controller
The controller now includes two new implementations for managing broker heartbeats: DefaultBrokerHeartbeatManager, which handles broker registration, liveness checks, and lifecycle notifications using local state, and RaftBrokerHeartBeatManager, which delegates heartbeat events to the JRaft controller for distributed consistency. Supporting data classes BrokerIdentityInfo and BrokerLiveInfo have been added to represent broker identity and live status respectively, enabling the controller to track broker epochs, offsets, and election priorities.
controller/src/main/java/org/apache/rocketmq/controller/impl/heartbeat · high confidence
New metadata service abstraction for proxy configuration and ACL caching
The proxy now introduces a \MetadataService\ interface with two implementations: \LocalMetadataService\ for direct broker metadata access and \ClusterMetadataService\ for cluster-wide metadata retrieval with caching. \ClusterMetadataService\ adds Guava-based caching for topic configurations, subscription group configs, users, and ACLs, with configurable cache sizes, expiration, and refresh intervals. This enables the proxy to efficiently cache and retrieve metadata across the cluster, improving performance for metadata lookups.
proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata · high confidence
New offset management components for Lite Topics and broadcast consumption
The broker now includes new classes in the offset package to support advanced consumption models. BroadcastOffsetStore provides in-memory offset tracking for broadcast consumption mode. LmqConsumerOffsetManager extends the standard offset manager to handle Light Message Queue (LMQ) specific offset storage and reset logic. MemoryConsumerOrderInfoManager introduces a lightweight, non-persistent manager for tracking consumer order information and queue suspension in Lite Topics, prioritizing low latency over persistence for these specific use cases.
broker/src/main/java/org/apache/rocketmq/broker/offset · high confidence
New operation-mode Producer and Consumer examples added
Added new example classes for the operation package that demonstrate how to send and receive messages using the DefaultMQProducer and DefaultMQPushConsumer. These examples include command-line argument parsing via Apache Commons CLI to configure group names, topics, tags, and message counts, allowing users to quickly test basic messaging scenarios.
example/src/main/java/org/apache/rocketmq/example/operation · high confidence
New pagecache transfer classes for message I/O
The broker's pagecache package now includes three new classes—ManyMessageTransfer, OneMessageTransfer, and QueryMessageTransfer—that implement Netty's FileRegion interface to handle zero-copy file transfers. These classes manage the transfer of message data (headers and payloads) from internal buffer results to network channels, supporting scenarios involving multiple messages, single messages, and query results.
broker/src/main/java/org/apache/rocketmq/broker/pagecache · high confidence
New proxy channel management and invocation context infrastructure
The proxy now introduces a dedicated channel management layer to handle client connections and internal invocations more robustly. A new ChannelManager centralizes the creation and lifecycle of SimpleChannel instances, using client identifiers to map and track active connections while automatically cleaning up inactive channels. Additionally, new InvocationContext and InvocationContextInterface classes provide a structured way to manage asynchronous remoting command responses, including expiration handling and response completion, replacing ad-hoc context management in the proxy's remoting processors.
proxy/src/main/java/org/apache/rocketmq/proxy/service/channel · high confidence
New proxy message service interface and local command support
The proxy now exposes a new \MessageService\ interface in the \proxy.service.message\ package, defining asynchronous methods for core messaging operations including sending messages (with support for multiple results), popping, pulling, acknowledging (including batch acks), updating consumer offsets, locking/unlocking message queues, and recalling delay messages. This interface is accompanied by a \LocalRemotingCommand\ helper for creating local remoting requests and a \ReceiptHandleMessage\ wrapper to bundle receipt handles with message IDs, providing a structured foundation for the proxy's message handling logic.
proxy/src/main/java/org/apache/rocketmq/proxy/service/message · high confidence
New proxy utility classes for tag filtering, gRPC metadata, and constants
Added three new utility classes to the proxy common module: FilterUtils provides a method to check if a message's tags match a consumer group's subscription data; GrpcUtils offers helpers to safely add or replace gRPC metadata headers and retrieve server call attributes; and ProxyUtils defines constants for maximum message numbers in pop requests and broker addresses.
proxy/src/main/java/org/apache/rocketmq/proxy/common/utils · high confidence
New remoting message conversion utility
A new RemotingConverter utility class has been added to the proxy remoting common package. This singleton component provides a method to serialize RocketMQ messages into bytes, specifically handling store size recalculation and validating topic length constraints during the encoding process.
proxy/src/main/java/org/apache/rocketmq/proxy/remoting/common · high confidence
New simple example clients for RocketMQ features
Added a suite of Java example programs in the simple examples package to demonstrate core RocketMQ capabilities. These include AclClient for authenticated access, AsyncProducer and OnewayProducer for different send modes, PopConsumer for the POP consumption protocol, and multiple LitePullConsumer variants (LitePullConsumerAssign, LitePullConsumerAssignWithSubExpression, LitePullConsumerSubscribe) for flexible pull-based consumption. The set also includes legacy and specialized consumers like PullConsumer, PullScheduleService, and PushConsumer, alongside helper classes like CachedQueue and RandomAsyncCommit to support these examples.
example/src/main/java/org/apache/rocketmq/example/simple · high confidence
New timer benchmark examples for timing messages
Added TimerConsumer and TimerProducer example programs in the benchmark timer module to demonstrate and measure performance for timing messages with arbitrary time delays. The producer sends messages with specific delivery timestamps using the timer deliver property, while the consumer tracks and reports metrics such as consume TPS, average/percentile delayed durations, and send latency.
example/src/main/java/org/apache/rocketmq/example/benchmark/timer · high confidence
New utility for identifying channel types and protocols
A new ChannelHelper utility class has been added to the proxy's common channel package. This component provides static methods to determine if a network channel is remote (synced from another proxy) and to identify the specific protocol type (gRPC v2, Remoting, or other remote types) associated with a connection, enabling more precise channel handling within the proxy infrastructure.
proxy/src/main/java/org/apache/rocketmq/proxy/common/channel · high confidence
Proxy metrics framework with OpenTelemetry support
The proxy now exposes operational metrics via a new metrics framework built on OpenTelemetry. A new ProxyMetricsManager initializes and manages metric exporters (OTLP gRPC, Prometheus, or logging) based on configuration, while ProxyMetricsConstant defines standard labels and gauges such as proxy status. This enables users to monitor proxy health and performance through standard observability backends.
proxy/src/main/java/org/apache/rocketmq/proxy/metrics · high confidence
Proxy now synchronizes consumer registration state across the cluster
The proxy introduces a new heartbeat synchronization mechanism to keep consumer group information consistent across multiple proxy instances. A new \ClusterConsumerManager\ intercepts consumer register and unregister events and delegates them to a \HeartbeatSyncer\. This syncer broadcasts client connection details (such as client ID, language, version, and subscription data) to a dedicated broadcast topic, allowing other proxies in the cluster to receive and apply these changes. This ensures that consumer routing and load balancing decisions remain accurate even when clients connect to different proxy nodes.
proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage · high confidence
SQL message filtering now supports CONTAINS, STARTSWITH, and ENDSWITH operators
The SQL92-based message selector parser in the filter module has been updated to recognize and process the CONTAINS, STARTSWITH, and ENDSWITH string operators. This change, implemented via the regenerated JavaCC parser files (SelectorParser.jj and its generated Java counterparts), allows users to write more expressive filter expressions that check for substring presence or string prefix/suffix matches, extending the previous capabilities which were limited to standard comparison and logical operators.
filter/src/main/java/org/apache/rocketmq/filter/parser · high confidence
Support for forwarding messages to dead-letter queues and recalling delay messages
The gRPC v2 proxy now exposes two new producer operations: forwarding messages to a dead-letter queue (DLQ) and recalling delay messages. Users can now send a message to the DLQ via the new ForwardMessageToDeadLetterQueue activity, which validates the topic and consumer group, retrieves the receipt handle, and forwards the request to the messaging processor. Additionally, the new RecallMessage activity allows users to recall a previously sent delay message by providing the topic and a recall handle, with the operation completing asynchronously and returning the recalled message ID.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer · high confidence
Tiered storage module core implementation and stability fixes
This release introduces the core implementation of the RocketMQ tiered storage module, including the MessageStoreDispatcher for handling message dispatch and memory backpressure, the MessageStoreFetcher with a configurable dual-TTL read-ahead cache, and the DefaultMetadataStore for managing topic and queue metadata. Alongside these new capabilities, the update includes numerous bug fixes addressing concurrency issues, resource leaks, commit failure reporting, and index service logic to ensure stable operation.
rocketmq-tiered-store · high confidence
gRPC v2 consumer support for message acknowledgment and visibility management
The proxy now exposes gRPC v2 consumer operations for acknowledging messages and modifying their visibility duration. AckMessageActivity handles both batch and individual message acknowledgments, returning detailed per-message results and a combined status (including MULTIPLE\_RESULTS when outcomes vary). ChangeInvisibleDurationActivity allows clients to extend or shorten the time a message remains invisible, supporting Lite Topic and consumer suspension flags. ReceiveMessageResponseStreamWriter manages the streaming response for message retrieval, handling various pop statuses (FOUND, POLLING\_FULL, NO\_NEW\_MSG) and automatically adjusting visibility on write errors.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer · high confidence
Architecture
Introduction of the Message Store submodule
The \rocketmq-store\ module has been introduced as a standalone submodule, containing the core message persistence and retrieval logic. This includes the \CommitLog\ for durable message storage, \DefaultMessageStore\ for coordinating store components, \IndexService\ for message indexing, and the \CompactionLog\ system for topic compaction. The diff also adds supporting infrastructure such as \MultiPathMappedFileQueue\ for multi-path storage support and \FlushDiskWatcher\ for managing asynchronous disk flush operations.
rocketmq-store · high confidence
Proxy processor layer refactored into dedicated service components
The proxy's messaging logic has been reorganized into a set of specialized processor classes—ProducerProcessor, ConsumerProcessor, TransactionProcessor, ReceiptHandleProcessor, ClientProcessor, and RequestBrokerProcessor—managed by a new DefaultMessagingProcessor. This change introduces a unified MessagingProcessor interface and an AbstractProcessor base class to standardize access to the ServiceManager and handle lifecycle management. For users, this represents a structural shift in how the proxy routes and processes messages, acknowledging, and transactions, laying the groundwork for improved modularity and support for new features like Lite Topic and v2 protocol handling.
proxy/src/main/java/org/apache/rocketmq/proxy/processor · high confidence
Proxy service layer refactored into pluggable local and cluster modes
The proxy's service management has been restructured to support distinct local and cluster execution modes via a new \ServiceManager\ interface and \ServiceManagerFactory\. The \ClusterServiceManager\ now initializes dedicated \MQClientAPIFactory\ instances for messaging, operations, transactions, and lite subscriptions, while the \LocalServiceManager\ delegates to the embedded \BrokerController\ for local testing or embedded scenarios. This change introduces a \LiteSubscriptionService\ to handle subscription synchronization for the new Lite Topic model and centralizes lifecycle management through the factory, allowing the proxy to dynamically switch its underlying remoting and client infrastructure based on the deployment mode.
proxy/src/main/java/org/apache/rocketmq/proxy/service · high confidence
Remoting module refactored into new package structure with Bazel build support
The remoting module has been reorganized into a new package structure (org.apache.rocketmq.remoting) and is now buildable with Bazel. This change introduces new core interfaces for the remoting layer, including RemotingClient, RemotingServer, and RemotingService, along with supporting types like InvokeCallback, RPCHook, and Configuration. It also adds new exception classes (RemotingException, RemotingConnectException, etc.) and utility classes (RemotingHelper, ServiceThread) to support the refactored architecture.
remoting · high confidence
Behavioural changes
Batched broker unregistration for faster offline processing
The NameServer now processes broker unregistration requests in batches rather than individually. A new \BatchUnregistrationService\ runs as a background thread, collecting \UnRegisterBrokerRequestHeader\ messages into a queue and processing them together via \RouteInfoManager.unRegisterBroker\. This change speeds up the broker-offline process, reducing latency when multiple brokers disconnect or are deregistered simultaneously.
namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo · high confidence
Broker initialization and configuration loading restructured
The broker startup process has been refactored to use a new ConfigContext object for passing configuration, and the BrokerController now initializes a comprehensive set of services including the new BrokerPreOnlineService for coordinated HA handshakes, a PID-based adaptive cold-read control strategy, and a v2 RocksDB-based configuration storage layer (ConfigHelper, TableId, SerializationType) alongside the existing v1 managers.
broker/src/main/java/org/apache/rocketmq/broker · high confidence
Controller introduces dedicated replica and sync-state management classes
The controller's broker metadata handling is refactored to use new internal classes: BrokerReplicaInfo, ReplicasInfoManager, and SyncStateInfo. BrokerReplicaInfo now tracks broker instances by ID rather than just IP, storing IP addresses and registration check codes. ReplicasInfoManager acts as the central state machine for managing these replicas and their sync state sets, enforcing epoch checks for master and sync state set updates to prevent stale configurations. SyncStateInfo manages the current master broker ID, master epoch, and sync state set epoch for each broker group, ensuring consistent leader election and failover logic within the controller.
controller/src/main/java/org/apache/rocketmq/controller/impl/manager · high confidence
DefaultMonitorListener now uses shaded SLF4J logging
The DefaultMonitorListener in the RocketMQ tools module has been updated to use the internal shaded SLF4J logging implementation (org.apache.rocketmq.logging.org.slf4j) instead of the standard external SLF4J API. This change ensures that monitoring logs are handled consistently with the project's internal logging strategy, avoiding potential conflicts with user-provided logging frameworks.
rocketmq-tools · high confidence
Distribution packaging restructured with dledger configs and logback XMLs
The distribution module now includes example configuration files for dledger (Raft-based) broker nodes (broker-n0.conf, broker-n1.conf, broker-n2.conf) to support high-availability setups. Additionally, the release assembly (release.xml) has been updated to bundle specific logback XML configuration files (for broker, client, controller, namesrv, tools, and proxy) directly into the distribution's conf directory, ensuring consistent logging behavior out-of-the-box. The release process also explicitly excludes Jaeger tracing dependencies from the final package.
distribution · high confidence
Introduction of queue-based transactional message implementation
The broker now uses a new queue-based architecture for transactional messages, replacing the previous implementation. This change introduces a dedicated batch service (TransactionalOpBatchService) to handle operation message batching, a new listener (DefaultTransactionalMessageCheckListener) that moves half-messages exceeding the maximum check limit to the TRANS\_CHECK\_MAXTIME\_TOPIC system topic, and a refactored service layer (TransactionalMessageServiceImpl) with updated bridge logic for fetching and processing half and operation messages.
broker/src/main/java/org/apache/rocketmq/broker/transaction/queue · high confidence
NamesrvController refactoring and embedded controller support
The NameServer startup logic has been restructured to introduce a dedicated NamesrvController component that centralizes initialization, network setup, and scheduled tasks. This change adds support for an embedded controller mode: when the enableControllerInNamesrv configuration is enabled, the startup process now also initializes and starts a ControllerManager alongside the NameServer. Additionally, the command-line interface now uses the DefaultParser instead of the deprecated PosixParser, and the -p flag has been updated to print configuration properties for the embedded controller when that mode is active.
namesrv/src/main/java/org/apache/rocketmq/namesrv · high confidence
New gRPC interceptors and request mapping in the proxy
The proxy now includes a set of new gRPC server interceptors (ContextInterceptor, GlobalExceptionInterceptor, HeaderInterceptor) and a RequestMapping utility. ContextInterceptor propagates metadata into the gRPC context, while GlobalExceptionInterceptor provides centralized error handling and logging for gRPC calls. HeaderInterceptor extracts and injects headers such as authorization keys, remote/local addresses, and proxy protocol information into the gRPC metadata. RequestMapping defines the correspondence between gRPC v2 request types (e.g., SendMessage, ReceiveMessage, AckMessage) and internal RocketMQ request codes, enabling the proxy to route gRPC requests to the appropriate remoting handlers.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor · high confidence
New remoting proxy activity handlers and request pipeline
The proxy now includes a new set of remoting activity handlers (AckMessage, ChangeInvisibleTime, ConsumerManager, GetTopicRoute, PopMessage, PullMessage, RecallMessage, SendMessage, Transaction) and a request pipeline (ContextInitPipeline) to process client requests. These components handle core messaging operations such as sending, pulling, popping, and acknowledging messages, managing consumer groups and connections, retrieving topic routes, and processing transactional messages. The pipeline initializes the proxy context with channel and client information. This change introduces new behavior for how the proxy processes these specific remoting commands, including subscription validation for pull messages and topic type validation for send and recall operations.
proxy/src/main/java/org/apache/rocketmq/proxy/remoting/activity · high confidence
Proxy transaction state management refactored to use in-memory storage
The proxy's transaction service implementation has been rewritten to store transaction state in memory rather than encoding it in the transaction ID. This change introduces a new \TransactionDataManager\ and associated data models (\TransactionData\, \EndTransactionRequestData\) to handle transaction lifecycle, including configurable expiration and isolation by producer group. For users, this improves transaction handling reliability within the proxy and aligns with the v2 transaction protocol support.
proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction · high confidence
Proxy validates topic message type consistency
The proxy now enforces that the message type specified in a request matches the topic's configured type. A new validation layer (TopicMessageTypeValidator) checks the TopicMessageType attribute and throws a MESSAGE\_PROPERTY\_CONFLICT\_WITH\_TYPE error if there is a mismatch, with the exception of topics configured as MIXED which accept any type. This prevents users from sending messages with an incompatible type to a topic that requires a specific one.
proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator · high confidence
Refactor gRPC v2 messaging layer with new activity interfaces and base classes
The proxy's gRPC v2 messaging implementation has been restructured to improve code organization and maintainability. A new \GrpcMessagingActivity\ interface defines the contract for messaging operations (such as sending, receiving, and acknowledging messages), while \AbstractMessagingActivity\ provides a base class with shared validation logic for topics, consumer groups, and invisible times. \DefaultGrpcMessagingActivity\ implements this interface, delegating specific operations to dedicated activity classes (e.g., \SendMessageActivity\, \ReceiveMessageActivity\) and managing shared resources like \GrpcClientSettingsManager\ and \GrpcChannelManager\. Additionally, a new \ContextStreamObserver\ interface ensures that gRPC stream callbacks consistently receive the \ProxyContext\, enabling better context propagation throughout the messaging pipeline.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2 · high confidence
Refactor proxy relay service to support cluster mode and new result types
The proxy's relay service layer has been refactored to introduce a new \ProxyRelayService\ interface and a \ProxyRelayResult\ wrapper, replacing the previous \ProxyOutService\ and \ProxyOutResult\ naming conventions. A new \ClusterProxyRelayService\ implementation has been added to handle relay operations in cluster mode, ensuring that relay futures are correctly bridged to the caller's \CompletableFuture\ so that results (including exceptions) are observed rather than lost. This change also introduces a \RelayData\ class to encapsulate process results and relay futures, improving how transaction and consumer info requests are processed and returned to clients.
proxy/src/main/java/org/apache/rocketmq/proxy/service/relay · high confidence
Refactored NameServer request handling into dedicated processor classes
The NameServer's request processing logic has been reorganized from a monolithic structure into specific processor components. A new \ClientRequestProcessor\ now handles client route queries, introducing logic to reject requests if the NameServer is not yet ready (based on startup time and wait configuration) and supporting standard JSON responses for newer clients. A \ClusterTestRequestProcessor\ extends this to fetch route info from a product environment if local data is missing. The \DefaultRequestProcessor\ now centrally manages administrative and broker-related requests (such as registration, KV config, and topic deletion) via a request-code switch, replacing the previous inline handling. This change improves code modularity and adds readiness checks for client requests.
namesrv/src/main/java/org/apache/rocketmq/namesrv/processor · high confidence
Refactored Pop Lite orderly consumption with new management interface
The broker's Pop Lite orderly consumption logic has been refactored to support more flexible concurrency strategies. A new \ConsumerOrderInfoManager\ interface defines the contract for managing ordered consumption state, replacing the previous monolithic implementation. The existing queue-level ordering logic is now encapsulated in \QueueLevelConsumerManager\, which implements this interface to handle message status updates, blocking checks, and offset commits specifically for queue-level ordering. This change isolates the ordering mechanism, facilitating future expansion to support message-group-level ordering and custom strategies without altering the core broker flow.
broker/src/main/java/org/apache/rocketmq/broker/pop/orderly · high confidence
Refactored broker client connection management with optimized channel handling
The broker's client connection management has been restructured to improve concurrency and performance. A new \ClientChannelAttributeHelper\ now stores consumer and producer group associations directly on Netty \Channel\ attributes, enabling faster lookups during channel close events. The \ConsumerManager\ and \ProducerManager\ have been updated to leverage this helper when \enableFastChannelEventProcess\ is enabled, reducing the overhead of scanning global tables. Additionally, a new \ClientHousekeepingService\ runs a scheduled task to periodically scan for and clean up inactive client channels, ensuring timely resource reclamation. The \DefaultConsumerIdsChangeListener\ also introduces a batching mechanism for consumer change notifications, allowing for more efficient updates when real-time notification is disabled.
broker/src/main/java/org/apache/rocketmq/broker/client · high confidence
Refactored proxy route service with queue penalization and priority selection
The proxy's topic routing logic in the \route\ package has been restructured to support more flexible and fault-aware message queue selection. A new \AddressableMessageQueue\ class now wraps standard queues with explicit broker addresses, enabling the \ClusterTopicRouteService\ to correctly resolve broker locations in cluster mode. Queue selection is now driven by a \MessageQueuePenalizer\ system that applies penalty scores (aggregated from multiple sources like fault strategies) to rank queues, and a \MessageQueuePriorityProvider\ that groups queues by priority levels. The \MessageQueueView\ and \MessageQueueSelector\ components now utilize these mechanisms to select the least-penalty, highest-priority queue for reads and writes, while \ProxyTopicRouteData\ and \TopicRouteHelper\ provide supporting data structures and error handling for topic route resolution.
proxy/src/main/java/org/apache/rocketmq/proxy/service/route · high confidence
gRPC proxy now supports Proxy Protocol and configurable server tuning
The gRPC proxy now supports the Proxy Protocol (v1 and v2), allowing it to correctly handle and expose the original client IP and port when operating behind a load balancer or proxy that terminates TCP connections. Additionally, the gRPC server builder has been refactored to expose several new configuration options for better performance and stability, including limits on concurrent calls per connection, keep-alive settings, and customizable event loop thread counts. TLS certificate reloading has also been improved to prevent native memory leaks by properly releasing old SSL contexts.
proxy/src/main/java/org/apache/rocketmq/proxy/grpc · high confidence
Test coverage
Added BaseServiceTest base class for proxy service tests; Added test certificates for password-encrypted private keys; Added test for InvocationChannel context handling; Added test for master election offset overflow handling; Added test resources for ACL and logging configuration; Added test resources for RMQ proxy configuration; Added tests for FAQUrl utility methods; Added tests for GenericMapSuperclassDeserializer with fastjson2; Added tests for HTTP/2 proxy protocol handlers; Added tests for RecallMessageHandle encoding and decoding; Added tests for RemotingProtocolServer and AuthorizationPipeline; Added tests for TlsCertificateManager; Added tests for gRPC proxy server lifecycle, TLS negotiation, and admin service; Added tests for proxy startup and command-line argument parsing; Added tests for topic and group validation logic; Added unit tests for Action and RocketMQAction; Added unit tests for BitsArray utility; Added unit tests for BrokerContainer lifecycle and configuration; Added unit tests for ClientConfig and Validators; Added unit tests for ClusterMessageService and LocalMessageService; Added unit tests for ClusterMetadataService; Added unit tests for ConfigHelper; Added unit tests for DefaultAdminService; Added unit tests for DefaultReceiptHandleManager; Added unit tests for FilterUtil tag matching and subscription data building; Added unit tests for GrpcClientChannel behavior; Added unit tests for LatencyFaultToleranceImpl; Added unit tests for LiteSubscription; Added unit tests for LocalProxyRelayService and ProxyChannel; Added unit tests for MQClientAPIExt and MQClientAPIFactory; Added unit tests for MQClientInstance; Added unit tests for MqClientAdminImpl; Added unit tests for NamespaceRpcHook; Added unit tests for OpenMessaging RocketMQ module; Added unit tests for RemoteChannel encoding and decoding; Added unit tests for RemotingChannel and RemotingChannelManager; Added unit tests for RocketMQ consumer implementations; Added unit tests for StatsItemSet and MomentStatsItemSet; Added unit tests for ThreadLocalIndex; Added unit tests for broker identity and heartbeat management; Added unit tests for client implementation components; Added unit tests for client message tracing and OpenTracing integration; Added unit tests for client-side ACL common utilities; Added unit tests for common chain, cold counter, and consumer receipt handle components; Added unit tests for common message components; Added unit tests for common module components; Added unit tests for common utility classes; Added unit tests for compression and pull consumer system flags; Added unit tests for consumer implementation classes; Added unit tests for consumer offset storage components; Added unit tests for controller replica and heartbeat management; Added unit tests for gRPC ClientActivity; Added unit tests for gRPC proxy authentication and authorization pipelines; Added unit tests for gRPC v2 proxy messaging components; Added unit tests for gRPC v2 route activity; Added unit tests for message queue allocation strategies; Added unit tests for message queue selectors and producer implementation; Added unit tests for message reply utilities; Added unit tests for proxy client services; Added unit tests for proxy configuration and initialization; Added unit tests for proxy remoting activities; Added unit tests for proxy system message handling; Added unit tests for the NameServer module; Added unit tests for the auth module's authentication and authorization components; Added unit tests for the controller module; Added unit tests for the gRPC EndTransaction activity; Added unit tests for the message filtering module; Expanded unit test coverage for broker core components; New Bazel build infrastructure and test harness for integration testing; Updated Mockito configuration for proxy tests.
Dependencies
Enable Bazel build system for the store module
The store module now supports building and testing via Bazel. A new BUILD.bazel file defines the Java library and test targets, explicitly listing dependencies such as fastjson2, Netty, Guava, RocksDB, and OpenTelemetry. The build configuration also includes a GenTestRules macro to manage test execution, specifically excluding known flaky or slow tests (AutoSwitchHATest, DLedgerCommitlogTest, DLedgerMultiPathTest) and classifying others as medium tests to optimize CI performance.
store · high confidence
Introduce Maven multi-module structure with dependency management
The project has been restructured into a Maven multi-module build, adding dedicated pom.xml files for core components including auth, broker, client, common, container, controller, distribution, example, filter, namesrv, openmessaging, and tools. This change centralizes dependency management in the root pom.xml and explicitly declares inter-module dependencies (e.g., broker depending on auth and tiered-store), ensuring consistent versioning and cleaner build isolation across the RocketMQ codebase.
(dependencies) · high confidence
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
How this codebase got here
Score
- CAI 55 → 61 (+6.0)
- Rubric changed (rubric-2026.08.19 → rubric-2026.09.15) — scores are not directly comparable.
Lenses
- Code Health 72 → 82 (+10.0)
- Architecture 99 → 98 (-0.5)
- Maturity 70 → 69 (-0.3)
- Readiness 59 → 61 (+1.9)
- Security 38 → 49 (+10.4)
- Domain Modelling 100 → 100 (+0.0)
Resolved (622)
- 1.run (cognitive 28) (example/src/main/java/org/apache/rocketmq/example/benchmark/BatchProducer.java)
- 1.run (cognitive 29) (example/src/main/java/org/apache/rocketmq/example/simple/PullConsumer.java)
- 1.run (cyclomatic 17) (example/src/main/java/org/apache/rocketmq/example/simple/PullConsumer.java)
- 3.run (cognitive 17) (example/src/main/java/org/apache/rocketmq/example/benchmark/TransactionProducer.java)
- 3.run (cognitive 25) (example/src/main/java/org/apache/rocketmq/example/benchmark/timer/TimerProducer.java)
- 3.run (cognitive 45) (example/src/main/java/org/apache/rocketmq/example/benchmark/Producer.java)
- 3.run (cyclomatic 20) (example/src/main/java/org/apache/rocketmq/example/benchmark/Producer.java)
- BatchProducer.main (cognitive 18) (example/src/main/java/org/apache/rocketmq/example/benchmark/BatchProducer.java)
- Boundary-crossing change coupling: MQClientAPIImpl.java ↔ DefaultMQAdminExt.java (client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java)
- Boundary-crossing change coupling: PositiveAtomicCounter.java ↔ DefaultMQProducer.java (broker/src/main/java/org/apache/rocketmq/broker/util/PositiveAtomicCounter.java)
- Boundary-crossing change coupling: PositiveAtomicCounter.java ↔ DefaultMQProducerImpl.java (broker/src/main/java/org/apache/rocketmq/broker/util/PositiveAtomicCounter.java)
- Boundary-crossing change coupling: PositiveAtomicCounter.java ↔ MQProducer.java (broker/src/main/java/org/apache/rocketmq/broker/util/PositiveAtomicCounter.java)
- Consumer.main (cognitive 24) (example/src/main/java/org/apache/rocketmq/example/benchmark/Consumer.java)
- Consumer.main (cyclomatic 22) (example/src/main/java/org/apache/rocketmq/example/benchmark/Consumer.java)
- Coverage not included — suite not readable by the collector
- DefaultAuthorizationContextBuilder.build (cognitive 24) (auth/src/main/java/org/apache/rocketmq/auth/authorization/builder/DefaultAuthorizationContextBuilder.java)
- DefaultAuthorizationContextBuilder.build (cognitive 70) (auth/src/main/java/org/apache/rocketmq/auth/authorization/builder/DefaultAuthorizationContextBuilder.java)
- DefaultAuthorizationContextBuilder.build (cyclomatic 20) (auth/src/main/java/org/apache/rocketmq/auth/authorization/builder/DefaultAuthorizationContextBuilder.java)
- DefaultAuthorizationContextBuilder.build (cyclomatic 47) (auth/src/main/java/org/apache/rocketmq/auth/authorization/builder/DefaultAuthorizationContextBuilder.java)
- Dependency hygiene not measured — dependency manifest found but not parsed for hygiene
- …and 602 more
New (1464)
- AdminModelConverter.toAccumulation (cognitive 26) (proxy/src/main/java/org/apache/rocketmq/proxy/grpc/admin/AdminModelConverter.java)
- AdminModelConverter.toConsumerRunningInfo (cognitive 19) (proxy/src/main/java/org/apache/rocketmq/proxy/grpc/admin/AdminModelConverter.java)
- Change coupling: BrokerConfig.java ↔ ControllerConfig.java (common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java)
- Change coupling: ConsumerProgressSubCommand.java ↔ TopicStatusSubCommand.java (tools/src/main/java/org/apache/rocketmq/tools/command/consumer/ConsumerProgressSubCommand.java)
- Change coupling: ContextVariable.java ↔ GrpcMessagingApplication.java (proxy/src/main/java/org/apache/rocketmq/proxy/common/ContextVariable.java)
- Change coupling: ControllerConfig.java ↔ NamesrvConfig.java (common/src/main/java/org/apache/rocketmq/common/ControllerConfig.java)
- Change coupling: DefaultMQPushConsumer.java ↔ ConsumeMessageService.java (client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java)
- Change coupling: NotificationProcessor.java ↔ PopMessageProcessor.java (broker/src/main/java/org/apache/rocketmq/broker/processor/NotificationProcessor.java)
- Change coupling: RemotingClient.java ↔ NettyRemotingAbstract.java (remoting/src/main/java/org/apache/rocketmq/remoting/RemotingClient.java)
- Change coupling: StoreCheckpoint.java ↔ IndexService.java (store/src/main/java/org/apache/rocketmq/store/StoreCheckpoint.java)
- Change-coupling hub: ChangeInvisibleTimeProcessor.java → AckMessageProcessor.java, PopBufferMergeService.java, PopReviveService.java (broker/src/main/java/org/apache/rocketmq/broker/processor/ChangeInvisibleTimeProcessor.java)
- Change-coupling hub: ClusterMessageService.java → DefaultMessagingProcessor.java, MessagingProcessor.java, LocalMessageService.java (proxy/src/main/java/org/apache/rocketmq/proxy/service/message/ClusterMessageService.java)
- ClassTooLong: AbstractRocksDBStorage (common/src/main/java/org/apache/rocketmq/common/config/AbstractRocksDBStorage.java)
- ClassTooLong: AdminBrokerProcessor (broker/src/main/java/org/apache/rocketmq/broker/processor/AdminBrokerProcessor.java)
- ClassTooLong: AutoSwitchHAConnection (store/src/main/java/org/apache/rocketmq/store/ha/autoswitch/AutoSwitchHAConnection.java)
- ClassTooLong: BatchConsumeQueue (store/src/main/java/org/apache/rocketmq/store/queue/BatchConsumeQueue.java)
- ClassTooLong: BrokerConfig (common/src/main/java/org/apache/rocketmq/common/BrokerConfig.java)
- ClassTooLong: BrokerController (broker/src/main/java/org/apache/rocketmq/broker/BrokerController.java)
- ClassTooLong: BrokerMetricsManager (broker/src/main/java/org/apache/rocketmq/broker/metrics/BrokerMetricsManager.java)
- ClassTooLong: BrokerOuterAPI (broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java)
- …and 1444 more
Changes since last survey
- 49 commits — 26 feature/other, 23 fixes
By area
- broker/src — 15 commits
- store/src — 8 commits
- proxy/src — 5 commits
- remoting/src — 5 commits
- (root) — 3 commits
- client/src — 3 commits
- auth/src — 2 commits
- common/src — 2 commits
- controller/src — 2 commits
- filter/src — 1 commit
- namesrv/src — 1 commit
- test/src — 1 commit
- tieredstore/src — 1 commit
Notable commits
- fix: [ISSUE #10383] Fix flaky store and proxy tests (#10407)
- fix: [ISSUE #10575] Fix race condition between scanResponseTable and processResponseCommand (#10576)
- fix: [ISSUE #10827] fix(broker): spin for the lock on same-attemptId pop orderly retry to avoid empty response (#10828)
- fix: [ISSUE #10853] fix(store): validate HA state ordinals (#11200)
- fix: [ISSUE #10877] fix(filter): convert boolean string operands safely
- fix: [ISSUE #10959] Fix RocksDBConsumeQueue.iterateFrom should reject offset below min offset (#10960)
- fix: [ISSUE #11007] Fix flaky CreateAndUpdateTopicIT by awaiting route propagation (#11008)
- fix: [ISSUE #11011] Fix misplaced license header in ConfigManagerTest (#11012)
- fix: [ISSUE #11091] Fix remoting sub-server timeout scanning and scheduler shutdown (#11092)
- fix: [ISSUE #11129] Fix HashedWheelTimer leak of Lite pop on broker shutdown (#11130)
- fix: [ISSUE #11156]fix(timer): persist TimelineRollService checkpoint to avoid repeated and loss (#11157)
- fix: [ISSUE #11161] Fix closeChannel table eviction and ChannelWrapper.close lock ordering (#11162)
- fix: [ISSUE #11163] Fix lite topic prefix index not maintained on first message (#11164)
- fix: [ISSUE #11168] Fix tiered storage commit failure reporting and related hazards (#11169)
- fix: [ISSUE #11176] Fix incorrect filter appender name in broker logging configuration (#11177)
- fix: [ISSUE #11191] Fix Lite topic consumer losing messages on event re-dispatch during a pop (#11192)
- fix: fix(auth): address ACL follow-up regressions (#11009)
- fix: fix(auth): tolerate blank subscription topics in heartbeats (#11082)
- fix: fix(broker): reject timer delays that overflow the absolute delivery time (#10965)
- fix: fix(common): compare message queue ids safely (#10884)
- …and 29 more
Architecture
- Containers 0 added · 0 removed · contexts 14 added · 0 removed · edges 41 added · 0 removed
Added bounded contexts (14)
- rocketmq-auth
- rocketmq-broker
- rocketmq-client
- rocketmq-common
- rocketmq-container
- rocketmq-controller
- rocketmq-filter
- rocketmq-namesrv
- rocketmq-openmessaging
- rocketmq-proxy
- rocketmq-remoting
- rocketmq-srvutil
- rocketmq-store
- rocketmq-tiered-store
Added dependency edges (41)
- rocketmq-auth → rocketmq-common
- rocketmq-auth → rocketmq-remoting
- rocketmq-broker → rocketmq-auth
- rocketmq-broker → rocketmq-client
- rocketmq-broker → rocketmq-common
- rocketmq-broker → rocketmq-filter
- rocketmq-broker → rocketmq-remoting
- rocketmq-broker → rocketmq-srvutil
- rocketmq-broker → rocketmq-store
- rocketmq-client → rocketmq-common
- rocketmq-client → rocketmq-remoting
- rocketmq-container → rocketmq-auth
- rocketmq-container → rocketmq-broker
- rocketmq-container → rocketmq-common
- rocketmq-container → rocketmq-remoting
- rocketmq-container → rocketmq-store
- rocketmq-controller → rocketmq-common
- rocketmq-controller → rocketmq-remoting
- rocketmq-namesrv → rocketmq-common
- rocketmq-namesrv → rocketmq-controller
- …and 21 more
Written by watchdog.canine.dev from the codebase's own history, inside the signed delivery this page is composed from.
Survey your own repository
apache/rocketmq was measured the same way every project in this corpus was: the same rubric, at a pinned commit, with the result published in full. Point a surveyor at a repository you know and see whether you agree with it.
About this page
- The score is its most recent published measurement, taken on 24 September 2026 at a pinned commit. It is not a live figure and does not change until the project is measured again.
- Measured at commit 78b96bc5e21216cd7896efae08f90c5cde4cae53 — the exact code this score is about.
- Scored under rubric-2026.09.15 — the same rubric and the same method as every other entry in this index.
- Measured by watchdog.canine.dev using codehealth-analyzer preprod-ae95d6cad036.