yuanzhongqiao/rocketmq
56.2
Adequate · 22 September 2026
181.6k
lines of production code
Java
primary language
3
measurements over time
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.