{"record":{"id":"02b2cd8c19ed52cf","repo":"apache/rocketmq","slug":"the-broker-brokername-brokerversion-does-not","errorCode":null,"errorMessage":"The broker[{brokerName}, {brokerVersion}] does not upgrade to support for filter message by {expressionType}","messagePattern":"The broker\\[(.+?), (.+?)\\] does not upgrade to support for filter message by (.+?)","errorType":"exception","errorClass":"MQClientException","httpStatus":null,"severity":"error","filePath":"client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java","lineNumber":211,"sourceCode":"        final PullCallback pullCallback\n    ) throws MQClientException, RemotingException, MQBrokerException, InterruptedException {\n        FindBrokerResult findBrokerResult =\n            this.mQClientFactory.findBrokerAddressInSubscribe(this.mQClientFactory.getBrokerNameFromMessageQueue(mq),\n                this.recalculatePullFromWhichNode(mq), false);\n        if (null == findBrokerResult) {\n            this.mQClientFactory.updateTopicRouteInfoFromNameServer(mq.getTopic());\n            findBrokerResult =\n                this.mQClientFactory.findBrokerAddressInSubscribe(this.mQClientFactory.getBrokerNameFromMessageQueue(mq),\n                    this.recalculatePullFromWhichNode(mq), false);\n        }\n\n\n        if (findBrokerResult != null) {\n            {\n                // check version\n                if (!ExpressionType.isTagType(expressionType)\n                    && findBrokerResult.getBrokerVersion() < MQVersion.Version.V4_1_0_SNAPSHOT.ordinal()) {\n                    throw new MQClientException(\"The broker[\" + mq.getBrokerName() + \", \"\n                        + findBrokerResult.getBrokerVersion() + \"] does not upgrade to support for filter message by \" + expressionType, null);\n                }\n            }\n            int sysFlagInner = sysFlag;\n\n            if (findBrokerResult.isSlave()) {\n                sysFlagInner = PullSysFlag.clearCommitOffsetFlag(sysFlagInner);\n            }\n\n            PullMessageRequestHeader requestHeader = new PullMessageRequestHeader();\n            requestHeader.setConsumerGroup(this.consumerGroup);\n            requestHeader.setTopic(mq.getTopic());\n            requestHeader.setQueueId(mq.getQueueId());\n            requestHeader.setQueueOffset(offset);\n            requestHeader.setMaxMsgNums(maxNums);\n            requestHeader.setSysFlag(sysFlagInner);\n            requestHeader.setCommitOffset(commitOffset);\n            requestHeader.setSuspendTimeoutMillis(brokerSuspendMaxTimeMillis);","sourceCodeStart":193,"sourceCodeEnd":229,"githubUrl":"https://github.com/apache/rocketmq/blob/293f5885719fc4aa3619446a1900f58ccfcfdd29/client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java#L193-L229","documentation":"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).","triggerScenarios":"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.","commonSituations":"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.","solutions":["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","If upgrade is not possible, switch the subscription to TAG filtering (MessageSelector.byTag or a tag expression), which all broker versions support","During rolling upgrades, hold off enabling SQL92 filters until every broker reports >= 4.1.0 (check brokerVersionDesc via admin tools)"],"exampleFix":"// before\nconsumer.subscribe(topic, MessageSelector.bySql(\"amount = 100\")); // broker 4.0.x -> throws on pull\n// after (compat fallback)\nconsumer.subscribe(topic, MessageSelector.byTag(\"pay\")); // tag filtering works pre-4.1.0\n// or: upgrade brokers to >= 4.1.0 then keep the SQL92 selector","handlingStrategy":"validation","validationCode":"// after start, before relying on SQL92 pulls:\nfor (MessageQueue q : consumer.fetchMessageQueues(topic)) {\n    MQClientInstance i = consumer.getDefaultMQPushConsumerImpl().getRebalanceImpl().getmQClientFactory();\n    // simpler: use admin API to assert broker versions >= 4.1.0\n}\nDefaultMQAdminExt admin = new DefaultMQAdminExt(); admin.start();\n// check brokerVersionDesc in admin.examineBrokerClusterInfo(); require >= 4.1.0 for SQL92","typeGuard":null,"tryCatchPattern":"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)); } }","preventionTips":["Pin a minimum broker version in deployment checks when using SQL92","Complete rolling upgrades before enabling expression filters","Keep a tag-based fallback selector in configuration for mixed fleets"],"tags":["rocketmq","consumer","broker-version","sql92-filter","compatibility"],"backgroundTag":null,"analyzedSha":"293f5885719fc4aa3619446a1900f58ccfcfdd29","analyzedAt":"2026-08-14T11:50:13.822Z","schemaVersion":2},"datasetVersion":"2026-08-15T22:17:37.221Z"}