This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-11576-0d9f9e2303da0b40b20e55fd341d564757e86dc0
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 38a37104e59ee3317cf0bc299c9bfee22149391c
Author: Jast <[email protected]>
AuthorDate: Tue Sep 15 04:16:29 2026 +0000

    [Fix][Connector-V2] Support secure passive RabbitMQ connections (#11576)
    
    Co-authored-by: zhangshenghang <[email protected]>
    Co-authored-by: zhangshenghang <[email protected]>
---
 docs/en/connectors/sink/Rabbitmq.md                | 20 ++++++
 docs/en/connectors/source/Rabbitmq.md              | 20 ++++++
 .../introduction/concepts/incompatible-changes.md  | 13 ++++
 docs/zh/connectors/sink/Rabbitmq.md                | 20 ++++++
 docs/zh/connectors/source/Rabbitmq.md              | 20 ++++++
 .../introduction/concepts/incompatible-changes.md  | 11 +++
 .../seatunnel/rabbitmq/client/RabbitmqClient.java  | 55 +++++++++++---
 .../rabbitmq/config/RabbitmqBaseOptions.java       | 20 ++++++
 .../seatunnel/rabbitmq/config/RabbitmqConfig.java  | 25 ++++++-
 .../exception/RabbitmqConnectorErrorCode.java      |  3 +-
 .../rabbitmq/sink/RabbitmqSinkFactory.java         |  4 ++
 .../rabbitmq/source/RabbitmqSourceFactory.java     |  4 ++
 .../seatunnel/rabbitmq/RabbitmqFactoryTest.java    | 42 +++++++++++
 .../rabbitmq/client/RabbitmqClientTest.java        | 83 ++++++++++++++++++++++
 14 files changed, 328 insertions(+), 12 deletions(-)

diff --git a/docs/en/connectors/sink/Rabbitmq.md 
b/docs/en/connectors/sink/Rabbitmq.md
index 70a507bb9a..bab1147b37 100644
--- a/docs/en/connectors/sink/Rabbitmq.md
+++ b/docs/en/connectors/sink/Rabbitmq.md
@@ -32,6 +32,8 @@ Used to write data to RabbitMQ queues.
 | protobuf_schema            | string  | no       | -             |
 | protobuf_message_name      | string  | no       | -             |
 | url                        | string  | no       | -             |
+| uri                        | string  | no       | -             |
+| ssl                        | boolean | no       | false         |
 | routing_key                | string  | no       | -             |
 | exchange                   | string  | no       | -             |
 | network_recovery_interval  | int     | no       | -             |
@@ -43,6 +45,7 @@ Used to write data to RabbitMQ queues.
 | durable                    | boolean | no       | true          |
 | exclusive                  | boolean | no       | false         |
 | auto_delete                | boolean | no       | false         |
+| passive                    | boolean | no       | false         |
 
 ### host [string]
 
@@ -70,6 +73,16 @@ the password to use when connecting to the broker
 
 convenience method for setting the fields in an AMQP URI: host, port, 
username, password and virtual host
 
+### uri [string]
+
+Legacy alias for `url`. Configure only one of `url` and `uri`.
+
+### ssl [boolean]
+
+Enables SSL/TLS for host-and-port configuration. Use `url` with an `amqps://` 
URI when the URI itself supplies the connection settings.
+
+When `url` uses an `amqps://` URI, the broker certificate is verified against 
the JVM trust store with hostname verification enabled. Connections that 
previously relied on the implicit trust-all behavior with self-signed or 
private-CA certificates must import the broker certificate into the trust 
store, or they will fail to connect.
+
 ### queue_name [string]
 
 the queue to write the message to. If `routing_key` is not configured, the 
connector publishes messages to this queue through the default exchange.
@@ -135,10 +148,17 @@ Sink plugin common parameters, please refer to [Sink 
Common Options](../common-o
 - true: The queue will be deleted automatically when the last consumer 
unsubscribes.
 - false: The queue will not be automatically deleted.
 
+### passive
+
+- false: Declare the queue with the configured durable, exclusive, and 
auto-delete settings.
+- true: Verify that the queue already exists without creating or modifying it. 
Use this for accounts that can publish but cannot declare queues.
+
 
 ## Configuration Notes
 
 - If you configure `username`, you must also configure `password`, and vice 
versa.
+- Configure only one of `url` and `uri`. `uri` is retained for existing 
configurations; use `url` in new configurations.
+- Set `ssl = true` when connecting to an AMQPS endpoint with `host` and `port` 
settings.
 - `host`, `port`, `virtual_host`, and `queue_name` are required connector 
options. `url` can additionally provide the AMQP URI used by the RabbitMQ 
client.
 - `durable`, `exclusive`, and `auto_delete` are used when the connector 
declares the target queue.
 - When `format` is `protobuf`, configure both `protobuf_schema` and 
`protobuf_message_name`.
diff --git a/docs/en/connectors/source/Rabbitmq.md 
b/docs/en/connectors/source/Rabbitmq.md
index 3cd9d38bbe..1c50b4b0ba 100644
--- a/docs/en/connectors/source/Rabbitmq.md
+++ b/docs/en/connectors/source/Rabbitmq.md
@@ -46,6 +46,8 @@ The source must be non-parallel (parallelism set to 1) in 
order to achieve exact
 | protobuf_schema            | string  | no       | -             |
 | protobuf_message_name      | string  | no       | -             |
 | url                        | string  | no       | -             |
+| uri                        | string  | no       | -             |
+| ssl                        | boolean | no       | false         |
 | routing_key                | string  | no       | -             |
 | exchange                   | string  | no       | -             |
 | network_recovery_interval  | int     | no       | -             |
@@ -62,6 +64,7 @@ The source must be non-parallel (parallelism set to 1) in 
order to achieve exact
 | durable                    | boolean | no       | true          |
 | exclusive                  | boolean | no       | false         |
 | auto_delete                | boolean | no       | false         |
+| passive                    | boolean | no       | false         |
 
 ### host [string]
 
@@ -89,6 +92,16 @@ the password to use when connecting to the broker
 
 convenience method for setting the fields in an AMQP URI: host, port, 
username, password and virtual host
 
+### uri [string]
+
+Legacy alias for `url`. Configure only one of `url` and `uri`.
+
+### ssl [boolean]
+
+Enables SSL/TLS for host-and-port configuration. Use `url` with an `amqps://` 
URI when the URI itself supplies the connection settings.
+
+When `url` uses an `amqps://` URI, the broker certificate is verified against 
the JVM trust store with hostname verification enabled. Connections that 
previously relied on the implicit trust-all behavior with self-signed or 
private-CA certificates must import the broker certificate into the trust 
store, or they will fail to connect.
+
 ### queue_name [string]
 
 the queue to consume messages from. *Note: Required if `tables_configs` is not 
configured.*
@@ -186,6 +199,11 @@ Source plugin common parameters, please refer to [Source 
Common Options](../comm
 - true: The queue will be deleted automatically when the last consumer 
unsubscribes.
 - false: The queue will not be automatically deleted.
 
+### passive
+
+- false: Declare the queue with the configured durable, exclusive, and 
auto-delete settings.
+- true: Verify that the queue already exists without creating or modifying it. 
Use this for consumer accounts without queue-declaration permission.
+
 ## Migration Guide & Configuration Rules
 
 If you are upgrading from a previous version that only supported single-table 
reads, your existing configuration will work without any changes.
@@ -197,6 +215,8 @@ If you are upgrading from a previous version that only 
supported single-table re
 - In multi-table mode, put each queue's `schema` inside its own 
`tables_configs` item.
 - When `format` is `protobuf`, configure both `protobuf_schema` and 
`protobuf_message_name` at the same level as the queue configuration.
 - If you configure `username`, you must also configure `password`, and vice 
versa.
+- Configure only one of `url` and `uri`. `uri` is retained for existing 
configurations; use `url` in new configurations.
+- Set `ssl = true` when connecting to an AMQPS endpoint with `host` and `port` 
settings.
 - `host` and `port` are always required. `virtual_host` is optional unless 
your RabbitMQ deployment requires a non-default virtual host.
 
 ## Example
diff --git a/docs/en/introduction/concepts/incompatible-changes.md 
b/docs/en/introduction/concepts/incompatible-changes.md
index 0dc9cb8392..ebe8cebd76 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -5,6 +5,19 @@ You need to check this document before you upgrade to related 
version.
 
 ## dev
 
+### RabbitMQ Connector
+
+- **Breaking Change: `amqps://` connections now verify broker certificates**
+  - **Affected component**: `seatunnel-connectors-v2/connector-rabbitmq`
+  - **Description**: Previously, connecting with an `amqps://` `url`/`uri` 
implicitly installed a
+    trust-all trust manager without hostname verification. Certificate 
verification is now
+    enforced for `amqps://` connections, consistent with the `ssl = true` 
host/port path.
+  - **Impact**: Jobs that connect with `amqps://` URLs to brokers using 
self-signed or private-CA
+    certificates will fail to connect after upgrading.
+  - **Migration Guide**: Import the broker certificate (or your private CA 
chain) into the JVM
+    trust store of the SeaTunnel runtime, or switch to the `host`/`port` + 
`ssl = true`
+    configuration with a properly configured trust store.
+
 ### Zeta REST Pagination Parameter Validation
 
 - **Behavior change: `page` and `rows` are validated on paginated endpoints**
diff --git a/docs/zh/connectors/sink/Rabbitmq.md 
b/docs/zh/connectors/sink/Rabbitmq.md
index e8028f7586..12b9222797 100644
--- a/docs/zh/connectors/sink/Rabbitmq.md
+++ b/docs/zh/connectors/sink/Rabbitmq.md
@@ -33,6 +33,8 @@ import ChangeLog from '../changelog/connector-rabbitmq.md';
 | protobuf_schema            | string  | 否    | -     |
 | protobuf_message_name      | string  | 否    | -     |
 | url                        | string  | 否    | -     |
+| uri                        | string  | 否    | -     |
+| ssl                        | boolean | 否    | false |
 | routing_key                | string  | 否    | -     |
 | exchange                   | string  | 否    | -     |
 | network_recovery_interval  | int     | 否    | -     |
@@ -43,6 +45,7 @@ import ChangeLog from '../changelog/connector-rabbitmq.md';
 | durable                    | boolean | 否    | true  |
 | exclusive                  | boolean | 否    | false |
 | auto_delete                | boolean | 否    | false |
+| passive                    | boolean | 否    | false |
 | common-options             |         | 否    | -     |
 
 ### host [string]
@@ -71,6 +74,16 @@ virtual host,连接 broker 使用的 vhost
 
 设置host、port、username、password和virtual host的简便方式。
 
+### uri [string]
+
+`url` 的兼容别名。`url` 和 `uri` 只能配置一个。
+
+### ssl [boolean]
+
+使用 `host` 和 `port` 配置连接时启用 SSL/TLS。若 URI 本身提供连接信息,请使用 `amqps://` 开头的 `url`。
+
+当 `url` 使用 `amqps://` 时,将按 JVM 信任库校验 Broker 证书并启用主机名校验。此前依赖隐式信任所有证书、使用自签名或私有 
CA 证书的连接,需要将 Broker 证书导入信任库,否则将无法建立连接。
+
 ### queue_name [string]
 
 数据写入的队列名。如果没有配置 `routing_key`,连接器会通过默认 exchange 将消息直接写入该队列。
@@ -110,6 +123,11 @@ virtual host,连接 broker 使用的 vhost
 - true:队列将在最后一个消费者取消订阅时自动删除。
 - false:队列不会自动删除。
 
+### passive [boolean]
+
+- false:按已配置的 durable、exclusive 和 auto_delete 参数声明队列。
+- true:只校验队列已存在,不创建或修改队列。适用于可发布但没有队列声明权限的账号。
+
 ### network_recovery_interval [int]
 
 自动恢复需等待多长时间才尝试重连,单位为毫秒。
@@ -139,6 +157,8 @@ Sink插件常用参数,请参考[Sink常用选项](../common-options/sink-comm
 ## 配置说明
 
 - 如果配置了 `username`,也必须配置 `password`,反过来也一样。
+- `url` 和 `uri` 只能配置一个。`uri` 为兼容已有配置保留,新配置请使用 `url`。
+- 使用 `host` 和 `port` 连接 AMQPS 端点时,请设置 `ssl = true`。
 - `host`、`port`、`virtual_host` 和 `queue_name` 是连接器必填项。`url` 可额外提供 RabbitMQ 
客户端使用的 AMQP URI。
 - `durable`、`exclusive` 和 `auto_delete` 用于连接器声明目标队列。
 - 当 `format` 为 `protobuf` 时,需要同时配置 `protobuf_schema` 和 `protobuf_message_name`。
diff --git a/docs/zh/connectors/source/Rabbitmq.md 
b/docs/zh/connectors/source/Rabbitmq.md
index 1be0f4ef75..5975930117 100644
--- a/docs/zh/connectors/source/Rabbitmq.md
+++ b/docs/zh/connectors/source/Rabbitmq.md
@@ -46,6 +46,8 @@ import ChangeLog from '../changelog/connector-rabbitmq.md';
 | protobuf_schema            | string  | 否  | -     | 当 format 为 protobuf 
时生效,用于解析消息体的 Protobuf Schema                              |
 | protobuf_message_name      | string  | 否  | -     | 当 format 为 protobuf 
时生效,指定要解析的 Protobuf Message 名称                            |
 | url                        | string  | 否  | -     | 便捷方法,用于设置 AMQP URI 
中的字段:主机、端口、用户名、密码和虚拟主机                                   |
+| uri                        | string  | 否  | -     | `url` 的兼容别名              
                                                 |
+| ssl                        | boolean | 否  | false | 使用 host 和 port 连接时是否启用 
SSL/TLS                                      |
 | routing_key                | string  | 否  | -     | RabbitMQ 共享配置中的可选路由键     
                                                    |
 | exchange                   | string  | 否  | -     | RabbitMQ 共享配置中的可选 
exchange                                                     |
 | network_recovery_interval  | int     | 否  | -     | 自动恢复在尝试重新连接之前等待多长时间(毫秒)  
                                                   |
@@ -61,6 +63,7 @@ import ChangeLog from '../changelog/connector-rabbitmq.md';
 | durable                    | boolean | 否  | true  | 队列是否在服务器重启时保留            
                                                   |
 | exclusive                  | boolean | 否  | false | 队列是否仅由当前连接使用             
                                                   |
 | auto_delete                | boolean | 否  | false | 队列是否在最后一个消费者取消订阅时自动删除    
                                                   |
+| passive                    | boolean | 否  | false | 是否只校验已有队列而不声明或创建队列       
                                             |
 | common-options             |         | 否  | -     | 源插件通用参数                  
                                                   |
 
 ### host [string]
@@ -89,6 +92,16 @@ import ChangeLog from '../changelog/connector-rabbitmq.md';
 
 便捷方法,用于设置 AMQP URI 中的字段:主机、端口、用户名、密码和虚拟主机
 
+### uri [string]
+
+`url` 的兼容别名。`url` 和 `uri` 只能配置一个。
+
+### ssl [boolean]
+
+使用 `host` 和 `port` 配置连接时启用 SSL/TLS。若 URI 本身提供连接信息,请使用 `amqps://` 开头的 `url`。
+
+当 `url` 使用 `amqps://` 时,将按 JVM 信任库校验 Broker 证书并启用主机名校验。此前依赖隐式信任所有证书、使用自签名或私有 
CA 证书的连接,需要将 Broker 证书导入信任库,否则将无法建立连接。
+
 ### queue_name [string]
 
 要消费消息的队列。*注意:如果未配置 `tables_configs`,则为必填项。*
@@ -184,6 +197,11 @@ RabbitMQ 共享配置中的可选 exchange。普通队列消费不需要配置
 - true:队列将在最后一个消费者取消订阅时自动删除。
 - false:队列不会自动删除。
 
+### passive
+
+- false:按已配置的 durable、exclusive 和 auto_delete 参数声明队列。
+- true:只校验队列已存在,不创建或修改队列。适用于没有队列声明权限的消费者账号。
+
 ## 迁移指南与配置规则
 
 如果您从仅支持单表读取的早期版本升级,您现有的配置无需任何更改即可正常工作。
@@ -195,6 +213,8 @@ RabbitMQ 共享配置中的可选 exchange。普通队列消费不需要配置
 - 多表模式下,每个队列自己的 `schema` 应放在对应的 `tables_configs` 条目里。
 - 当 `format` 为 `protobuf` 时,需要在队列配置所在层级同时配置 `protobuf_schema` 和 
`protobuf_message_name`。
 - 如果配置了 `username`,也必须配置 `password`,反过来也一样。
+- `url` 和 `uri` 只能配置一个。`uri` 为兼容已有配置保留,新配置请使用 `url`。
+- 使用 `host` 和 `port` 连接 AMQPS 端点时,请设置 `ssl = true`。
 - `host` 和 `port` 总是必填。`virtual_host` 是可选项,除非您的 RabbitMQ 环境要求使用非默认虚拟主机。
 
 ## 示例
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md 
b/docs/zh/introduction/concepts/incompatible-changes.md
index 21a5217a28..ea8fc9d80b 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -4,6 +4,17 @@
 
 ## dev
 
+### RabbitMQ Connector
+
+- **破坏性变更:`amqps://` 连接现在会校验 Broker 证书**
+  - **影响范围**:`seatunnel-connectors-v2/connector-rabbitmq`
+  - **变更说明**:此前使用 `amqps://` 的 `url`/`uri` 建立连接时,会隐式启用“信任所有证书”的
+    TrustManager 且不校验主机名。现在 `amqps://` 连接会强制校验证书,与 `ssl = true` 的
+    host/port 路径行为保持一致。
+  - **影响**:使用自签名或私有 CA 证书的 Broker,升级后通过 `amqps://` 建立的连接将失败。
+  - **迁移指南**:将 Broker 证书(或私有 CA 证书链)导入 SeaTunnel 运行时的 JVM 信任库,或改用
+    `host`/`port` + `ssl = true` 配置并正确设置信任库。
+
 ### Zeta REST 分页参数校验
 
 - **行为变更:分页接口开始校验 `page` 与 `rows`**
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java
index 85c6700678..af7ca2ef62 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java
@@ -29,6 +29,8 @@ import com.rabbitmq.client.ConnectionFactory;
 import com.rabbitmq.client.DefaultConsumer;
 import lombok.extern.slf4j.Slf4j;
 
+import javax.net.ssl.SSLContext;
+
 import java.io.IOException;
 import java.net.URISyntaxException;
 import java.security.KeyManagementException;
@@ -38,6 +40,7 @@ import java.util.concurrent.TimeoutException;
 
 import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.CLOSE_CONNECTION_FAILED;
 import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.CREATE_RABBITMQ_CLIENT_FAILED;
+import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.ILLEGAL_CONFIG;
 import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.INIT_SSL_CONTEXT_FAILED;
 import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.PARSE_URI_FAILED;
 import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.SEND_MESSAGE_FAILED;
@@ -105,7 +108,9 @@ public class RabbitmqClient implements AutoCloseable {
             try {
                 factory.setUri(config.getUri());
             } catch (URISyntaxException e) {
-                throw new RabbitmqConnectorException(PARSE_URI_FAILED, e);
+                throw new RabbitmqConnectorException(
+                        PARSE_URI_FAILED,
+                        "Failed to parse RabbitMQ uri. Check the URI syntax 
without exposing credentials.");
             } catch (KeyManagementException e) {
                 // this should never happen
                 throw new RabbitmqConnectorException(INIT_SSL_CONTEXT_FAILED, 
e);
@@ -113,6 +118,9 @@ public class RabbitmqClient implements AutoCloseable {
                 // this should never happen
                 throw new RabbitmqConnectorException(SETUP_SSL_FACTORY_FAILED, 
e);
             }
+            if (factory.isSSL()) {
+                configureSsl(factory);
+            }
         } else {
             factory.setHost(config.getHost());
             factory.setPort(config.getPort());
@@ -121,6 +129,9 @@ public class RabbitmqClient implements AutoCloseable {
             }
             factory.setUsername(config.getUsername());
             factory.setPassword(config.getPassword());
+            if (config.isSsl()) {
+                configureSsl(factory);
+            }
         }
 
         if (config.getAutomaticRecovery() != null) {
@@ -147,6 +158,16 @@ public class RabbitmqClient implements AutoCloseable {
         return factory;
     }
 
+    /** Configures TLS with the JVM trust store and hostname verification 
enabled. */
+    private void configureSsl(ConnectionFactory factory) {
+        try {
+            factory.useSslProtocol(SSLContext.getDefault());
+            factory.enableHostnameVerification();
+        } catch (NoSuchAlgorithmException e) {
+            throw new RabbitmqConnectorException(INIT_SSL_CONTEXT_FAILED, e);
+        }
+    }
+
     /**
      * Write data to RabbitMQ.
      *
@@ -221,6 +242,29 @@ public class RabbitmqClient implements AutoCloseable {
     /** Declare a specific queue */
     public void setupQueue(String queueName) throws IOException {
         if (StringUtils.isNotEmpty(queueName)) {
+            declareQueue(channel, config, queueName);
+        }
+    }
+
+    private void declareQueueDefaults(Channel channel, RabbitmqConfig config) 
throws IOException {
+        declareQueue(channel, config, config.getQueueName());
+    }
+
+    static void declareQueue(Channel channel, RabbitmqConfig config, String 
queueName)
+            throws IOException {
+        if (config.isPassive()) {
+            try {
+                channel.queueDeclarePassive(queueName);
+            } catch (IOException e) {
+                throw new RabbitmqConnectorException(
+                        ILLEGAL_CONFIG,
+                        String.format(
+                                "Cannot passively declare RabbitMQ queue '%s'. 
The queue may not exist "
+                                        + "or the account may lack access; 
create it first or set passive=false.",
+                                queueName),
+                        e);
+            }
+        } else {
             channel.queueDeclare(
                     queueName,
                     config.getDurable(),
@@ -229,13 +273,4 @@ public class RabbitmqClient implements AutoCloseable {
                     null);
         }
     }
-
-    private void declareQueueDefaults(Channel channel, RabbitmqConfig config) 
throws IOException {
-        channel.queueDeclare(
-                config.getQueueName(),
-                config.getDurable(),
-                config.getExclusive(),
-                config.getAutoDelete(),
-                null);
-    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java
index 68d3ee15d0..1b93f5ccec 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java
@@ -66,6 +66,19 @@ public class RabbitmqBaseOptions extends 
ConnectorCommonOptions {
                     .withDescription(
                             "convenience method for setting the fields in an 
AMQP URI: host, port, username, password and virtual host");
 
+    public static final Option<String> URI =
+            Options.key("uri")
+                    .stringType()
+                    .noDefaultValue()
+                    .withDescription("legacy alias of url for an AMQP URI");
+
+    public static final Option<Boolean> SSL =
+            Options.key("ssl")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "whether to enable SSL/TLS when connecting with 
host and port");
+
     public static final Option<String> ROUTING_KEY =
             Options.key("routing_key")
                     .stringType()
@@ -133,6 +146,13 @@ public class RabbitmqBaseOptions extends 
ConnectorCommonOptions {
                             "true: The queue will be deleted automatically 
when the last consumer unsubscribes."
                                     + "false: The queue will not be 
automatically deleted.");
 
+    public static final Option<Boolean> PASSIVE =
+            Options.key("passive")
+                    .booleanType()
+                    .defaultValue(false)
+                    .withDescription(
+                            "whether to verify an existing queue without 
declaring or creating it");
+
     public static final Option<RabbitmqMessageFormat> FORMAT =
             Options.key("format")
                     .enumType(RabbitmqMessageFormat.class)
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java
index 1bab65f778..eabc9d4018 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java
@@ -18,6 +18,7 @@
 package org.apache.seatunnel.connectors.seatunnel.rabbitmq.config;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException;
 
 import lombok.AllArgsConstructor;
 import lombok.Getter;
@@ -28,17 +29,27 @@ import java.io.Serializable;
 import java.util.HashMap;
 import java.util.Map;
 
+import static 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.ILLEGAL_CONFIG;
+
 @Setter
 @Getter
 @NoArgsConstructor
 @AllArgsConstructor
 public class RabbitmqConfig implements Serializable {
+    /**
+     * Pinned to the default computed UID of the class before ssl/passive were 
added, so execution
+     * plans and checkpoint state serialized by older versions still 
deserialize without throwing
+     * InvalidClassException.
+     */
+    private static final long serialVersionUID = -6715216959598971323L;
+
     private String host;
     private Integer port;
     private String virtualHost;
     private String username;
     private String password;
     private String uri;
+    private boolean ssl;
     private Integer networkRecoveryInterval;
     private Boolean automaticRecovery;
     private Boolean topologyRecovery;
@@ -52,6 +63,7 @@ public class RabbitmqConfig implements Serializable {
     private Boolean durable;
     private Boolean exclusive;
     private Boolean autoDelete;
+    private boolean passive;
     private RabbitmqMessageFormat format;
     private String protobufSchema;
     private String protobufMessageName;
@@ -127,6 +139,8 @@ public class RabbitmqConfig implements Serializable {
         this.durable = config.get(RabbitmqBaseOptions.DURABLE);
         this.exclusive = config.get(RabbitmqBaseOptions.EXCLUSIVE);
         this.autoDelete = config.get(RabbitmqBaseOptions.AUTO_DELETE);
+        this.ssl = config.get(RabbitmqBaseOptions.SSL);
+        this.passive = config.get(RabbitmqBaseOptions.PASSIVE);
         this.format = config.get(RabbitmqBaseOptions.FORMAT);
         if 
(config.getOptional(RabbitmqBaseOptions.PROTOBUF_SCHEMA).isPresent()) {
             this.protobufSchema = 
config.get(RabbitmqBaseOptions.PROTOBUF_SCHEMA);
@@ -137,8 +151,17 @@ public class RabbitmqConfig implements Serializable {
         if 
(config.getOptional(RabbitmqSinkOptions.RABBITMQ_CONFIG).isPresent()) {
             this.sinkOptionProps = 
config.get(RabbitmqSinkOptions.RABBITMQ_CONFIG);
         }
-        if (config.getOptional(RabbitmqBaseOptions.URL).isPresent()) {
+        boolean hasUrl = 
config.getOptional(RabbitmqBaseOptions.URL).isPresent();
+        boolean hasUri = 
config.getOptional(RabbitmqBaseOptions.URI).isPresent();
+        if (hasUrl && hasUri) {
+            throw new RabbitmqConnectorException(
+                    ILLEGAL_CONFIG,
+                    "RabbitMQ connector options 'url' and 'uri' are mutually 
exclusive; please configure only one of them. 'uri' is a legacy alias kept for 
backward compatibility, prefer 'url' for new configurations.");
+        }
+        if (hasUrl) {
             this.uri = config.get(RabbitmqBaseOptions.URL);
+        } else if (hasUri) {
+            this.uri = config.get(RabbitmqBaseOptions.URI);
         }
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java
index 35e8eedd43..d31472d6f7 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java
@@ -29,7 +29,8 @@ public enum RabbitmqConnectorErrorCode implements 
SeaTunnelErrorCode {
     MESSAGE_ACK_REJECTED("RABBITMQ-06", "messages could not be acknowledged 
with basicReject"),
     PARSE_URI_FAILED("RABBITMQ-07", "parse uri failed"),
     INIT_SSL_CONTEXT_FAILED("RABBITMQ-08", "initialize ssl context failed"),
-    SETUP_SSL_FACTORY_FAILED("RABBITMQ-09", "setup ssl factory failed");
+    SETUP_SSL_FACTORY_FAILED("RABBITMQ-09", "setup ssl factory failed"),
+    ILLEGAL_CONFIG("RABBITMQ-10", "illegal connector configuration");
 
     private final String code;
     private final String description;
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java
index 906b5425a6..c41a0ea930 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java
@@ -22,6 +22,7 @@ import org.apache.seatunnel.api.table.connector.TableSink;
 import org.apache.seatunnel.api.table.factory.Factory;
 import org.apache.seatunnel.api.table.factory.TableSinkFactory;
 import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqMessageFormat;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSinkOptions;
@@ -53,6 +54,8 @@ public class RabbitmqSinkFactory implements TableSinkFactory {
                         RabbitmqSinkOptions.PROTOBUF_MESSAGE_NAME)
                 .optional(
                         RabbitmqSinkOptions.URL,
+                        RabbitmqBaseOptions.URI,
+                        RabbitmqSinkOptions.SSL,
                         RabbitmqSinkOptions.ROUTING_KEY,
                         RabbitmqSinkOptions.EXCHANGE,
                         RabbitmqSinkOptions.NETWORK_RECOVERY_INTERVAL,
@@ -63,6 +66,7 @@ public class RabbitmqSinkFactory implements TableSinkFactory {
                         RabbitmqSinkOptions.DURABLE,
                         RabbitmqSinkOptions.EXCLUSIVE,
                         RabbitmqSinkOptions.AUTO_DELETE,
+                        RabbitmqSinkOptions.PASSIVE,
                         RabbitmqSinkOptions.RABBITMQ_CONFIG)
                 .build();
     }
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java
index 48a248fe74..969a3880b8 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java
@@ -25,6 +25,7 @@ import org.apache.seatunnel.api.table.connector.TableSource;
 import org.apache.seatunnel.api.table.factory.Factory;
 import org.apache.seatunnel.api.table.factory.TableSourceFactory;
 import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqMessageFormat;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSingleTableValidator;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSinkOptions;
@@ -69,6 +70,8 @@ public class RabbitmqSourceFactory implements 
TableSourceFactory {
                 .optional(
                         RabbitmqSourceOptions.VIRTUAL_HOST,
                         RabbitmqSourceOptions.URL,
+                        RabbitmqBaseOptions.URI,
+                        RabbitmqSourceOptions.SSL,
                         RabbitmqSourceOptions.ROUTING_KEY,
                         RabbitmqSourceOptions.EXCHANGE,
                         RabbitmqSourceOptions.NETWORK_RECOVERY_INTERVAL,
@@ -79,6 +82,7 @@ public class RabbitmqSourceFactory implements 
TableSourceFactory {
                         RabbitmqSinkOptions.DURABLE,
                         RabbitmqSinkOptions.EXCLUSIVE,
                         RabbitmqSinkOptions.AUTO_DELETE,
+                        RabbitmqSourceOptions.PASSIVE,
                         RabbitmqSourceOptions.REQUESTED_CHANNEL_MAX,
                         RabbitmqSourceOptions.REQUESTED_FRAME_MAX,
                         RabbitmqSourceOptions.REQUESTED_HEARTBEAT,
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java
index 36641e585f..e5ba329aa8 100644
--- 
a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java
@@ -17,12 +17,19 @@
 
 package org.apache.seatunnel.connectors.seatunnel.rabbitmq;
 
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.sink.RabbitmqSinkFactory;
 import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.source.RabbitmqSourceFactory;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.util.HashMap;
+import java.util.Map;
+
 class RabbitmqFactoryTest {
 
     @Test
@@ -30,4 +37,39 @@ class RabbitmqFactoryTest {
         Assertions.assertNotNull((new RabbitmqSourceFactory()).optionRule());
         Assertions.assertNotNull((new RabbitmqSinkFactory()).optionRule());
     }
+
+    @Test
+    void readsLegacyUriOption() {
+        Map<String, Object> options = new HashMap<>();
+        options.put(RabbitmqBaseOptions.URI.key(), 
"amqps://guest:guest@localhost:5671/%2F");
+
+        RabbitmqConfig config = new 
RabbitmqConfig(ReadonlyConfig.fromMap(options));
+
+        Assertions.assertEquals("amqps://guest:guest@localhost:5671/%2F", 
config.getUri());
+    }
+
+    @Test
+    void readsSecurePassiveOptions() {
+        Map<String, Object> options = new HashMap<>();
+        options.put(RabbitmqBaseOptions.HOST.key(), "localhost");
+        options.put(RabbitmqBaseOptions.PORT.key(), 5671);
+        options.put(RabbitmqBaseOptions.SSL.key(), true);
+        options.put(RabbitmqBaseOptions.PASSIVE.key(), true);
+
+        RabbitmqConfig config = new 
RabbitmqConfig(ReadonlyConfig.fromMap(options));
+
+        Assertions.assertTrue(config.isSsl());
+        Assertions.assertTrue(config.isPassive());
+    }
+
+    @Test
+    void rejectsUrlAndUriSetTogether() {
+        Map<String, Object> options = new HashMap<>();
+        options.put(RabbitmqBaseOptions.URL.key(), "amqp://host-a:5672/%2F");
+        options.put(RabbitmqBaseOptions.URI.key(), "amqps://host-b:5671/%2F");
+
+        Assertions.assertThrows(
+                RabbitmqConnectorException.class,
+                () -> new RabbitmqConfig(ReadonlyConfig.fromMap(options)));
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java
 
b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java
new file mode 100644
index 0000000000..a2ab086521
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java
@@ -0,0 +1,83 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.rabbitmq.client;
+
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException;
+
+import org.junit.jupiter.api.Test;
+
+import com.rabbitmq.client.Channel;
+
+import java.io.IOException;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class RabbitmqClientTest {
+
+    @Test
+    void declaresExistingQueuePassivelyWhenConfigured() throws Exception {
+        Channel channel = mock(Channel.class);
+        RabbitmqConfig config = mock(RabbitmqConfig.class);
+        when(config.isPassive()).thenReturn(true);
+
+        RabbitmqClient.declareQueue(channel, config, "existing-queue");
+
+        verify(channel).queueDeclarePassive("existing-queue");
+        verify(channel, never()).queueDeclare("existing-queue", true, false, 
false, null);
+    }
+
+    @Test
+    void declaresQueueWithConfiguredPropertiesByDefault() throws Exception {
+        Channel channel = mock(Channel.class);
+        RabbitmqConfig config = mock(RabbitmqConfig.class);
+        when(config.isPassive()).thenReturn(false);
+        when(config.getDurable()).thenReturn(true);
+        when(config.getExclusive()).thenReturn(false);
+        when(config.getAutoDelete()).thenReturn(false);
+
+        RabbitmqClient.declareQueue(channel, config, "new-queue");
+
+        verify(channel).queueDeclare("new-queue", true, false, false, null);
+        verify(channel, never()).queueDeclarePassive("new-queue");
+    }
+
+    @Test
+    void explainsPassiveQueueDeclarationFailure() throws Exception {
+        Channel channel = mock(Channel.class);
+        RabbitmqConfig config = mock(RabbitmqConfig.class);
+        when(config.isPassive()).thenReturn(true);
+        doThrow(new IOException("queue not found"))
+                .when(channel)
+                .queueDeclarePassive("missing-queue");
+
+        RabbitmqConnectorException exception =
+                assertThrows(
+                        RabbitmqConnectorException.class,
+                        () -> RabbitmqClient.declareQueue(channel, config, 
"missing-queue"));
+
+        assertTrue(exception.getMessage().contains("missing-queue"));
+        assertTrue(exception.getMessage().contains("passive=false"));
+    }
+}

Reply via email to