apache/rocketmq · error · MQClientException

The broker[ , ] does not upgrade to support for filter…

Error message

The broker[{brokerName}, {brokerVersion}] does not upgrade to support for filter message by {expressionType}

What it means

Thrown by PullAPIWrapper.pullKernelImpl when a pull request uses a non-tag expression type (SQL92 or class filter) but the located broker's version is below V4_1_0_SNAPSHOT. Server-side expression filtering was only added to brokers in 4.1.0, so older brokers cannot evaluate the filter; the client refuses the pull rather than shipping unfiltered results (which would violate subscription semantics).

Solutions

  1. Upgrade all brokers in the cluster to 4.1.0 or later (preferably one consistent version), then retest — the client picks up versions from heartbeats/route data
  2. If upgrade is not possible, switch the subscription to TAG filtering (MessageSelector.byTag or a tag expression), which all broker versions support
  3. During rolling upgrades, hold off enabling SQL92 filters until every broker reports >= 4.1.0 (check brokerVersionDesc via admin tools)

Example fix

// before
consumer.subscribe(topic, MessageSelector.bySql("amount = 100")); // broker 4.0.x -> throws on pull
// after (compat fallback)
consumer.subscribe(topic, MessageSelector.byTag("pay")); // tag filtering works pre-4.1.0
// or: upgrade brokers to >= 4.1.0 then keep the SQL92 selector
Defensive patterns

Strategy: validation

Validate before calling

// after start, before relying on SQL92 pulls:
for (MessageQueue q : consumer.fetchMessageQueues(topic)) {
    MQClientInstance i = consumer.getDefaultMQPushConsumerImpl().getRebalanceImpl().getmQClientFactory();
    // simpler: use admin API to assert broker versions >= 4.1.0
}
DefaultMQAdminExt admin = new DefaultMQAdminExt(); admin.start();
// check brokerVersionDesc in admin.examineBrokerClusterInfo(); require >= 4.1.0 for SQL92

Try / catch

try { result = consumer.subscribe(topic, MessageSelector.bySql(sql)); } catch (MQClientException e) { if (e.getMessage().contains("does not upgrade")) { /* fallback to tags */ consumer.subscribe(topic, MessageSelector.byTag(tag)); } }

Prevention

When it happens

Trigger: A push consumer subscribed with MessageSelector.bySql(...) (expressionType = SQL92, or a class filter) pulls from a broker whose reported version ordinal < V4_1_0_SNAPSHOT.ordinal(). Happens as soon as rebalance assigns a queue and the pull thread issues its first request — not at subscribe() time.

Common situations: Mixed-version cluster during a rolling upgrade where one replica is still pre-4.1.0; SQL92 or class filtering introduced to an app deployed against an older broker fleet; downgrade of a broker after the client already cached route data.

Related errors


AI-assisted analysis of apache/rocketmq@293f588571 (2026-08-14). Data as JSON: /api/errors/02b2cd8c19ed52cf. Report an issue: GitHub.

Appendix: source

Thrown at client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java:211

        final PullCallback pullCallback
    ) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {
        FindBrokerResult findBrokerResult =
            this.mQClientFactory.findBrokerAddressInSubscribe(this.mQClientFactory.getBrokerNameFromMessageQueue(mq),
                this.recalculatePullFromWhichNode(mq), false);
        if (null == findBrokerResult) {
            this.mQClientFactory.updateTopicRouteInfoFromNameServer(mq.getTopic());
            findBrokerResult =
                this.mQClientFactory.findBrokerAddressInSubscribe(this.mQClientFactory.getBrokerNameFromMessageQueue(mq),
                    this.recalculatePullFromWhichNode(mq), false);
        }


        if (findBrokerResult != null) {
            {
                // check version
                if (!ExpressionType.isTagType(expressionType)
                    && findBrokerResult.getBrokerVersion() < MQVersion.Version.V4_1_0_SNAPSHOT.ordinal()) {
                    throw new MQClientException("The broker[" + mq.getBrokerName() + ", "
                        + findBrokerResult.getBrokerVersion() + "] does not upgrade to support for filter message by " + expressionType, null);
                }
            }
            int sysFlagInner = sysFlag;

            if (findBrokerResult.isSlave()) {
                sysFlagInner = PullSysFlag.clearCommitOffsetFlag(sysFlagInner);
            }

            PullMessageRequestHeader requestHeader = new PullMessageRequestHeader();
            requestHeader.setConsumerGroup(this.consumerGroup);
            requestHeader.setTopic(mq.getTopic());
            requestHeader.setQueueId(mq.getQueueId());
            requestHeader.setQueueOffset(offset);
            requestHeader.setMaxMsgNums(maxNums);
            requestHeader.setSysFlag(sysFlagInner);
            requestHeader.setCommitOffset(commitOffset);
            requestHeader.setSuspendTimeoutMillis(brokerSuspendMaxTimeMillis);

View on GitHub (pinned to 293f588571)