alibaba/canal · error · CanalClientException

mq get/ack not support concurrent & async ack

Error message

mq get/ack not support concurrent & async ack

What it means

Thrown by CanalRabbitMQConsumer.getMessage when a previous batch (lastGetBatchMessage) has not yet been acknowledged or rolled back. The connector enforces strict get-then-ack ordering: you must ack or rollback the current batch before requesting the next one, because it does not support concurrent or async acknowledgement.

Solutions

  1. Ensure every getMessage is followed by exactly one ack() or rollback() before the next getMessage call.
  2. Serialize getMessage/ack/rollback on a single thread or with an external lock — the connector is not concurrency-safe across these.
  3. Wrap batch processing in try/finally so rollback() runs if processing throws, clearing lastGetBatchMessage.

Example fix

// before
List<CommonMessage> batch = consumer.getMessage(timeout, unit);
process(batch);            // if this throws, next getMessage fails
doSomethingElse();
consumer.getMessage(...); // throws 'mq get/ack not support concurrent & async ack'

// after — ack or rollback before the next get
List<CommonMessage> batch = consumer.getMessage(timeout, unit);
try {
    process(batch);
    consumer.ack();
} catch (Exception e) {
    consumer.rollback();
    throw e;
}
Defensive patterns

Strategy: validation

Validate before calling

if (this.lastGetBatchMessage != null) {
    // ack or rollback the previous batch before getting the next
    rollback();
}

Type guard

boolean readyForNextBatch() {
    return this.lastGetBatchMessage == null;
}

Prevention

When it happens

Trigger: getMessage() is called while this.lastGetBatchMessage != null — i.e. a prior getMessage returned a batch and neither ack nor rollback has cleared the field. Re-entrant or multi-threaded calls to getMessage trip this immediately.

Common situations: Calling getMessage from multiple threads; forgetting to call ack()/rollback() after processing a batch; an exception in processing that skips the ack and the caller loops back to getMessage.

Related errors


AI-assisted analysis of alibaba/canal@87be50e876 (2026-08-14). Data as JSON: /api/errors/7c08b1f93ecbfabb. Report an issue: GitHub.

Appendix: source

Thrown at connector/rabbitmq-connector/src/main/java/com/alibaba/otter/canal/connector/rabbitmq/consumer/CanalRabbitMQConsumer.java:166

            logger.error("Put message to queue error", e);
            throw new RuntimeException(e);
        }
        boolean isCompleted;
        try {
            isCompleted = batchMessage.waitFinish(batchProcessTimeout);
        } catch (InterruptedException e) {
            logger.error("Interrupted when waiting messages to be finished.", e);
            throw new RuntimeException(e);
        }
        boolean isSuccess = batchMessage.isSuccess();
        return isCompleted && isSuccess;
    }

    @Override
    public List<CommonMessage> getMessage(Long timeout, TimeUnit unit) {
        try {
            if (this.lastGetBatchMessage != null) {
                throw new CanalClientException("mq get/ack not support concurrent & async ack");
            }

            ConsumerBatchMessage<CommonMessage> batchMessage = messageBlockingQueue.poll(timeout, unit);
            if (batchMessage != null) {
                this.lastGetBatchMessage = batchMessage;
                return batchMessage.getData();
            }
        } catch (InterruptedException ex) {
            logger.warn("Get message timeout", ex);
            throw new CanalClientException("Failed to fetch the data after: " + timeout);
        }
        return null;
    }

    @Override
    public void rollback() {
        try {
            if (this.lastGetBatchMessage != null) {

View on GitHub (pinned to 87be50e876)