Skip to content
CAI
Software that uses CAICheck a score

yuanzhongqiao/rocketmq

56.2

Adequate · 22 September 2026

181.6k

lines of production code

Java

primary language

3

measurements over time

CAI band scale
CAI trend line
CAI lens gauges

What this system is

This system is a distributed message queue platform that manages message storage, routing, and delivery through a broker and controller architecture. It provides comprehensive client-side APIs for producing and consuming messages, supporting patterns such as transactional messaging, ordered delivery, and high-performance POP consumption. The system also features a proxy layer for gRPC and remoting protocols, enabling flexible integration and cluster management.

How it got here

2013–2016 — broker and client feature expansion

53 changes.

This period focused on significantly expanding the RocketMQ broker and client capabilities, introducing new consumption modes like POP and LitePull, and enhancing broker features such as cold read control, transaction tracking, and persistent storage via RocksDB. The client library was also modernized with new APIs, hooks, and fault tolerance mechanisms, while the build system was upgraded to support Bazel alongside Maven.

2017–2018 — ACL, filtering, and tracing features

55 changes.

This period focused on implementing core security and observability features, specifically introducing Access Control List (ACL) support, SQL92 message filtering, and comprehensive message tracing with OpenTracing integration. The work also included significant expansion of unit test coverage across the client, broker, and filter modules to ensure the stability of these new capabilities.

2019–2022 — Proxy and Controller Architecture

75 changes.

This period focused on introducing the Controller module for high-availability broker management and building a comprehensive gRPC v2 Proxy layer. The work established a new service-oriented architecture for the proxy, supporting features like request-response, scheduled messages, and transactional messaging via both gRPC and remoting protocols.

2023 — observability and storage features

7 changes.

This period focused on introducing new capabilities for tiered storage and message receipt handling, alongside enhanced observability through OpenTelemetry metrics. The work also included significant test coverage for the new components and access control mechanisms.

Features

Add Bazel build system and project configuration files

The project introduces a Bazel-based build system, adding configuration files such as .bazelrc, .bazelversion (v5.2.0), and the root-level .asf.yaml for GitHub automation. The root-level .gitignore is updated to exclude Bazel-generated directories (bazel-out, bazel-bin, etc.), and a new .licenserc.yaml file is added to manage license header checks. These changes provide an alternative, reproducible build environment alongside the existing Maven setup.

(repo-wide) · high confidence

Add MessageRequestModeManager for pop consuming mode

A new \MessageRequestModeManager\ class has been added to the broker's load balance package. This manager handles the storage and retrieval of message request modes for topics and consumer groups, supporting the 'pop consuming' feature introduced in RIP-19. It maintains a concurrent map of topic and consumer group to request body, allowing the broker to track and manage message request modes persistently.

broker/src/main/java/org/apache/rocketmq/broker/loadbalance · high confidence

Add OpenMessaging domain classes for message handling and results

The codebase now includes new domain classes in the \io.openmessaging.rocketmq.domain\ package to support the OpenMessaging specification. This includes \BytesMessageImpl\ for handling byte array message bodies, \SendResultImpl\ to wrap message send outcomes, \ConsumeRequest\ to manage message consumption requests, and constants/interfaces like \NonStandardKeys\ and \RocketMQConstants\ to define RocketMQ-specific headers and system keys.

openmessaging/src/main/java/io/openmessaging/rocketmq/domain · medium confidence

Add OpenMessaging promise implementation

The \openmessaging\ module now includes a new \DefaultPromise\ class and a \FutureState\ enum to implement the OpenMessaging specification. This adds support for asynchronous operations and future/promise patterns within the RocketMQ client, allowing users to handle asynchronous results and timeouts more effectively.

openmessaging/src/main/java/io/openmessaging/rocketmq/promise · high confidence

Add OpenTelemetry-based metrics collection for the controller

The controller now exposes operational metrics (such as request counts, latencies, and broker status) via OpenTelemetry. This change introduces new constants and a metrics manager that supports exporting data to OTLP gRPC, Prometheus, or log-based endpoints, enabling users to monitor controller health and performance through standard observability tools.

controller/src/main/java/org/apache/rocketmq/controller/metrics · high confidence

Add SQL92 message filtering support

The filter module now includes a new SQL92 message filtering capability. This change introduces a plugin-based architecture with a \FilterFactory\ for managing filter implementations, a \FilterSpi\ interface for defining filters, and a \SqlFilter\ that wraps the \SelectorParser\ to enable SQL92-based message filtering. The update also adds a comprehensive set of expression classes (e.g., \BinaryExpression\, \LogicExpression\, \ComparisonExpression\) to support complex filtering conditions, including \CONTAINS\, \STARTSWITH\, and \ENDsWith\ operations on message properties.

filter/src/main/java/org/apache/rocketmq/filter/expression · high confidence

Add SQL92 message filtering support with CONTAINS, STARTSwith, and ENDsWITh operators

The RocketMQ filter parser now supports SQL92-style message filtering. Users can filter messages using the CONTAINS, STARTSWITH, and ENDSWITH operators in their selector expressions, enabling more flexible content-based message routing.

filter/src/main/java/org/apache/rocketmq/filter/parser · high confidence

Add SlaveSynchronize for master-slave data replication

A new SlaveSynchronize class has been introduced in the broker module to handle synchronization of configuration and state from the master broker to the slave. This includes syncing topic configurations, consumer offsets, delay offsets, and subscription group settings, ensuring the slave broker maintains a consistent state with the master.

broker/src/main/java/org/apache/rocketmq/broker/slave · high confidence

Add benchmark scripts for producer, consumer, and shutdown

New shell scripts are added to the benchmark directory to manage the lifecycle of benchmark tools. The \runclass.sh\ script provides a unified way to execute benchmark classes with appropriate Java options and GC logging, handling platform-specific log directories (RAM disk on Darwin, /dev/shm on Linux). Convenience scripts (\producer.sh\, \consumer.sh\, \tproducer.sh\, \batchproducer.sh\) are introduced to launch the respective benchmark implementations. Additionally, a \shutdown.sh\ script is added to gracefully stop running benchmark processes by their specific class names.

distribution/benchmark · high confidence

Add broadcast push consumer example

A new PushConsumer example is added to the broadcast package, demonstrating how to configure a consumer for broadcasting mode using the default namesrv address and a sample topic subscription.

example/src/main/java/org/apache/rocketmq/example/broadcast · high confidence

Add configuration files for 2-master, no-slave deployment

New configuration files (broker-a.properties, broker-b.properties, and broker-trace.properties) are added to the distribution/conf/2m-noslave directory, enabling users to deploy a two-master, no-slave message queue cluster with asynchronous flushing and master roles.

distribution/conf/2m-noslave · high confidence

Add core proxy context and message receipt handling classes

The proxy module introduces several new classes to support message receipt handling and context management. This includes \ProxyContext\ for managing request-scoped data (such as client ID, channel, and action), \MessageReceiptHandle\ to track message state and renewal attempts, and \RenewEvent\/\RenewStrategyPolicy\ to manage message renewal strategies. Additionally, \ContextVariable\ defines keys for context storage, \Address\ handles host/port logic, and \ProxyException\/\ProxyExceptionCode\ provide structured error handling.

proxy/src/main/java/org/apache/rocketmq/proxy/common · high confidence

Add default broker election policy for controller mode

A new DefaultElectPolicy implementation is introduced for the RocketMQ controller, enabling broker priority-based leader election. The policy filters live brokers using a validity predicate, preserves the current master if it remains valid, respects a preferred broker ID when available, and otherwise selects the best candidate by sorting on epoch, offset, and election priority.

controller/src/main/java/org/apache/rocketmq/controller/elect/impl · high confidence

Add default configuration files for broker, ACL, proxy, and tools

The distribution/conf directory now includes four new configuration files: broker.conf, which sets default broker cluster and role settings; plain\_acl.yml, which provides a sample Access Control List configuration with white-listed remote addresses and account permissions; rmq-proxy.json, which specifies the RocketMQ cluster name for the proxy; and tools.yml, which contains default access and secret keys for tools. These files serve as the baseline configuration for these components.

distribution/conf · high confidence

Add dledger fast-try.sh script for quick start and stop

A new shell script, fast-try.sh, is added to the dledger distribution. This script provides a convenient way for users to quickly start and stop the dledger environment, including the NameServer and multiple Broker instances, by running simple start or stop commands.

distribution/bin/dledger · high confidence

Add example implementations for RocketMQ producer and consumer operations

New example classes, Consumer.java and Producer.java, are added to the org.apache.rocketmq.example.operation package. These provide reference implementations for sending and receiving messages using the RocketMQ client libraries, including command-line argument parsing for configuration.

example/src/main/java/org/apache/rocketmq/example/operation · high confidence

Add gRPC server implementation with Proxy Protocol and TLS support

