apache/seatunnel · critical · RuntimeException

Failed to setup RabbitMQ queue

Error message

Failed to setup RabbitMQ queue

What it means

RabbitmqSinkWriter's constructor throws RuntimeException("Failed to setup RabbitMQ queue") when RabbitmqClient.setupQueue() (queue declare/bind) throws any exception. Note this is a plain RuntimeException, not a SeaTunnel connector exception, and it aborts sink initialization.

Solutions

  1. Read the wrapped cause (`Caused by:`) to see the actual broker error, then fix connectivity, credentials, or permissions accordingly
  2. Confirm queue/exchange declaration parameters match the existing queue (durable/auto-delete/arguments) — PRECONDITION_FAILED means an existing queue differs
  3. Check the user's permissions with rabbitmqctl list_permissions for the vhost
  4. Validate host/port/virtualHost in the sink config and test with a small standalone client

Example fix

// before (config mismatch causing 406 on declare)
queue-name = "orders"
durable = false // existing queue is durable
// after
queue-name = "orders"
durable = true // matches existing declaration
Defensive patterns

Strategy: try-catch

Validate before calling

// before the job: test declare on the target vhost
try (com.rabbitmq.client.Connection c = factory.newConnection();
     com.rabbitmq.client.Channel ch = c.createChannel()) {
  ch.queueDeclarePassive("orders");
  System.out.println("Queue reachable and permitted");
}

Try / catch

try {
  client.setupQueue();
} catch (Exception e) {
  // inspect root cause: connectivity vs permission vs PRECONDITION_FAILED (queue args mismatch)
  Throwable root = e;
  while (root.getCause() != null) root = root.getCause();
  throw new IllegalStateException("Queue setup failed: " + root.getMessage(), e);
}

Prevention

When it happens

Trigger: setupQueue() throws during sink writer construction — e.g. broker unreachable (ConnectException), access-refused for declare, channel already closed, or queue already exists with different parameters (406 PRECONDITION_FAILED).

Common situations: RabbitMQ not running or wrong host/port; user lacks configure permission to declare the queue; queue already exists with different durability/arguments causing an inequality error; vhost typo.

Related errors


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

Appendix: source

Thrown at seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkWriter.java:38

import org.apache.seatunnel.api.table.type.SeaTunnelRow;
import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
import org.apache.seatunnel.connectors.seatunnel.common.sink.AbstractSinkWriter;
import org.apache.seatunnel.connectors.seatunnel.rabbitmq.client.RabbitmqClient;
import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig;
import org.apache.seatunnel.format.json.JsonSerializationSchema;

import java.util.Optional;

public class RabbitmqSinkWriter extends AbstractSinkWriter<SeaTunnelRow, Void> {
    private RabbitmqClient rabbitMQClient;
    private final JsonSerializationSchema jsonSerializationSchema;

    public RabbitmqSinkWriter(RabbitmqConfig config, SeaTunnelRowType seaTunnelRowType) {
        this.rabbitMQClient = new RabbitmqClient(config);
        try {
            this.rabbitMQClient.setupQueue();
        } catch (Exception e) {
            throw new RuntimeException("Failed to setup RabbitMQ queue", e);
        }
        this.jsonSerializationSchema = new JsonSerializationSchema(seaTunnelRowType);
    }

    @Override
    public void write(SeaTunnelRow element) {
        rabbitMQClient.write(jsonSerializationSchema.serialize(element));
    }

    @Override
    public Optional prepareCommit() {
        return Optional.empty();
    }

    @Override
    public void close() {
        if (rabbitMQClient != null) {
            rabbitMQClient.close();

View on GitHub (pinned to cf67b549a7)