apache/seatunnel · warning

Consumer for queue already exists, skipping

Error message

Consumer for queue {} already exists, skipping

What it means

When the RabbitMQ source reader receives splits (queue assignments) via addSplits, it skips any split whose queue already has an active consumer registered in activeConsumers. This warning indicates a duplicate split assignment for the same queue, which is ignored to prevent creating duplicate consumers that would leak resources or double-consume.

Solutions

  1. Check why the enumerator re-assigned an already-active split (failover/restoration path) and whether deduplication belongs in the enumerator
  2. Deduplicate the splits list before calling addSplits
  3. If persistent, verify reader/enumerator split state consistency after recovery; update the connector if the enumerator duplicates splits
Defensive patterns

Strategy: validation

Validate before calling

List<RabbitmqSplit> deduped = splits.stream().collect(Collectors.collectingAndThen(toMap(RabbitmqSplit::splitId, s -> s, (a, b) -> a), m -> new ArrayList<>(m.values())));

Prevention

When it happens

Trigger: addSplits is called with a RabbitmqSplit whose splitId() is already a key in activeConsumers — e.g. the enumerator re-assigned the same queue after recovery/restoration, or the split list contains duplicates.

Common situations: Split restoration after task failover re-delivering already-active splits; enumerator bug producing duplicate splits; tests explicitly covering consumer resource-leak/duplicate behavior.

Understand the failure class

Background: "Invalid state transition" errors: "status must be X, actually Y", "already rejected/charging/uninstalled", "cannot ... while running" — what they mean when a library rejects your call — this error's family across 31 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/5b910023e7b49862. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java:209

        // Bounded mode logic: Stop the job if all splits have been consumed and the queue is empty
        if (Boundedness.BOUNDED.equals(context.getBoundedness()) && noMoreSplitsAssigned) {
            if (message == null && queue.isEmpty()) {
                log.info(
                        "No more splits assigned, queue is empty, and polling timed out. Signaling end of input.");
                context.signalNoMoreElement();
            }
        }
    }

    @Override
    public void addSplits(List<RabbitmqSplit> splits) {
        // Dynamically start consuming from newly assigned queues (splits)
        for (RabbitmqSplit split : splits) {
            log.info("Received split for queue: {}", split.splitId());
            try {
                if (activeConsumers.containsKey(split.splitId())) {
                    log.warn("Consumer for queue {} already exists, skipping", split.splitId());
                    continue;
                }

                // Create a new consumer that feeds messages into the shared internal 'queue'
                DefaultConsumer consumer =
                        rabbitMQClient.getQueueingConsumer(queue, split.splitId());
                rabbitMQClient.setupQueue(split.splitId());
                channel.basicConsume(split.splitId(), autoAck, consumer);
                activeConsumers.put(split.splitId(), consumer);
                sourceSplits.add(split);

                log.info("Started consuming from queue: {}", split.splitId());
            } catch (IOException e) {
                throw new RabbitmqConnectorException(
                        org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception
                                .RabbitmqConnectorErrorCode.CREATE_RABBITMQ_CLIENT_FAILED,
                        e);
            }

View on GitHub (pinned to cf67b549a7)