The gRPC server is now implemented in the proxy module, introducing a new GrpcServer class that manages the lifecycle of the gRPC server instance. The GrpcServerBuilder configures the server with configurable thread pools, message size limits, and connection idle timeouts. A custom ProxyAndTlsProtocolNegotiator is added to handle both TLS and HAProxy Protocol v2 connections, allowing the proxy to extract client IP and port information from PROXY protocol headers. Additionally, the server is configured with standard interceptors for authentication, context, headers, and global exception handling.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc · high confidence

Add gRPC v2 consumer activities for message acknowledgment and visibility duration changes

The proxy now exposes gRPC v2 endpoints for consumers to acknowledge messages and change their invisible duration. AckMessageActivity handles both batch and individual message acknowledgments, returning per-message results with appropriate status codes. ChangeInvisibleDurationActivity allows clients to update message visibility timeouts. ReceiveMessageResponseStreamWriter manages the streaming response for message retrieval, handling various pop statuses and error conditions.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer · high confidence

Add gRPC v2 route handling for topic and assignment queries

Introduces RouteActivity in the proxy's gRPC v2 route package to handle topic routing and message queue assignment queries. The new implementation converts gRPC requests into internal proxy calls, constructs response objects containing message queues and broker information, and validates topics and consumer groups. This enables clients to query route and assignment data via the gRPC v2 interface.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route · high confidence

Add gRPC v2 support for ending transactions

The proxy now exposes a new gRPC v2 endpoint for ending transactions. The \EndTransactionActivity\ class handles incoming \EndTransactionRequest\ messages, validates the transaction ID, maps the resolution (commit/rollback) to a transaction status, and calls the underlying \messagingProcessor.endTransaction\ method. This adds the capability for clients to explicitly commit or roll back a transaction via the gRPC interface.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/transaction · high confidence

Add message filtering support for retry topics

The broker now includes new classes (ExpressionForRetryMessageFilter, ConsumerFilterData, ConsumerFilterManager, etc.) that enable expression-based message filtering for retry topics. This allows consumers to apply complex filtering expressions to messages in retry queues, ensuring that only relevant messages are delivered to subscribers, while also handling the decoding of message properties to resolve the actual topic for filtering.

broker/src/main/java/org/apache/rocketmq/broker/filter · high confidence

Add new example classes for RocketMQ features

The example/src/main/java/org/apache/rocketmq/example/simple directory now includes several new example files that demonstrate specific RocketMQ capabilities. These include AclClient for authentication, AsyncProducer for asynchronous messaging, OnewayProducer for one-way sends, PopConsumer for POP consumption mode, and various LitePullConsumer examples (Assign, Subscribe, and AssignWithSubExpression) for pull-based consumption. Additionally, the directory contains supporting classes like CachedQueue and RandomAsyncCommit to facilitate these examples.

example/src/main/java/org/apache/rocketmq/example/simple · high confidence

Add order message example for RocketMQ

The ordermessage example now includes both a Producer and a Consumer implementation. The Producer sends messages with specific tags and uses a MessageQueueSelector to route messages to specific queues based on an order ID, while the Consumer processes these messages in an orderly fashion, handling success and suspension scenarios.

example/src/main/java/org/apache/rocketmq/example/ordermessage · high confidence

Add proxy metrics framework and OpenTelemetry integration

The proxy module now includes a new metrics framework that exposes proxy status and configuration via OpenTelemetry. A new \ProxyMetricsManager\ initializes and manages metrics export to OTLP gRPC, Prometheus, or logging endpoints, while \ProxyMetricsConstant\ defines the associated constants. This change enables users to monitor proxy health and metrics through standard OpenTelemetry exporters.

proxy/src/main/java/org/apache/rocketmq/proxy/metrics · high confidence

Add quickstart examples for producing and consuming messages

New Java example files, Consumer.java and Producer.java, have been added to the quickstart directory. These provide reference implementations for sending and receiving messages using the RocketMQ client libraries, including configuration for the name server address, topic, and consumer group.

example/src/main/java/org/apache/rocketmq/example/quickstart · high confidence

Add startup, shutdown, and utility scripts for RocketMQ components

The distribution/bin directory now includes a comprehensive set of shell and Windows batch scripts to manage RocketMQ services. This includes launchers for the NameServer, Broker, Proxy, and Controller, as well as a unified shutdown script. Additionally, utility scripts are provided for OS tuning (os.sh), cache management (cachedog.sh, cleancache.sh), and a quick-start script (play.sh) to launch a basic cluster. The broker launcher (mqbroker) has been updated to support the --enable-proxy flag, allowing users to start the broker with the proxy enabled directly from the command line.

distribution/bin · high confidence

Add timer benchmark examples for message delay and timing

Added new benchmark examples (TimerProducer and TimerConsumer) in the RocketMQ example directory to demonstrate and measure performance for timing messages with arbitrary time delays. These examples allow users to benchmark send/receive throughput, average and percentile delayed durations, and other metrics related to timer functionality.

example/src/main/java/org/apache/rocketmq/example/benchmark/timer · high confidence

Add transactional message example

The RocketMQ example project now includes a complete transactional message demonstration, consisting of a \TransactionProducer\ that sends messages using the \TransactionMQProducer\ API and a \TransactionListenerImpl\ that handles local transaction states. This addition provides users with a reference implementation for building reliable transactional messaging workflows with RocketMQ.

example/src/main/java/org/apache/rocketmq/example/transaction · high confidence

Added Java examples for message filtering

Added four new Java example files to the RocketMQ codebase: SqlFilterProducer, SqlFilterConsumer, TagFilterProducer, and TagFilterConsumer. These examples demonstrate how to use tag-based and SQL92-based message filtering in both producer and consumer scenarios, providing users with reference implementations for these features.

example/src/main/java/org/apache/rocketmq/example/filter · high confidence

Added OpenMessaging (OMS) example code for producers and consumers

New example files were added to demonstrate the OpenMessaging 0.1.0-alpha specification, including a SimpleProducer that sends messages synchronously, asynchronously, and oneway, as well as SimplePullConsumer and SimplePushConsumer implementations that receive and acknowledge messages using the OMS API.

example/src/main/java/org/apache/rocketmq/example/openmessaging · high confidence

Added OpenTracing and message trace hooks for message consumption, sending, and transaction completion

The client now includes new hook implementations for distributed tracing and message tracking. For OpenTracing, hooks are added for consuming messages, sending messages, and ending transactions, each creating and managing OpenTracing spans with relevant message metadata. Additionally, new trace hooks are introduced for consuming and sending messages, as well as ending transactions, which capture trace context and dispatch trace data via the local dispatcher. These changes enable end-to-end tracing and detailed message tracking for RocketMQ clients.

client/src/main/java/org/apache/rocketmq/client/trace/hook · high confidence

Added RebalanceLockManager for message queue locking

A new RebalanceLockManager class has been introduced in the broker's rebalance package to manage distributed locks for message queues. This component tracks lock states per consumer group and message queue, supporting both single and batch lock acquisition with expiration handling, enabling more robust coordination during rebalance operations.

broker/src/main/java/org/apache/rocketmq/broker/client/rebalance · high confidence

Added Windows and Unix shell scripts for quick-starting and managing controller, namesrv, and broker instances

New launch and management scripts are now available in the controller bin directory to simplify setting up and controlling the message queue infrastructure. The 'fast-try' scripts allow users to quickly start a Namesrv and two Brokers (Windows .cmd and Unix .sh). Additionally, 'fast-try-independent-deployment' scripts enable starting three controller nodes, while 'fast-try-namesrv-plugin' scripts start three Namesrv instances. These scripts handle configuration checks, process management, and cleanup, providing a convenient way to spin up the full stack for testing or development.

distribution/bin/controller · high confidence

Added automated pull request merge script

A new Python utility, dev/merge\_rocketmq\_pr.py, has been added to the repository. This script automates the process of merging Apache RocketMQ pull requests by handling branch management, conflict resolution, and commit message formatting. It integrates with GitHub and JIRA to streamline the contribution workflow for developers.

dev · high confidence

Added batch message sending examples

Two new example classes, SimpleBatchProducer and SplitBatchProducer, have been added to demonstrate how to send messages in batches using the RocketMQ client. SimpleBatchProducer shows a basic batch send, while SplitBatchProducer includes a ListSplitter utility to handle large batches by splitting them into smaller chunks to stay within size limits.

example/src/main/java/org/apache/rocketmq/example/batch · high confidence

Added configuration files for 2m-2s broker topologies

New configuration files have been added to the distribution module for two broker topologies: 2m-2s-async and 2m-2s-sync. These files define the settings for broker-a and broker-b, including their roles (ASYNC\_MASTER, SYNC\_MASTER, and SLAVE) and disk flush types (ASYNC\_FLUSH).

