{"record":{"id":"5b910023e7b49862","repo":"apache/seatunnel","slug":"consumer-for-queue-already-exists-skipping","errorCode":null,"errorMessage":"Consumer for queue {} already exists, skipping","messagePattern":"Consumer for queue (.+?) already exists, skipping","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java","lineNumber":209,"sourceCode":"\n        // Bounded mode logic: Stop the job if all splits have been consumed and the queue is empty\n        if (Boundedness.BOUNDED.equals(context.getBoundedness()) && noMoreSplitsAssigned) {\n            if (message == null && queue.isEmpty()) {\n                log.info(\n                        \"No more splits assigned, queue is empty, and polling timed out. Signaling end of input.\");\n                context.signalNoMoreElement();\n            }\n        }\n    }\n\n    @Override\n    public void addSplits(List<RabbitmqSplit> splits) {\n        // Dynamically start consuming from newly assigned queues (splits)\n        for (RabbitmqSplit split : splits) {\n            log.info(\"Received split for queue: {}\", split.splitId());\n            try {\n                if (activeConsumers.containsKey(split.splitId())) {\n                    log.warn(\"Consumer for queue {} already exists, skipping\", split.splitId());\n                    continue;\n                }\n\n                // Create a new consumer that feeds messages into the shared internal 'queue'\n                DefaultConsumer consumer =\n                        rabbitMQClient.getQueueingConsumer(queue, split.splitId());\n                rabbitMQClient.setupQueue(split.splitId());\n                channel.basicConsume(split.splitId(), autoAck, consumer);\n                activeConsumers.put(split.splitId(), consumer);\n                sourceSplits.add(split);\n\n                log.info(\"Started consuming from queue: {}\", split.splitId());\n            } catch (IOException e) {\n                throw new RabbitmqConnectorException(\n                        org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception\n                                .RabbitmqConnectorErrorCode.CREATE_RABBITMQ_CLIENT_FAILED,\n                        e);\n            }","sourceCodeStart":191,"sourceCodeEnd":227,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceReader.java#L191-L227","documentation":"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.","triggerScenarios":"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.","commonSituations":"Split restoration after task failover re-delivering already-active splits; enumerator bug producing duplicate splits; tests explicitly covering consumer resource-leak/duplicate behavior.","solutions":["Check why the enumerator re-assigned an already-active split (failover/restoration path) and whether deduplication belongs in the enumerator","Deduplicate the splits list before calling addSplits","If persistent, verify reader/enumerator split state consistency after recovery; update the connector if the enumerator duplicates splits"],"exampleFix":null,"handlingStrategy":"validation","validationCode":"List<RabbitmqSplit> deduped = splits.stream().collect(Collectors.collectingAndThen(toMap(RabbitmqSplit::splitId, s -> s, (a, b) -> a), m -> new ArrayList<>(m.values())));","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Deduplicate splits before addSplits","Verify enumerator does not re-issue active splits after failover","Log and diff split assignments on recovery","Add tests for split restoration paths"],"tags":["rabbitmq","duplicate-splits","recovery"],"backgroundTag":"invalid-state-transition","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}