This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 38a37104e5 [Fix][Connector-V2] Support secure passive RabbitMQ
connections (#11576)
38a37104e5 is described below
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"));
+ }
+}