distribution/conf/2m-2s-async, distribution/conf/2m-2s-sync · high confidence

Added containerized broker configuration files for high-availability setups

New configuration files have been added to the distribution/conf/container/2container-2m-2s/ directory to support the RocketMQ BrokerContainer feature. This includes individual broker configurations for master and slave nodes (broker-a and broker-b) as well as container-level configurations that link brokers to the container. These files provide the necessary settings for running a two-container, two-master, two-slave topology, enabling users to deploy a highly available RocketMQ cluster using the new containerized architecture.

distribution/conf/container · high confidence

Added core ACL interfaces for access validation and permission checking

New interfaces \AccessResource\, \AccessValidator\, and \PermissionChecker\ were added to the \acl\ module. \AccessValidator\ defines the contract for parsing access resources from requests and validating them, including methods to update and delete access configurations. \PermissionChecker\ provides a method to verify permissions between checked and owned access resources. These interfaces form the foundation for the new ACL (Access Control List) functionality.

acl/src/main/java/org/apache/rocketmq/acl · high confidence

Added gRPC v2 proxy common utilities

Introduced new classes in the proxy's gRPC v2 common package to support the v2 protocol: GrpcConverter for mapping internal message structures to gRPC messages, GrpcValidator for validating topics, consumer groups, and tags, ResponseBuilder for mapping error codes, and ResponseWriter for handling gRPC responses. These components provide the foundational conversion, validation, and response-building logic required for the gRPC v2 interface.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/common · high confidence

Added namespace example code for producer and consumer

The example/src/main/java/org/apache/rocketmq/example/namespace directory now includes three new Java files: ProducerWithNamespace, PullConsumerWithNamespace, and PushConsumerWithNamespace. These files demonstrate how to configure and use the namespace feature for sending and receiving messages in RocketMQ.

example/src/main/java/org/apache/rocketmq/example/namespace · high confidence

Added utility classes for OpenMessaging conversion and bean population

Added BeanUtils and OMSUtil helper classes in the openmessaging module. BeanUtils provides reflection-based population of JavaBeans from Properties or KeyValue objects. OMSUtil handles conversion between OpenMessaging message types (BytesMessage) and RocketMQ internal message formats, including header mapping, send result conversion, and key-value building.

openmessaging/src/main/java/io/openmessaging/rocketmq/utils · high confidence

Added utility classes for message filtering

Added new utility classes BitsArray, BloomFilter, and BloomFilterData in the filter package. BitsArray provides a wrapper for byte arrays to enable efficient single-bit operations. BloomFilter implements a simple Bloom filter algorithm using MurmurHash3, allowing for probabilistic set membership testing. BloomFilterData holds the bit positions and bit number for the filter. These classes support the SQL92 message filtering feature.

filter/src/main/java/org/apache/rocketmq/filter/util · high confidence

Added utility for creating reply messages in request-response model

A new MessageUtil class was added to the client module to support the request-response (RPC) model. It provides a createReplyMessage method that constructs a reply message by copying properties such as cluster, correlation ID, and TTL from the original request, and includes error handling for invalid request messages.

client/src/main/java/org/apache/rocketmq/client/utils · medium confidence

Benchmark suite gains ACL support and new test scenarios

The benchmark examples in example/src/main/java/org/apache/rocketmq/example/benchmark now support authentication via the new AclClient helper class, allowing users to enable ACL for producers and consumers. Additionally, the suite introduces new benchmark scripts for batch message production, transactional message handling, and delay message testing, while also adding options for message compression and configurable report intervals.

example/src/main/java/org/apache/rocketmq/example/benchmark · high confidence

Initial implementation of OpenMessaging 0.1.0-alpha specification for RocketMQ consumers

Added new client configuration and consumer implementations (Pull and Push) that map the OpenMessaging API to RocketMQ's internal APIs. This includes the ClientConfig class for managing consumer settings and the PullConsumerImpl and PushConsumerImpl classes that handle message retrieval and event-driven consumption, effectively implementing the OpenMessaging specification for this component.

openmessaging/src/main/java/io/openmessaging/rocketmq/consumer · high confidence

Introduce Broker2Client class for broker-to-client communication

A new Broker2Client class has been added to the broker module to handle communication with connected clients. This class encapsulates remote procedure calls for checking transaction states, notifying consumers of ID changes, and resetting consumer offsets, providing a centralized way for the broker to interact with client channels.

broker/src/main/java/org/apache/rocketmq/broker/client/net · high confidence

Introduce BrokerContainer to host and manage multiple broker instances

The container module now provides a new \BrokerContainer\ architecture that allows a single process to manage multiple broker instances (masters, slaves, and DLedger brokers) simultaneously. This includes a \BrokerContainerConfig\ for container-level settings, a \BrokerContainerProcessor\ to handle remote requests for adding/removing brokers, and startup logic in \BrokerContainerStartup\. This change enables users to run multiple brokers within one JVM, simplifying deployment and resource sharing.

container/src/main · high confidence

Introduce BrokerFastFailure for request timeout and system-busy flow control

A new BrokerFastFailure component has been added to the broker's latency package. It runs a scheduled task that periodically scans send, pull, heartbeat, and transaction thread-pool queues, removing expired requests and returning SYSTEM\_BUSY responses when the OS page cache is busy or when requests exceed configured wait times. This enables the broker to proactively apply flow control and prevent resource exhaustion under heavy load.

broker/src/main/java/org/apache/rocketmq/broker/latency · high confidence

Introduce BrokerReplicaInfo and ReplicasInfoManager for controller-side broker state management

The controller now maintains detailed replica and sync-state information for each broker via new classes: BrokerReplicaInfo, which tracks broker IDs, IP addresses, and register check codes; and ReplicasInfoManager, which manages the state of all broker replicas, including handling sync state set updates, master election, and broker registration. This change enables the controller to track and manage broker replicas more effectively, supporting features like automatic master election and dynamic sync state set updates.

controller/src/main/java/org/apache/rocketmq/controller/impl/manager · high confidence

Introduce Cold Read Control and Broker Pre-Online Service

The broker now includes a Cold Read Control mechanism (ColdCtrStrategy and PIDAdaptiveColdCtrStrategy) to dynamically adjust read speeds for cold data based on system load, alongside a new BrokerPreOnlineService that coordinates the broker's startup sequence by waiting for HA handshakes and synchronizing metadata (consumer offsets, delay offsets, timer checkpoints) from the master before coming online.

broker/src/main/java/org/apache/rocketmq/broker · high confidence

Introduce ControllerRequestProcessor to handle controller-specific RPC requests

A new ControllerRequestProcessor class has been added to the controller module, implementing NettyRequestProcessor to route and handle controller-specific requests such as broker registration, heartbeat, elect-master, and configuration updates. This processor centralizes the handling of these RPC calls, integrating with the ControllerManager and BrokerHeartbeatManager to manage broker state and configuration.

controller/src/main/java/org/apache/rocketmq/controller/processor · high confidence

Introduce EscapeBridge for remote message escaping

Added EscapeBridge.java, a new component in the broker's failover package that handles routing messages to remote brokers. This class provides both synchronous (putMessage) and asynchronous (asyncPutMessage) methods to send messages to a remote broker, supporting features like slave-acting-master mode and remote escape configurations. It manages internal producer/consumer groups and thread pools for asynchronous operations.

broker/src/main/java/org/apache/rocketmq/broker/failover · high confidence

Introduce MQClientAPIExt and MQClientAPIFactory for improved client API management

Added new classes DoNothingClientRemotingProcessor, MQClientAPIExt, and MQClientAPIFactory to the client module. MQClientAPIExt extends the existing MQClientAPIImpl to provide additional asynchronous and one-way methods for heartbeats, message sending, and other remoting operations, while MQClientAPIFactory manages a pool of these client instances, handling initialization, lifecycle management, and name server address updates.

client/src/main/java/org/apache/rocketmq/client/impl/mqclient · high confidence

Introduce MetadataService for topic and subscription group configuration

The proxy module now includes a new MetadataService interface and its implementations (LocalMetadataService and ClusterMetadataService) to retrieve topic type and subscription group configurations. LocalMetadataService fetches metadata directly from the local broker controller, while ClusterMetadataService caches and refreshes metadata from cluster brokers, supporting both local and distributed metadata access patterns.

proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata · high confidence

Introduce ReceiptHandleManager for managing message receipt handles

The proxy layer now includes a new ReceiptHandleManager interface and its DefaultReceiptHandleManager implementation to manage message receipt handles. This component tracks message states, handles client unregistration cleanup, and schedules renewal tasks to ensure messages are not lost during processing. This change supports the broader receipt handling feature by providing a dedicated manager for these operations.

proxy/src/main/java/org/apache/rocketmq/proxy/service/receipt · high confidence

Introduce RocksDB-backed subscription group storage and LMQ support

