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
- 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
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
- 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
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
- API-01
- API-01
- Both channel and connection closing failed. Logging channel…
- Cannot specify both ' ' and root-level ' '.
- CorrelationId is missing but required, rejecting message tag
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)