{"record":{"id":"50f2fa6259131665","repo":"alibaba/canal","slug":"mq-get-ack-not-support-concurrent-async-ack-50f2fa","errorCode":null,"errorMessage":"mq get/ack not support concurrent & async ack","messagePattern":"mq get/ack not support concurrent & async ack","errorType":"exception","errorClass":"CanalClientException","httpStatus":null,"severity":"error","filePath":"connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java","lineNumber":187,"sourceCode":"            logger.error(\"Put message to queue error\", e);\n            throw new RuntimeException(e);\n        }\n        boolean isCompleted;\n        try {\n            isCompleted = batchMessage.waitFinish(batchProcessTimeout);\n        } catch (InterruptedException e) {\n            logger.error(\"Interrupted when waiting messages to be finished.\", e);\n            throw new RuntimeException(e);\n        }\n        boolean isSuccess = batchMessage.isSuccess();\n        return isCompleted && isSuccess;\n    }\n\n    @Override\n    public List<CommonMessage> getMessage(Long timeout, TimeUnit unit) {\n        try {\n            if (this.lastGetBatchMessage != null) {\n                throw new CanalClientException(\"mq get/ack not support concurrent & async ack\");\n            }\n\n            ConsumerBatchMessage<CommonMessage> batchMessage = messageBlockingQueue.poll(timeout, unit);\n            if (batchMessage != null) {\n                this.lastGetBatchMessage = batchMessage;\n                return batchMessage.getData();\n            }\n        } catch (InterruptedException ex) {\n            logger.warn(\"Get message timeout\", ex);\n            throw new CanalClientException(\"Failed to fetch the data after: \" + timeout);\n        }\n        return null;\n    }\n\n    @Override\n    public void rollback() {\n        try {\n            if (this.lastGetBatchMessage != null) {","sourceCodeStart":169,"sourceCodeEnd":205,"githubUrl":"https://github.com/alibaba/canal/blob/87be50e87686a3e8af08c368d0e1ffd1f59eb04a/connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/consumer/CanalRocketMQConsumer.java#L169-L205","documentation":"Thrown by CanalRocketMQConsumer.getMessage when a previous batch (lastGetBatchMessage) has not been acknowledged or rolled back before the next getMessage call. The RocketMQ connector mirrors the RabbitMQ one: it enforces strict sequential get→ack/rollback ordering and does not support concurrent or asynchronous acknowledgement.","triggerScenarios":"getMessage() is called while this.lastGetBatchMessage != null — a prior batch is still outstanding because neither ack() nor rollback() has cleared it. Multi-threaded or re-entrant getMessage calls trip this immediately.","commonSituations":"Processing a batch in a separate thread while the main loop calls getMessage again; an exception between getMessage and ack that skips rollback; forgetting to ack after successful processing.","solutions":["Always call ack() or rollback() before the next getMessage().","Run getMessage/ack/rollback on a single thread or under a lock — the connector is not concurrency-safe.","Use try/finally to rollback on processing failure so lastGetBatchMessage is cleared."],"exampleFix":"// before\nList<CommonMessage> batch = consumer.getMessage(timeout, unit);\nprocess(batch);\nconsumer.getMessage(...); // throws\n\n// after\nList<CommonMessage> batch = consumer.getMessage(timeout, unit);\ntry {\n    process(batch);\n    consumer.ack();\n} catch (RuntimeException e) {\n    consumer.rollback();\n    throw e;\n}","handlingStrategy":"validation","validationCode":"if (this.lastGetBatchMessage != null) {\n    rollback(); // clear the outstanding batch first\n}","typeGuard":"boolean readyForNextBatch() {\n    return this.lastGetBatchMessage == null;\n}","tryCatchPattern":null,"preventionTips":["Always ack() or rollback() before the next getMessage().","Run getMessage/ack/rollback single-threaded.","Use try/finally to rollback on failure."],"tags":["rocketmq","consumer","concurrency","canal-connector","acknowledgement"],"backgroundTag":null,"analyzedSha":"87be50e87686a3e8af08c368d0e1ffd1f59eb04a","analyzedAt":"2026-08-14T04:30:11.918Z","schemaVersion":2},"datasetVersion":"2026-08-14T05:17:29.042Z"}