The broker now supports persistent subscription group configuration using RocksDB, implemented via the new \RocksDBSubscriptionGroupManager\ and \RocksDBLmqSubscriptionGroupManager\ classes. This change enables durable storage of subscription group metadata, improving data persistence and recovery. Additionally, the \LmqSubscriptionGroupManager\ adds support for Light Message Queue (LMQ) groups, allowing the broker to handle LMQ-specific subscription logic separately from standard groups.

broker/src/main/java/org/apache/rocketmq/broker/subscription · high confidence

Introduce centralized proxy configuration management

The proxy now uses a centralized configuration system (Configuration, ConfigurationManager, and ProxyConfig) to load and manage settings from a JSON file (rmq-proxy.json). This change allows users to customize proxy behavior—such as thread pool sizes, timeout values, and message limits—through a single configuration file, simplifying deployment and tuning.

proxy/src/main/java/org/apache/rocketmq/proxy/config · high confidence

Introduce channel management and invocation context for client connections

Added new classes to the proxy's channel service layer to manage client connections and handle invocation contexts. The new \ChannelManager\ coordinates the creation and lifecycle of \SimpleChannel\ instances, which provide a simplified view of Netty channels for write operations. Additionally, \InvocationContext\ and its interface are introduced to track pending RPC responses and handle timeouts, enabling the proxy to manage client connection states and invocation contexts more effectively.

proxy/src/main/java/org/apache/rocketmq/proxy/service/channel · high confidence

Introduce client-side configuration and admin interfaces

Added new client-side configuration and administrative interfaces to the RocketMQ Java client. This includes the \AccessChannel\ enum to distinguish between LOCAL and CLOUD access modes, and the \ClientConfig\ class to centralize client settings such as namespace, proxy, and timeout configurations. Additionally, new admin interfaces (\MQAdmin\, \MqClientAdmin\, \MQAdminExtInner\) and helper classes (\MQHelper\, \Validators\, \QueryResult\) are introduced to support management operations like topic creation, message querying, and offset resetting.

client/src/main/java/org/apache/rocketmq/client · high confidence

Introduce cold data flow control and enhanced remote address matching for ACL

Added a new cold data flow control mechanism that tracks consumer group cold read activity and applies adaptive or simple throttling strategies to manage hot data access. Additionally, the ACL module now supports flexible remote address matching, allowing administrators to configure IP-based access control using single addresses, ranges, or multiple addresses via the configuration file.

rocketmq-tools · high confidence

Introduce controller module with heartbeat management and master election

The controller module is introduced, providing the core logic for managing broker heartbeats and handling master election. This includes the BrokerHeartbeatManager for tracking broker liveness and triggering elections when brokers become inactive, the Controller interface and its DLedger-based implementation for managing broker registration and replica info, and the ControllerManager to orchestrate these components. Additionally, an ElectPolicy interface and its DefaultElectPolicy implementation are added to support configurable master election strategies.

controller/src/main/java/org/apache/rocketmq/controller · high confidence

Introduce core ACL utility classes for authentication and authorization

Added new classes in the \acl/common\ package to support RocketMQ's access control features. This includes \AclConstants\ for configuration keys and permission codes, \AclException\ for error handling, \AclSigner\ for computing request signatures, \AclUtils\ for IP address validation and byte manipulation, \AuthenticationHeader\ and \AuthorizationHeader\ for managing header metadata, \Permission\ for checking access rights, \SessionCredentials\ for managing access/secret keys, and \SigningAlgorithm\ for supported hashing methods. These components provide the foundational logic for authenticating and authorizing requests.

acl/src/main/java/org/apache/rocketmq/acl/common · high confidence

Introduce event-based controller architecture

The controller module now uses a new event-driven model for handling broker state changes. This introduces a set of event classes (such as CleanBrokerDataEvent, UpdateBrokerAddressEvent, and others) that implement the EventMessage interface, each associated with a specific EventType. A ControllerResult class wraps responses and associated events, while an EventSerializer handles serialization and deserialization of these events using FastJson. This change shifts the controller's internal communication from direct method calls or raw protocol messages to a structured event system, enabling better decoupling and extensibility for broker lifecycle and state management.

controller/src/main/java/org/apache/rocketmq/controller/impl/event · high confidence

Introduce latency-based fault tolerance for message queue selection

The client now includes a new latency-aware fault tolerance mechanism in the \client.latency\ package. \LatencyFaultTolerance\ and its implementation \LatencyFaultToleranceImpl\ track broker health, latency, and reachability, while \MQFaultStrategy\ uses these signals to filter out unavailable or high-latency brokers during message queue selection. This enables the client to dynamically avoid degraded brokers and improve message delivery reliability.

client/src/main/java/org/apache/rocketmq/client/latency · high confidence

Introduce new admin API and utility classes for the MQAdminExt interface

The tools module now includes a new \MQAdminExt\ interface and its \DefaultMQAdminExt\ implementation, providing a centralized API for administrative operations such as managing broker configurations, topics, and consumer groups. This change introduces supporting classes like \AdminToolResult\ and \AdminToolsResultCodeEnum\ to standardize command execution results, and adds utility classes such as \MQAdminUtils\ and \CommandUtil\ to assist with cluster and topic route metadata retrieval. These additions enable more consistent and structured interactions with the broker cluster for administrative tasks.

tools · high confidence

Introduce new client-side exception and implementation classes

The client module now includes new exception classes (MQBrokerException, MQClientException, OffsetNotFoundException, RequestTimeoutException) and implementation classes (CommunicationMode, FindBrokerResult, MQAdminImpl, MQClientAPIImpl, MQClientManager) that define the core client runtime behavior, error handling, and API interactions. These changes provide a more robust error reporting mechanism for broker and client errors, introduce a communication mode enum, and implement the primary client-side API and manager components.

client/src/main/java/org/apache/rocketmq/client/impl · high confidence

Introduce new producer client API and transactional message support

The client module now includes a new \DefaultMQProducer\ class and \MQProducer\ interface that provide a unified entry point for sending messages, including support for asynchronous sends, batch sending, and request-response (RPC) patterns. A new \TransactionMQProducer\ class and \TransactionListener\ interface enable the sending and management of transactional (half) messages, allowing users to execute local transactions and handle transaction state checks. The update also introduces \SendResult\ and \TransactionSendResult\ classes to encapsulate the outcomes of send operations, and adds \LocalTransactionState\ to represent the status of transactional messages.

client/src/main/java/org/apache/rocketmq/client/producer · high confidence

Introduce page-cache transfer classes for message retrieval

The broker now includes three new classes in the \org.apache.rocketmq.broker.pagecache\ package: \ManyMessageTransfer\, \OneMessageTransfer\, and \QueryMessageTransfer\. These classes implement Netty's \FileRegion\ interface to handle the transfer of message data from the broker's page cache to the network channel, improving how message retrieval operations are executed.

broker/src/main/java/org/apache/rocketmq/broker/pagecache · high confidence

Introduce pop consumption support for push consumers

Added new classes to support the 'pop' consumption mode for push consumers. This includes a new \ConsumeMessagePopOrderlyService\ to handle orderly consumption via the pop protocol, a \PopProcessQueue\ to track state for pop-based message processing, and a \PopRequest\ to encapsulate pop-related request data. These changes enable push consumers to utilize the pop protocol for message retrieval.

client/src/main/java/org/apache/rocketmq/client/impl/consumer · high confidence

Introduce queue-based transactional message processing components

Added new classes to the broker's transactional message handling: DefaultTransactionalMessageCheckListener, GetResult, MessageQueueOpContext, TransactionalMessageBridge, TransactionalMessageServiceImpl, TransactionalMessageUtil, and TransactionalOpBatchService. These components implement the queue-based approach for managing half messages, operation logs, and batch processing of transactional messages, replacing or supplementing previous implementations.

broker/src/main/java/org/apache/rocketmq/broker/transaction/queue · high confidence

Introduce static topic queue mapping and RocksDB-backed topic config storage

The broker now supports static topic queue mapping, allowing topics to be mapped to specific queues on specific brokers. This is implemented via new components: \TopicQueueMappingManager\ to manage the mapping details, \TopicQueueMappingCleanService\ to clean up expired mappings, and \TopicRouteInfoManager\ to handle route information updates. Additionally, topic configuration can now be persisted in a RocksDB store via \RocksDBTopicConfigManager\ and \RocksDBLmqTopicConfigManager\, providing an alternative to the default in-memory/config-file-based storage. LMQ (Lightweight Message Queue) topics are handled by specialized managers (\LmqTopicConfigManager\) that ensure consistent behavior for these special topics.

broker/src/main/java/org/apache/rocketmq/broker/topic · high confidence

Introduce the Tiered Storage module

The tiered storage module is introduced, enabling users to offload message data from local disk to cheaper, larger storage mediums to extend message retention at lower cost. This change adds the core implementation files, configuration classes, and build definitions for the tiered store, including the dispatcher, fetcher, and associated utilities.

tieredstore · high confidence

Introduce unified ServiceManager abstraction for proxy service initialization

The proxy module now uses a new ServiceManager interface and factory to initialize proxy services, with separate implementations for local and cluster modes. This change introduces ClusterServiceManager and LocalServiceManager classes that encapsulate the creation and lifecycle management of core proxy services (message, transaction, metadata, relay, and admin services) using a consistent interface, enabling cleaner separation between local broker-integrated and distributed cluster deployments.

proxy/src/main/java/org/apache/rocketmq/proxy/service · high confidence

Introduce zone-based route filtering for NameServer

The NameServer now supports zone-aware routing. A new \ZoneRouteRPCHook\ intercepts route queries and filters broker data to return only the brokers located in the requested zone, provided the \zoneMode\ flag is enabled. This allows clients to receive route information restricted to a specific availability zone, improving traffic isolation and latency for zone-specific deployments.

namesrv/src/main/java/org/apache/rocketmq/namesrv · high confidence

Introduced BrokerOuterAPI for external communication

Added BrokerOuterAPI, a new class that encapsulates the broker's external communication logic, including name server address resolution, RPC client initialization, and thread pool management. This refactors and centralizes the broker's interaction with the rest of the system, providing a cleaner interface for operations like fetching name server addresses and managing the remoting client lifecycle.

broker/src/main/java/org/apache/rocketmq/broker/out · high confidence

Introduces ReplicasManager for broker-controller coordination

Adds the ReplicasManager class to handle broker registration, metadata synchronization, and role management in controller mode. This new component manages the lifecycle of broker replicas, including periodic syncing of controller metadata, handling broker registration and re-registration, and managing the transition between master and slave roles. The implementation includes thread pools for scheduled tasks and background processing, along with state management for the registration process.

broker/src/main/java/org/apache/rocketmq/broker/controller · high confidence

Introduces a new plain-text ACL configuration system

The ACL module now supports a new 'plain' implementation for managing access control lists via YAML configuration files. This change introduces new classes including PlainAccessData, PlainAccessResource, PlainAccessValidator, PlainPermissionChecker, and PlainPermissionManager, which handle parsing, validation, and permission checking for both standard and gRPC requests. Users can now configure ACL rules in plain text files, enabling easier management and debugging of access policies.

acl/src/main/java/org/apache/rocketmq/acl/plain · high confidence

Introduces new POP consumption mode processors and plugin interfaces

The broker now includes new processor classes to support the POP (Pull-Only-Once) consumption mode, including \AckMessageProcessor\ for handling message acknowledgments, \PeekMessageProcessor\ for peeking messages, and \PopBufferMergeService\ to manage in-flight message tracking and offset commits. Additionally, the \BrokerAttachedPlugin\ and \PullMessageResultHandler\ interfaces are introduced to allow external plugins to hook into the broker's lifecycle and pull message processing, enabling custom logic for these operations.

broker/src/main/java/org/apache/rocketmq/broker/processor · high confidence

Introduces new proxy relay service interface and result types

The proxy module now includes a new \ProxyRelayService\ interface and a corresponding \ProxyRelayResult\ wrapper class, alongside supporting types like \RelayData\ and a stub implementation \ClusterProxyRelayService\. This change introduces the core abstraction for relaying proxy requests, providing a structured way to handle asynchronous relay operations and their results.

proxy/src/main/java/org/apache/rocketmq/proxy/service/relay · high confidence

Introduces server-side offset management and new offset storage backends

The broker now supports server-side offset management for pop messages and broadcast consumption modes, enabling the broker to track and reset consumer offsets directly. This change introduces new classes to handle offset persistence, including a base \ConsumerOffsetManager\ and specialized managers for Light Message Queue (LMQ) and RocksDB-backed storage. Additionally, a new \BroadcastOffsetStore\ is added to manage offsets for broadcast consumption, while \ConsumerOrderInfoManager\ tracks order information for pop consumption. These components allow the broker to maintain state for pop, broadcast, and LMQ scenarios, with options for in-memory or RocksDB-based persistence.

broker/src/main/java/org/apache/rocketmq/broker/offset · high confidence

New ChannelHelper utility for channel type detection

A new ChannelHelper class has been added to the proxy module, providing static methods to determine if a Netty Channel is remote and to identify the protocol type (gRPC v2, Remoting, or other) associated with a given channel instance.

proxy/src/main/java/org/apache/rocketmq/proxy/common/channel · medium confidence

New ConsumerStatsManager for tracking consumer metrics

A new ConsumerStatsManager class has been added to the RocketMQ client, introducing structured tracking of consumer-side metrics. This component monitors pull and consume performance by recording round-trip times (RT) and throughput (TPS) for both pull and consume operations, as well as tracking failed message counts. It exposes a ConsumeStatus object containing aggregated statistics for a specific consumer group and topic, enabling users to monitor and debug consumer health and performance.

client/src/main/java/org/apache/rocketmq/client/stat · high confidence

New KVConfigManager for namespace-based configuration storage

The Namesrv module now includes a new KVConfigManager class that manages key-value configuration data organized by namespace. This manager handles loading, persisting, and querying these configurations using a thread-safe in-memory map, with serialization to JSON via a new KVConfigSerializeWrapper helper class. This change introduces a new mechanism for storing and retrieving namespace-scoped configuration items, which may affect how the name server persists and accesses its configuration data.

namesrv/src/main/java/org/apache/rocketmq/namesrv/kvconfig · high confidence

New RPC request-response example implementations

Added three new Java example files in the RPC package: AsyncRequestProducer, RequestProducer, and ResponseConsumer. These demonstrate the request-response messaging pattern, showing how to send synchronous and asynchronous requests and handle replies using the MessageUtil helper for creating reply messages.

example/src/main/java/org/apache/rocketmq/example/rpc · high confidence

New SQL schema for transaction tracking

A new SQL script, transaction.sql, has been added to the broker resources, defining a t\_transaction table to store transaction offsets and producer groups. This introduces a new database schema component for transaction management within the broker module.

broker/src/main/resources · high confidence

New broker metrics for consumer lag, pop, and request tracking

The broker now exposes a comprehensive set of OpenTelemetry-based metrics for monitoring broker health and consumer behavior. This includes consumer lag and latency gauges, producer and consumer connection counts, message throughput counters, and histogram data for message sizes. Additionally, specific metrics are provided for the Pop consumption mode, tracking revive service lag, latency, and buffer sizes. These metrics are exported via OpenTelemetry (gRPC, Prometheus, or logging) and include labels for cluster, node, topic, consumer group, and message type, enabling detailed observability into broker performance and consumer processing states.

broker/src/main/java/org/apache/rocketmq/broker/metrics · high confidence

New client connection management and housekeeping for consumers and producers

The broker introduces a new client connection management model, replacing the previous monolithic client channel tables with dedicated managers for consumers and producers. A new \ClientChannelInfo\ class encapsulates channel metadata, while \ConsumerManager\ and \ProducerManager\ handle registration, unregistration, and expiration of client connections. A \ClientHousekeepingService\ runs a scheduled task to scan for and remove expired or inactive channels. Additionally, a \ConsumerIdsChangeListener\ interface and its default implementation are added to notify subscribers when consumer group membership changes, enabling real-time or batched notifications about consumer topology changes.

broker/src/main/java/org/apache/rocketmq/broker/client · high confidence

New client error codes and nameserver access configuration

The client module now includes a new ClientErrorCode class defining specific error codes for broker connection, timeout, and topic lookup failures, alongside a new NameserverAccessConfig class to manage nameserver address and domain settings.

client/src/main/java/org/apache/rocketmq/client/common · high confidence

New client-side hooking framework for message lifecycle events

The client module now includes a new hooking framework in the \org.apache.rocketmq.client.hook\ package, introducing a standardized way to intercept and observe message lifecycle events. This change adds several new context classes (\SendMessageContext\, \ConsumeMessageContext\, \EndTransactionContext\, \FilterMessageContext\, \CheckForbiddenContext\) and their corresponding hook interfaces (\SendMessageHook\, \ConsumeMessageHook\, \EndTransactionHook\, \FilterMessageHook\, \CheckForbiddenHook\). These components allow developers to register custom logic that runs before and after sending, consuming, or filtering messages, as well as during transaction management and forbidden checks.

client/src/main/java/org/apache/rocketmq/client/hook · high confidence

New configuration files for controller mode deployments

Added configuration files for controller mode, including standalone, 3-node cluster, namesrv-plugin, and quick-start setups. These files provide ready-to-use templates for deploying the controller in various topologies, enabling users to quickly set up controller-based architectures.

distribution/conf/controller · high confidence

New consumer interfaces and implementations for pull-based consumption

The client library introduces a new \LitePullConsumer\ interface and its \DefaultLitePullConsumer\ implementation, providing a simplified, modern API for actively pulling messages. This includes support for manual queue assignment, polling, and auto-committing offsets. Alongside this, the older \DefaultMQPullConsumer\ and \DefaultMQPushConsumer\ classes are introduced (or re-introduced in the diff) with \@Deprecated\ annotations, signaling a shift away from the legacy pull/push models toward the new lightweight consumer. Additionally, supporting types such as \AckCallback\, \AckResult\, \AckStatus\, and \AllocateMessageQueueStrategy\ are added to facilitate acknowledgment handling and queue allocation strategies.

client/src/main/java/org/apache/rocketmq/client/consumer · high confidence

New consumer listener interfaces and context classes for concurrent and orderly consumption

The client library now provides dedicated interfaces and context classes for message consumption: \MessageListenerConcurrently\ and \MessageListenerOrderly\ define how consumers handle messages in parallel or sequential order respectively. These are supported by new context classes (\ConsumeConcurrentlyContext\, \ConsumeOrderlyContext\) that allow consumers to manage retry delays, acknowledgment indices, and queue suspension. Status enums (\ConsumeConcurrentlyStatus\, \ConsumeOrderlyStatus\, \ConsumeReturnType\) are introduced to report consumption outcomes, enabling users to control retry behavior and handle exceptions or timeouts explicitly.

client/src/main/java/org/apache/rocketmq/client/consumer/listener · high confidence

New distribution packaging and configuration files for dledger and release assemblies

The distribution module now includes new configuration files for dledger (broker-n0.conf, broker-n1.conf, broker-n2.conf) and new release assembly descriptors (release.xml, release-client.xml) that bundle the full set of components (broker, tools, client, namesrv, example, openmessaging, controller) along with logback configuration files for each component. Additionally, LICENSE-BIN and NOTICE-BIN files are added to the distribution directory, providing the Apache License 2.0 text and copyright notices for the binary distribution.

distribution · high confidence

New examples for scheduled and timer-based 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 time delay using the standard delay time level mechanism. TimerMessageProducer and TimerMessageConsumer demonstrate the newer timer-based scheduling API (available in RocketMQ 5.0+), allowing messages to be delivered at a specific future timestamp or after a specified delay in milliseconds. These examples help users understand how to implement delayed and scheduled message patterns in their applications.

example/src/main/java/org/apache/rocketmq/example/schedule · high confidence

New exception and configuration classes for broker and controller management

The common module introduces several new classes to support broker and controller operations. A new \AbortProcessException\ allows broker hooks to immediately return error responses to clients. Configuration classes \BrokerConfig\ and \ControllerConfig\ are added to manage broker and controller settings respectively, including thread pool sizes, queue capacities, and metrics export options. Additionally, \BrokerIdentity\ and \BrokerConfigSingleton\ provide identity management and singleton access for broker configurations. Other new classes include \AclConfig\ for access control settings, \ConfigManager\ for loading and persisting configuration, \BoundaryType\ for boundary operations, \CountDownLatch2\ for synchronization, \KeyBuilder\ for retry topic construction, and \LockCallback\ for lock operations.

common/src/main · high confidence

New gRPC v2 client activity implementation

Added ClientActivity.java in the proxy's gRPC v2 client package, implementing the core client lifecycle (heartbeat, termination, telemetry) and registration of producer/consumer listeners for the new gRPC v2 protocol.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/client · high confidence

New heartbeat syncer for proxy cluster coordination

The proxy now includes a new \HeartbeatSyncer\ component (along with supporting classes like \ClusterConsumerManager\ and \HeartbeatSyncerData\) that synchronizes consumer registration and unregistration events across the proxy cluster. This ensures that all proxy nodes in a cluster maintain a consistent view of connected clients and their subscription states, facilitating better cluster-wide management and routing.

proxy/src/main/java/org/apache/rocketmq/proxy/service/sysmessage · high confidence

New in-memory transaction data storage and cluster-aware transaction service

The proxy now stores transaction state in memory using a new \TransactionDataManager\ and \TransactionData\ model, allowing the proxy to track and manage transactional messages locally. A new \AbstractTransactionService\ and its \ClusterTransactionService\ implementation handle transaction subscriptions, heartbeat scanning, and end-transaction request generation, enabling the proxy to coordinate with brokers across clusters for reliable transactional messaging.

proxy/src/main/java/org/apache/rocketmq/proxy/service/transaction · high confidence

New local file-based consumer offset storage

The client now includes a new local file-based offset store implementation (LocalFileOffsetStore) alongside the existing remote broker-based store. This allows consumers to persist and retrieve their consumption offsets to and from local files, providing a fallback or alternative storage mechanism for consumer state.

client/src/main/java/org/apache/rocketmq/client/consumer/store · high confidence

New long-polling and hold-request classes for the broker

Added new classes in the broker's longpolling package to support message delivery and hold mechanisms: LmqPullRequestHoldService, ManyPullRequest, NotificationRequest, PollingHeader, PollingResult, PopLongPollingService, PullRequest, and PullRequestHoldService. These files introduce the data structures and services that manage long-polling requests, hold requests, and notification requests for the broker.

broker/src/main/java/org/apache/rocketmq/broker/longpolling · high confidence

New message queue allocation strategies for consumer rebalancing

The client now provides several new implementations for distributing message queues among consumers, enabling more flexible rebalancing strategies. These include average hashing (AllocateMessageQueueAveragely), circular average hashing (AllocateMessageQueueAveragelyByCircle), configuration-based allocation (AllocateMessageQueueByConfig), and machine-room-aware allocation (AllocateMessageQueueByMachineRoom). These strategies are built on a new abstract base class (AbstractAllocateMessageQueueStrategy) that enforces parameter validation and uses an internal logger.

client/src/main/java/org/apache/rocketmq/client/consumer/rebalance · high confidence

New message service interface and local remoting command for proxy messaging

The proxy module introduces a new \MessageService\ interface that defines the contract for core messaging operations including sending, popping, pulling, acknowledging, and managing consumer offsets and locks. This interface is supported by new helper classes \LocalRemotingCommand\ and \ReceiptHandleMessage\ within the \proxy/service/message\ package, establishing the foundation for the v2 proxy messaging processor.

proxy/src/main/java/org/apache/rocketmq/proxy/service/message · high confidence

New message tracing context and hook interfaces for send/consume operations

The broker now exposes new tracing infrastructure in the \broker.mqtrace\ package. \SendMessageContext\ and \ConsumeMessageContext\ classes carry detailed metadata for message send and consume operations, including namespace, topic, queue, broker address, and commercial/accounting statistics. Corresponding \SendMessageHook\ and \ConsumeMessageHook\ interfaces allow external code to intercept and observe these events before and after the respective broker operations, enabling deeper observability and custom monitoring.

broker/src/main/java/org/apache/rocketmq/broker/mqtrace · high confidence

New proxy channel and admin service interfaces for connection management

The proxy module introduces new interfaces and utilities to support client connection management and administrative operations. In the channel processor package, \ChannelProtocolType\ defines supported protocols (gRPC v1/v2, Remoting), while \ChannelExtendAttributeGetter\, \RemoteChannelConverter\, and \RemoteChannelSerializer\ provide the structure and serialization logic for remote channel data. Additionally, the \AdminService\ interface is added to the admin service package, exposing methods to check topic existence and create topics on brokers, enabling new administrative capabilities within the proxy layer.

proxy/src/main/java/org/apache/rocketmq/proxy/processor/channel, proxy/src/main/java/org/apache/rocketmq/proxy/service/admin · high confidence

New proxy processor architecture for message handling

The proxy module introduces a new processor-based architecture for handling messaging operations. A new \MessagingProcessor\ interface and its \DefaultMessagingProcessor\ implementation provide a unified entry point for producer, consumer, transaction, and client management. Specific processors (\ProducerProcessor\, \ConsumerProcessor\, \TransactionProcessor\, \ClientProcessor\, \ReceiptHandleProcessor\) encapsulate their respective logic, delegating to a \ServiceManager\ for core services. This refactors the previous monolithic processor structure into a modular design, introducing classes like \AbstractProcessor\ for common functionality and \BatchAckResult\ for batch acknowledgment results.

proxy/src/main/java/org/apache/rocketmq/proxy/processor · high confidence

New remoting proxy layer with protocol handlers and activity processors

The proxy now includes a new remoting layer that handles Netty-based remoting protocols. This adds a ProtocolHandler interface and implementation to manage protocol matching and configuration, a RequestPipeline interface to chain request processing, and specific activity processors (e.g., for sending, pulling, and popping messages, as well as transaction and consumer management). The implementation also introduces a MultiProtocolTlsHelper to manage TLS/SSL contexts for secure communication, supporting both OpenSSL and JDK providers, and a RemotingConverter for message serialization. These changes enable the proxy to directly process and route remoting commands to the appropriate internal processors.

proxy/src/main/java/org/apache/rocketmq/proxy/remoting · high confidence

New srvutil module with file watching and server utilities

Added the srvutil module containing the AclFileWatchService for monitoring ACL file changes, a ServerUtil class for command-line argument parsing, and a ShutdownHookThread for JVM shutdown callbacks. The module also includes a Bazel build configuration and associated unit tests.

srvutil · high confidence

New tracing examples for message tracking and OpenTracing

Added example implementations for message tracing and OpenTracing integration. The \tracemessage\ package now includes \TraceProducer\ and \TracePushConsumer\ for standard message tracking, as well as \OpenTracingProducer\, \OpenTracingPushConsumer\, and \OpenTracingTransactionProducer\ which demonstrate how to integrate RocketMQ with OpenTracing (specifically Jaeger) for distributed tracing.

example/src/main/java/org/apache/rocketmq/example/tracemessage · high confidence

New transactional message check service and listener abstraction

The broker now introduces a dedicated transactional message check service that periodically scans for uncommitted or unrolled-back half messages and sends check-back requests to producers to obtain transaction status. This is supported by a new abstract listener interface for handling check and discard logic, a service interface for managing prepare, commit, and rollback operations, and a service thread that drives the periodic checking process.

broker/src/main/java/org/apache/rocketmq/broker/transaction · high confidence

New utility classes for exception handling, future chaining, and proxy constants

Added four new utility classes in the proxy common package: ExceptionUtils for unwrapping and formatting exceptions, FutureUtils for chaining and error-propagating CompletableFuture operations, FilterUtils for tag matching logic, and ProxyUtils for shared constants like MAX\_MSG\_NUMS\_FOR\_POP\_REQUEST. These provide reusable helpers for error handling, asynchronous flow control, and configuration constants within the proxy module.

proxy/src/main/java/org/apache/rocketmq/proxy/common/utils · high confidence

New utility classes for message validation and timing

Added HookUtils and PositiveAtomicCounter to the broker's utility package. HookUtils introduces pre-flight checks for message submission, including validation of topic length (capped at 255 bytes), body presence, and write availability. It also handles scheduling logic for timer and delay messages, clearing temporary timer properties and transforming messages for the wheel timer. PositiveAtomicCounter provides a thread-safe counter that masks the sign bit to ensure positive integer values.

broker/src/main/java/org/apache/rocketmq/broker/util · high confidence

Proxy startup now supports command-line configuration for broker, proxy config, and mode

The RocketMQ proxy can now be started with command-line arguments to specify the broker configuration path, proxy configuration path, and proxy mode (local or cluster). This allows users to override configuration settings at startup without modifying files, and the startup process now logs the active configuration for verification.

proxy/src/main/java/org/apache/rocketmq/proxy · high confidence

Remoting module refactored with new interfaces and metrics support

The remoting module has been refactored to introduce new core interfaces including RemotingService, RemotingClient, and RemotingServer, alongside supporting classes like RPCHook, InvokeCallback, and Configuration. Additionally, a new metrics framework has been added to the remoting module, introducing OpenTelemetry-based metrics for RPC latency and request/response codes.

remoting · high confidence

Support forwarding messages to dead-letter queues via gRPC v2

The gRPC v2 proxy now supports forwarding messages to dead-letter queues. A new ForwardMessageToDLQActivity class handles the ForwardMessageToDeadLetterQueueRequest, validating the topic and consumer group, decoding the receipt handle, and invoking the underlying messaging processor to forward the message. This adds a new capability for users to manage failed messages by routing them to a dead-letter queue through the gRPC v2 interface.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/producer · high confidence

Architecture

Refactor message store module with Bazel build and new core interfaces

The store module has been refactored to use the Bazel build system, introducing a new BUILD.bazel file that defines the 'store' and 'tests' targets with their respective dependencies. This change introduces several new core interfaces and classes that define how messages are appended, dispatched, and filtered within the message store. Specifically, the diff adds AppendMessageCallback and AppendMessageResult to handle message serialization and results, CommitLogDispatcher for dispatching messages to build consume queues and indexes, and DispatchRequest to carry message metadata. Additionally, new classes like DefaultMessageFilter, ConsumeQueue, and ConsumeQueueExt are introduced to manage consume queue storage and extension data. These changes represent a significant architectural shift in how the message store handles message persistence and retrieval.

store · high confidence

Refactored NameServer request handling into dedicated processor classes

The NameServer's request handling logic has been reorganized into specific processor classes (ClientRequestProcessor, ClusterTestRequestProcessor, and DefaultRequestProcessor) within the namesrv module. This refactoring separates client-facing route queries, cluster-test operations, and default broker/topic management requests, improving code maintainability and allowing for specialized handling of different request types.

namesrv/src/main/java/org/apache/rocketmq/namesrv/processor · high confidence

Behavioural changes

Add OpenMessaging 0.3.0 compatible producer implementation

The OpenMessaging producer implementation has been updated to be compatible with OMS 0.3.0. This introduces new \AbstractOMSProducer\ and \ProducerImpl\ classes that handle message sending (sync, async, and oneway) and lifecycle management according to the OpenMessaging specification, while mapping internal exceptions to OMS-specific runtime and message format exceptions.

openmessaging/src/main/java/io/openmessaging/rocketmq/producer · medium confidence

Add topic message type validation

A new validation mechanism for topic message types has been introduced to the proxy layer. The change adds a TopicMessageTypeValidator interface and a DefaultTopicMessageTypeValidator implementation that checks whether a message's actual type matches the expected type, throwing a ProxyException if they do not match (with special handling for mixed topic types). This ensures that messages conform to their declared message types before processing.

proxy/src/main/java/org/apache/rocketmq/proxy/processor/validator · high confidence

Fix message queue selection logic for hash and random selectors

The message queue selection logic has been corrected to prevent errors during queue selection. Specifically, the hash-based selector now properly handles negative hash codes (such as Integer.MIN\_VALUE) by ensuring the calculated index is always non-negative, and the random selector's implementation has been refined to ensure consistent random selection. These changes affect how messages are distributed across available queues, impacting the reliability of message delivery.

client/src/main/java/org/apache/rocketmq/client/producer/selector · medium confidence

Introduce DLedger role change handler for broker HA switching

Added DLedgerRoleChangeHandler to manage broker role transitions (leader/follower) and ensure the fallen behind node does not become leader, improving high availability switching in containerized environments.

broker/src/main/java/org/apache/rocketmq/broker/dledger · high confidence

Introduce internal producer and topic publish info abstractions

The client module now exposes the internal \MQProducerInner\ interface and the \TopicPublishInfo\ class within the \org.apache.rocketmq.client.impl.producer\ package. These additions provide a structured way to manage topic routing information, queue selection strategies, and transactional message listener bindings, supporting the underlying producer implementation with dedicated state and queue-selection logic.

client/src/main/java/org/apache/rocketmq/client/impl/producer · high confidence

Introduce v2 gRPC proxy activity layer

The proxy's gRPC v2 implementation is restructured into a new \AbstractMessingActivity\ base class, a \ContextStreamObserver\ interface, and a \DefaultGrpcMessingActivity\ implementation that delegates to specific activity classes (e.g., \ReceiveMessageActivity\, \SendMessageActivity\). This change consolidates message processing logic into a cleaner, more maintainable structure for the v2 gRPC interface.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2 · high confidence

Introduces batch broker unregistration for improved performance

The NameServer now processes broker unregistration requests asynchronously via a new \BatchUnregistrationService\. This service queues unregistration requests and processes them in batches, which is intended to speed up the broker-offline process and reduce memory usage compared to handling each unregistration individually.

namesrv/src/main/java/org/apache/rocketmq/namesrv/routeinfo · high confidence

New gRPC interceptors for context, headers, and exception handling

The proxy's gRPC server now includes several new interceptors that improve request handling and error management. ContextInterceptor attaches metadata to the gRPC context, while HeaderInterceptor extracts and injects remote/local addresses and proxy protocol headers into gRPC headers. GlobalExceptionInterceptor provides centralized exception handling for gRPC calls, ensuring consistent error responses and logging. InterceptorConstants defines metadata keys for various RPC attributes, and RequestMapping provides a mapping of gRPC request types to internal request codes, enabling proper routing and processing of v2 protocol messages.

proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor · high confidence

Refactored proxy route service with new routing classes and Caffeine caching

The proxy's routing logic has been refactored to improve cluster mode address resolution and fault tolerance. A new \AddressableMessageQueue\ class wraps message queues with their broker addresses, enabling the \ClusterTopicRouteService\ to correctly retrieve broker addresses in cluster mode, fixing a bug where the wrong address was previously obtained. The \MessageQueueSelector\ now supports a pipeline-based selection strategy that leverages \MQFaultStrategy\ for latency and reachability filtering. Additionally, the \TopicRouteService\ has been updated to use Caffeine for topic route caching with configurable expiration and refresh intervals, replacing the previous caching mechanism.

proxy/src/main/java/org/apache/rocketmq/proxy/service/route · medium confidence

Restructure and optimize message trace implementation

The message trace module in the client has been restructured and optimized. This includes the introduction of new core classes such as TraceBean, TraceContext, TraceDataEncoder, and TraceView, which handle trace data encoding, decoding, and view generation. The change also removes the SubAfter field from the trace data structure and fixes issues related to incomplete trace data and incorrect client host information.

client/src/main/java/org/apache/rocketmq/client/trace · medium confidence

Test coverage

Add unit tests for RocketMQ client consumer implementation; Added Mockito inline mock maker configuration for tests; Added base test class for proxy service tests; Added test configuration files for logging and mocking; Added test configuration files for the RMQ proxy; Added test configuration for ACL, transactional messaging, and logging; Added test coverage for BrokerFastFailure; Added tests for RocketMQ ACL client-side access control; Added unit tests for BrokerMetricsManager; Added unit tests for DefaultReceiptHandleManager; Added unit tests for EscapeBridge; Added unit tests for FilterUtil and FilterAPI integration; Added unit tests for GrpcClientChannel; Added unit tests for LatencyFaultTolerance; Added unit tests for LocalProxyRelayService and ProxyChannel; Added unit tests for MQClientAPI and ProxyClientRemotingProcessor; Added unit tests for MQClientAPIImpl; Added unit tests for MQClientInstance; Added unit tests for MessageUtil reply message creation; Added unit tests for OpenMessaging components; Added unit tests for PullRequestHoldService; Added unit tests for ReceiptHandleGroup and RenewStrategyPolicy; Added unit tests for RemoteChannel encoding and decoding; Added unit tests for RocketMQ ACL common utilities; Added unit tests for RocketMQ BrokerContainer startup and lifecycle; Added unit tests for RocketMQ client consumers; Added unit tests for ScheduleMessageService; Added unit tests for StatsItemSet and MomentStatsItemSet; Added unit tests for ThreadLocalIndex; Added unit tests for TopicValidator; Added unit tests for attribute parsing and validation; Added unit tests for broker client managers; Added unit tests for broker controller mode registration and replica management; Added unit tests for broker message filtering components; Added unit tests for broker offset management components; Added unit tests for broker startup, controller, and path configuration; Added unit tests for broker topic and subscription group managers; Added unit tests for common module components; Added unit tests for compression and pull system flags; Added unit tests for consumer offset storage; Added unit tests for gRPC v2 ClientActivity; Added unit tests for gRPC v2 consumer activities; Added unit tests for gRPC v2 route activity; Added unit tests for message handling components; Added unit tests for message queue allocation strategies; Added unit tests for message queue selection logic; Added unit tests for message trace functionality; Added unit tests for page cache transfer classes; Added unit tests for proxy admin and heartbeat syncer; Added unit tests for proxy configuration and metric collector mode; Added unit tests for proxy message services; Added unit tests for proxy routing and metadata services; Added unit tests for proxy startup and command-line argument parsing; Added unit tests for the ACL plain access validator and permission manager; Added unit tests for the EndTransaction activity; Added unit tests for the RocketMQ NameServer components; Added unit tests for the RocketMQ client producer module; Added unit tests for the Validators class; Added unit tests for the message filtering engine; Added unit tests for the proxy remoting channel management; Added unit tests for transactional message handling; Added unit tests for utility classes; Adds test infrastructure for message tracking and pop consumption.

Dependencies

Maven build structure refactored into multi-module project

The project has been restructured from a single-module build into a multi-module Maven project. This change introduces separate \pom.xml\ files for each component (such as \acl\, \broker\, \client\, \common\, \container\, \controller\, \distribution\, \example\, \filter\, \namesrv\, \openmessaging\, and the root \pom.xml\), each declaring their specific dependencies and build configurations. This modularization allows for more granular dependency management and independent builds of individual RocketMQ components.

(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 54 → 56 (+2.0)
  • Rubric changed (rubric-2026.08.19 → rubric-2026.09.15) — scores are not directly comparable.

Lenses

  • Code Health 73 → 83 (+10.6)
  • Architecture 99 → 99 (-0.2)
  • Maturity 67 → 66 (-0.5)
  • Readiness 59 → 60 (+1.6)
  • Security 38 → 40 (+1.8)
  • Domain Modelling 100 → 100 (+0.0)

Resolved (503)

  • 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)
  • Change coupling: BrokerHeartbeatManager.java ↔ ReplicasInfoManager.java (controller/src/main/java/org/apache/rocketmq/controller/BrokerHeartbeatManager.java)
  • Change coupling: Controller.java ↔ BrokerReplicaInfo.java (controller/src/main/java/org/apache/rocketmq/controller/Controller.java)
  • Change coupling: Controller.java ↔ ReplicasInfoManager.java (controller/src/main/java/org/apache/rocketmq/controller/Controller.java)
  • Change coupling: ControllerManager.java ↔ BrokerReplicaInfo.java (controller/src/main/java/org/apache/rocketmq/controller/ControllerManager.java)
  • Change coupling: DLedgerController.java ↔ BrokerReplicaInfo.java (controller/src/main/java/org/apache/rocketmq/controller/impl/DLedgerController.java)
  • Change coupling: DefaultMessagingProcessor.java ↔ MessageService.java (proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java)
  • Change coupling: MessagingProcessor.java ↔ MessageService.java (proxy/src/main/java/org/apache/rocketmq/proxy/processor/MessagingProcessor.java)
  • Change coupling: PeekMessageProcessor.java ↔ PopMessageProcessor.java (broker/src/main/java/org/apache/rocketmq/broker/processor/PeekMessageProcessor.java)
  • Change coupling: ReplicasInfoManager.java ↔ ControllerRequestProcessor.java (controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.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
  • …and 483 more

New (1135)

  • Base-context workflow trigger runs with an unscoped token
  • Change coupling: CommitLog.java ↔ StoreCheckpoint.java (store/src/main/java/org/apache/rocketmq/store/CommitLog.java)
  • Change coupling: ConsumeQueue.java ↔ BatchConsumeQueue.java (store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java)
  • Change coupling: ConsumerOffsetManager.java ↔ AdminBrokerProcessor.java (broker/src/main/java/org/apache/rocketmq/broker/offset/ConsumerOffsetManager.java)
  • Change coupling: DefaultMessageStore.java ↔ AbstractPluginMessageStore.java (store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java)
  • Change coupling: DefaultMessageStore.java ↔ BatchConsumeQueue.java (store/src/main/java/org/apache/rocketmq/store/DefaultMessageStore.java)
  • Change coupling: KVConfigManager.java ↔ DefaultRequestProcessor.java (namesrv/src/main/java/org/apache/rocketmq/namesrv/kvconfig/KVConfigManager.java)
  • Change coupling: NettyDecoder.java ↔ NettyEncoder.java (remoting/src/main/java/org/apache/rocketmq/remoting/netty/NettyDecoder.java)
  • Change coupling: RemotingClient.java ↔ NettyRemotingAbstract.java (remoting/src/main/java/org/apache/rocketmq/remoting/RemotingClient.java)
  • Change-coupling hub: ReplicasInfoManager.java → Controller.java, ControllerManager.java, ControllerRequestProcessor.java (controller/src/main/java/org/apache/rocketmq/controller/impl/manager/ReplicasInfoManager.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: BrokerOuterAPI (broker/src/main/java/org/apache/rocketmq/broker/out/BrokerOuterAPI.java)
  • ClassTooLong: BrokerStatsManager (store/src/main/java/org/apache/rocketmq/store/stats/BrokerStatsManager.java)
  • ClassTooLong: CommitLog (store/src/main/java/org/apache/rocketmq/store/CommitLog.java)
  • ClassTooLong: CompactionLog (store/src/main/java/org/apache/rocketmq/store/kv/CompactionLog.java)
  • ClassTooLong: ConsumeQueue (store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java)
  • …and 1115 more

Architecture

  • Containers 0 added · 0 removed · contexts 14 added · 0 removed · edges 40 added · 0 removed

Added bounded contexts (14)

  • rocketmq-acl
  • 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 (40)

  • rocketmq-acl → rocketmq-common
  • rocketmq-acl → rocketmq-remoting
  • rocketmq-acl → rocketmq-srvutil
  • rocketmq-broker → rocketmq-acl
  • 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-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 20 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

yuanzhongqiao/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 22 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 faae64715d917bb5d64b8d72581172d26ebe9501 — 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-821afab8930d.