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-12402-398fdfac5c6f48e894ce94c7b0ba7eeb74c96f7c
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit c12764c05bdfb1f8b2d908d0e34ab2d27bb6999e
Author: Goutam Adwant <[email protected]>
AuthorDate: Tue Sep 22 02:57:56 2026 +0000

    [Fix][Connector-V2] Honor Redis named-user authentication (#12402)
    
    Signed-off-by: Goutam Adwant <[email protected]>
---
 docs/en/connectors/sink/Redis.md                   |  13 ++
 docs/en/connectors/source/Redis.md                 |  12 +-
 .../introduction/concepts/incompatible-changes.md  |  12 +
 docs/zh/connectors/sink/Redis.md                   |  12 +
 docs/zh/connectors/source/Redis.md                 |  10 +-
 .../introduction/concepts/incompatible-changes.md  |  11 +
 .../seatunnel/redis/config/RedisBaseOptions.java   |  11 +-
 .../seatunnel/redis/config/RedisParameters.java    |  65 ++++--
 .../connectors/seatunnel/redis/Redis7Test.java     | 257 +++++++++++++++++++++
 .../seatunnel/redis/RedisFactoryTest.java          |  74 ++++++
 .../seatunnel/e2e/connector/redis/Redis7IT.java    |  46 ++++
 .../connector/redis/RedisTestCaseTemplateIT.java   |   2 +-
 .../src/test/resources/redis-named-user.conf       |  50 ++++
 13 files changed, 555 insertions(+), 20 deletions(-)

diff --git a/docs/en/connectors/sink/Redis.md b/docs/en/connectors/sink/Redis.md
index eaa81845b0..760176e433 100644
--- a/docs/en/connectors/sink/Redis.md
+++ b/docs/en/connectors/sink/Redis.md
@@ -58,6 +58,19 @@ downloaded from Maven Central.
 | multi_table_sink_replica | int | No                          | 1       | 
Writer replica count for multi-table writes. |
 | common-options     | config  | No                          | -       | Sink 
plugin common parameters. See [Sink Common 
Options](../common-options/sink-common-options.md). |
 
+### Authentication
+
+In both `SINGLE` and `CLUSTER` mode, a nonblank `user` selects Redis ACL 
authentication
+(`AUTH user auth`, Redis 6 or later). The connector does not create or modify 
ACL users.
+Create the user and grant its required command and key permissions before 
starting the job, including
+`INFO` for connector initialization, `SELECT` in `SINGLE` mode, and `CLUSTER 
SLOTS` for topology
+discovery in `CLUSTER` mode.
+The password is passed unchanged, including whitespace; omitted or empty 
`auth` is sent as an empty
+password and only works if the ACL user accepts it (for example, a `nopass` 
user).
+
+If `user` is omitted, empty, or whitespace-only, nonblank `auth` uses 
password-only authentication
+as the default user. If both options are omitted or blank, no authentication 
command is sent.
+
 ## Write Rules
 
 ### key
diff --git a/docs/en/connectors/source/Redis.md 
b/docs/en/connectors/source/Redis.md
index 77a203d473..093b4e5cdb 100644
--- a/docs/en/connectors/source/Redis.md
+++ b/docs/en/connectors/source/Redis.md
@@ -232,11 +232,19 @@ redis data types, support `key` `string` `hash` `list` 
`set` `zset`
 
 ### user [string]
 
-redis authentication user, you need it when you connect to an encrypted cluster
+Redis ACL username (Redis 6 or later), supported in both `SINGLE` and 
`CLUSTER` mode.
+When nonblank, the connector authenticates with `AUTH user auth`; it does not 
create or modify ACL users.
+Create the user and grant the required command and key permissions before 
starting the job, including
+`INFO` for connector initialization, `SELECT` in `SINGLE` mode, and `CLUSTER 
SLOTS` for topology
+discovery in `CLUSTER` mode.
+If `user` is omitted, empty, or whitespace-only, a nonblank `auth` uses 
password-only authentication
+as the default user; otherwise no authentication command is sent.
 
 ### auth [string]
 
-redis authentication password, you need it when you connect to an encrypted 
cluster
+Redis authentication password. With a nonblank `user`, the password is passed 
unchanged, including
+whitespace. An omitted or empty password is sent as an empty string and works 
only if that ACL user
+accepts it (for example, a user configured with `nopass`).
 
 ### db_num [int]
 
diff --git a/docs/en/introduction/concepts/incompatible-changes.md 
b/docs/en/introduction/concepts/incompatible-changes.md
index 99181067e1..b90ad91efd 100644
--- a/docs/en/introduction/concepts/incompatible-changes.md
+++ b/docs/en/introduction/concepts/incompatible-changes.md
@@ -5,6 +5,18 @@ You need to check this document before you upgrade to related 
version.
 
 ## dev
 
+### Redis Authentication
+
+- Redis sources and sinks now authenticate as the configured nonblank `user` 
in both `SINGLE` and
+  `CLUSTER` mode. Previously, `SINGLE` used password-only authentication 
followed by `ACL SETUSER`,
+  and `CLUSTER` ignored `user`. Connection setup no longer creates or modifies 
ACL users.
+- Before upgrading, create the intended ACL user and grant its required 
command and key permissions,
+  including `INFO` for connector initialization, `SELECT` in `SINGLE` mode, 
and `CLUSTER SLOTS` for
+  topology discovery in `CLUSTER` mode. Set `auth` to that user's password. An 
omitted or empty password
+  is sent as an empty string when `user` is nonblank.
+- To keep using the default user, remove `user` and retain `auth` when a 
password is required.
+  Named users require Redis 6 or later. Legacy configurations without a 
username remain unchanged.
+
 ### RabbitMQ Connector
 
 - **Breaking Change: `amqps://` connections now verify broker certificates**
diff --git a/docs/zh/connectors/sink/Redis.md b/docs/zh/connectors/sink/Redis.md
index c430f69bb6..3c59c1c1ab 100644
--- a/docs/zh/connectors/sink/Redis.md
+++ b/docs/zh/connectors/sink/Redis.md
@@ -57,6 +57,18 @@ Redis 接收器连接器可以在批处理或流处理作业中把上游数据
 | multi_table_sink_replica | int | 否                          | 1      | 
多表写入时的写入器副本数。 |
 | common-options     | config  | 否                          | -      | 
接收器插件通用参数,详情请参考[接收器通用选项](../common-options/sink-common-options.md)。 |
 
+### 认证
+
+在 `SINGLE` 和 `CLUSTER` 模式下,非空白的 `user` 使用 Redis ACL 认证
+(`AUTH user auth`,需要 Redis 6 或更新版本)。连接器不会创建或修改 ACL 用户。
+启动作业前,请创建用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`,
+`SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。
+密码将原样传递,包括空白字符;省略 `auth` 或使用空字符串时将发送空密码,
+仅当该 ACL 用户允许时才能成功认证(例如配置了 `nopass` 的用户)。
+
+若省略 `user`,或其值为空字符串、仅包含空白字符,则非空白的 `auth` 将用于默认用户的密码认证。
+如果两个选项都省略或为空白,则不发送认证命令。
+
 ## 写入规则
 
 ### key
diff --git a/docs/zh/connectors/source/Redis.md 
b/docs/zh/connectors/source/Redis.md
index b930d1c824..77f84c018d 100644
--- a/docs/zh/connectors/source/Redis.md
+++ b/docs/zh/connectors/source/Redis.md
@@ -224,11 +224,17 @@ redis 数据类型, 支持 `key` `string` `hash` `list` `set` `zset`。
 
 ### user [string]
 
-Redis 认证身份用户,当连接到加密集群时需要使用
+Redis ACL 用户名(需要 Redis 6 或更新版本),支持 `SINGLE` 和 `CLUSTER` 模式。
+当用户名非空白时,连接器通过 `AUTH user auth` 认证,不会创建或修改 ACL 用户。
+启动作业前,请创建用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`,
+`SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。
+若省略 `user`,或其值为空字符串、仅包含空白字符,则非空白的 `auth` 将用于默认用户的密码认证;
+否则不发送认证命令。
 
 ### auth [string]
 
-Redis 认证密钥,当连接到加密集群时需要使用
+Redis 认证密码。当 `user` 非空白时,密码将原样传递,包括空白字符。
+省略密码或使用空字符串时,将发送空密码,仅当该 ACL 用户允许时才能成功认证(例如配置了 `nopass` 的用户)。
 
 ### db_num [int]
 
diff --git a/docs/zh/introduction/concepts/incompatible-changes.md 
b/docs/zh/introduction/concepts/incompatible-changes.md
index 8a3b55c330..fb0a541081 100644
--- a/docs/zh/introduction/concepts/incompatible-changes.md
+++ b/docs/zh/introduction/concepts/incompatible-changes.md
@@ -4,6 +4,17 @@
 
 ## dev
 
+### Redis 认证
+
+- Redis Source 和 Sink 现在会在 `SINGLE` 和 `CLUSTER` 模式下以非空白的 `user` 指定的用户认证。
+  此前,`SINGLE` 模式先使用仅密码认证,再执行 `ACL SETUSER`;`CLUSTER` 模式忽略 `user`。
+  连接初始化不再创建或修改 ACL 用户。
+- 升级前,请创建目标 ACL 用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`,
+  `SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。
+  将 `auth` 设置为该用户的密码。当 `user` 非空白时,省略密码或使用空字符串将发送空密码。
+- 如需继续使用默认用户,请移除 `user`,并在需要密码时保留 `auth`。
+  命名用户需要 Redis 6 或更新版本;未配置用户名的旧配置行为保持不变。
+
 ### RabbitMQ Connector
 
 - **破坏性变更:`amqps://` 连接现在会校验 Broker 证书**
diff --git 
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java
 
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java
index 690f69fd92..8ee37aae4a 100644
--- 
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java
+++ 
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java
@@ -50,7 +50,11 @@ public class RedisBaseOptions extends ConnectorCommonOptions 
{
                     .stringType()
                     .noDefaultValue()
                     .withDescription(
-                            "redis authentication password, you need it when 
you connect to an encrypted cluster");
+                            "Redis authentication password for SINGLE and 
CLUSTER modes. With a nonblank user,"
+                                    + " the password is passed unchanged; 
omitted or empty auth sends an empty"
+                                    + " password, which the ACL user must 
accept (for example, via nopass)."
+                                    + " With a blank user, nonblank auth 
authenticates as the default user;"
+                                    + " otherwise no authentication command is 
sent.");
 
     public static final Option<Integer> DB_NUM =
             Options.key("db_num")
@@ -64,7 +68,10 @@ public class RedisBaseOptions extends ConnectorCommonOptions 
{
                     .stringType()
                     .noDefaultValue()
                     .withDescription(
-                            "redis authentication user, you need it when you 
connect to an encrypted cluster");
+                            "Redis ACL username (Redis 6 or later) for SINGLE 
and CLUSTER modes."
+                                    + " When nonblank, the connector 
authenticates with AUTH user auth"
+                                    + " without creating or modifying ACL 
users. Omitted, empty, or"
+                                    + " whitespace-only user retains 
default-user authentication behavior.");
 
     public static final Option<String> KEY_PATTERN =
             Options.key("keys")
diff --git 
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java
 
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java
index cac06d2b47..d8458a282b 100644
--- 
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java
+++ 
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java
@@ -29,6 +29,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisConnectorE
 import lombok.Data;
 import lombok.extern.slf4j.Slf4j;
 import redis.clients.jedis.ConnectionPoolConfig;
+import redis.clients.jedis.DefaultJedisClientConfig;
 import redis.clients.jedis.HostAndPort;
 import redis.clients.jedis.Jedis;
 import redis.clients.jedis.JedisCluster;
@@ -144,11 +145,16 @@ public class RedisParameters implements Serializable {
 
     public RedisClient buildRedisClient() {
         Jedis jedis = this.buildJedis();
-        this.redisVersion = extractRedisVersion(jedis);
-        if (mode.equals(RedisBaseOptions.RedisMode.SINGLE)) {
-            return new RedisSingleClient(this, jedis, redisVersion);
-        } else {
-            return new RedisClusterClient(this, jedis, redisVersion);
+        try {
+            this.redisVersion = extractRedisVersion(jedis);
+            if (mode.equals(RedisBaseOptions.RedisMode.SINGLE)) {
+                return new RedisSingleClient(this, jedis, redisVersion);
+            } else {
+                return new RedisClusterClient(this, jedis, redisVersion);
+            }
+        } catch (RuntimeException | Error failure) {
+            closeAfterFailure(jedis, failure);
+            throw failure;
         }
     }
 
@@ -180,18 +186,27 @@ public class RedisParameters implements Serializable {
                 "Did not get the expected redis_version from the jedis.info() 
method");
     }
 
+    /**
+     * Uses named-user AUTH when user is nonblank; otherwise preserves 
password-only or no-auth
+     * behavior. Authentication selects an existing identity, never 
administers users via ACL
+     * SETUSER.
+     */
     public Jedis buildJedis() {
         switch (mode) {
             case SINGLE:
                 Jedis jedis = new Jedis(host, port);
-                if (StringUtils.isNotBlank(auth)) {
-                    jedis.auth(auth);
+                try {
+                    if (StringUtils.isNotBlank(user)) {
+                        jedis.auth(user, StringUtils.defaultString(auth));
+                    } else if (StringUtils.isNotBlank(auth)) {
+                        jedis.auth(auth);
+                    }
+                    jedis.select(dbNum);
+                    return jedis;
+                } catch (RuntimeException | Error failure) {
+                    closeAfterFailure(jedis, failure);
+                    throw failure;
                 }
-                if (StringUtils.isNotBlank(user)) {
-                    jedis.aclSetUser(user);
-                }
-                jedis.select(dbNum);
-                return jedis;
             case CLUSTER:
                 HashSet<HostAndPort> nodes = new HashSet<>();
                 for (String redisNode : redisNodes) {
@@ -202,7 +217,19 @@ public class RedisParameters implements Serializable {
                 }
                 ConnectionPoolConfig connectionPoolConfig = new 
ConnectionPoolConfig();
                 JedisCluster jedisCluster;
-                if (StringUtils.isNotBlank(auth)) {
+                if (StringUtils.isNotBlank(user)) {
+                    jedisCluster =
+                            new JedisCluster(
+                                    nodes,
+                                    DefaultJedisClientConfig.builder()
+                                            .user(user)
+                                            
.password(StringUtils.defaultString(auth))
+                                            
.connectionTimeoutMillis(JedisCluster.DEFAULT_TIMEOUT)
+                                            
.socketTimeoutMillis(JedisCluster.DEFAULT_TIMEOUT)
+                                            .build(),
+                                    JedisCluster.DEFAULT_MAX_ATTEMPTS,
+                                    connectionPoolConfig);
+                } else if (StringUtils.isNotBlank(auth)) {
                     jedisCluster =
                             new JedisCluster(
                                     nodes,
@@ -222,4 +249,16 @@ public class RedisParameters implements Serializable {
                         CommonErrorCode.OPERATION_NOT_SUPPORTED, "Not support 
this redis mode");
         }
     }
+
+    /**
+     * Releases a connection after failed initialization, including an Error, 
before the caller
+     * rethrows the original failure. Cleanup failures are suppressed to 
preserve that failure.
+     */
+    private void closeAfterFailure(Jedis jedis, Throwable failure) {
+        try {
+            jedis.close();
+        } catch (RuntimeException | Error closeFailure) {
+            failure.addSuppressed(closeFailure);
+        }
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java
 
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java
index a2f008e799..87b7e3c3a0 100644
--- 
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java
+++ 
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java
@@ -16,16 +16,273 @@
  */
 package org.apache.seatunnel.connectors.seatunnel.redis;
 
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.redis.client.RedisClient;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.JedisWrapper;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisBaseOptions;
 import 
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisContainerInfo;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
 
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.condition.DisabledOnOs;
 import org.junit.jupiter.api.condition.OS;
+import org.testcontainers.containers.GenericContainer;
+
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.exceptions.JedisDataException;
+
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.function.BooleanSupplier;
+import java.util.stream.IntStream;
 
 @DisabledOnOs(
         value = OS.WINDOWS,
         disabledReason = "There is no docker environment on the windows test 
system")
 public class Redis7Test extends RedisTemplateTest {
 
+    @Test
+    public void namedUserAuthentication() {
+        String user = "seatunnel_named";
+        // Sharing the default user's password must not silently select the 
default identity.
+        jedis.aclSetUser(
+                user,
+                "reset",
+                "on",
+                ">" + password,
+                "~auth:*",
+                "+select",
+                "+info",
+                "+get",
+                "+set",
+                "+acl|whoami");
+        List<String> aclBefore = jedis.aclList();
+        try {
+            RedisParameters parameters = connectionParameters(user, password);
+            parameters.setDbNum(2);
+            try (Jedis connection = parameters.buildJedis()) {
+                Assertions.assertEquals(user, connection.aclWhoAmI());
+                connection.set("auth:selected", "value");
+                Assertions.assertThrows(
+                        JedisDataException.class, () -> 
connection.set("outside:auth", "value"));
+                Assertions.assertThrows(
+                        JedisDataException.class, () -> 
connection.aclSetUser("unexpected"));
+            }
+            jedis.select(2);
+            Assertions.assertEquals("value", jedis.get("auth:selected"));
+            jedis.del("auth:selected");
+            jedis.select(0);
+            RedisClient client = parameters.buildRedisClient();
+            try {
+                Assertions.assertEquals(7, parameters.getRedisVersion());
+            } finally {
+                client.close();
+            }
+            Assertions.assertEquals(aclBefore, jedis.aclList());
+            jedis.aclSetUser(user, "resetpass", ">named-password");
+            try (Jedis connection = connectionParameters(user, 
"named-password").buildJedis()) {
+                Assertions.assertEquals(user, connection.aclWhoAmI());
+            }
+        } finally {
+            jedis.select(0);
+            jedis.aclDelUser(user);
+        }
+    }
+
+    @Test
+    public void namedUserMissingBlankAndWrongPassword() {
+        String user = "seatunnel_nopass";
+        jedis.aclSetUser(user, "reset", "on", "nopass", "+select", 
"+acl|whoami");
+        try {
+            for (String auth : new String[] {null, "", "  "}) {
+                try (Jedis connection = connectionParameters(user, 
auth).buildJedis()) {
+                    Assertions.assertEquals(user, connection.aclWhoAmI());
+                }
+            }
+            jedis.aclSetUser(user, "resetpass", ">named-password");
+            for (String auth : new String[] {null, "", "  ", 
"wrong-password"}) {
+                Assertions.assertThrows(
+                        JedisDataException.class,
+                        () -> connectionParameters(user, auth).buildJedis());
+            }
+            jedis.aclSetUser(user, "resetpass", ">  ");
+            try (Jedis connection = connectionParameters(user, "  
").buildJedis()) {
+                Assertions.assertEquals(user, connection.aclWhoAmI());
+            }
+        } finally {
+            jedis.aclDelUser(user);
+        }
+    }
+
+    @Test
+    public void legacyDefaultUserAuthentication() {
+        for (String user : new String[] {null, "", "  "}) {
+            try (Jedis connection = connectionParameters(user, 
password).buildJedis()) {
+                Assertions.assertEquals("default", connection.aclWhoAmI());
+            }
+        }
+        jedis.aclSetUser("default", "nopass");
+        try {
+            for (String user : new String[] {null, "", "  "}) {
+                for (String auth : new String[] {null, "", "  "}) {
+                    try (Jedis connection = connectionParameters(user, 
auth).buildJedis()) {
+                        Assertions.assertEquals("default", 
connection.aclWhoAmI());
+                    }
+                }
+            }
+            List<String> aclBefore = jedis.aclList();
+            Assertions.assertThrows(
+                    JedisDataException.class,
+                    () -> connectionParameters("nonexistent", 
null).buildJedis());
+            Assertions.assertEquals(aclBefore, jedis.aclList());
+        } finally {
+            jedis.aclSetUser("default", "resetpass", ">" + password);
+        }
+    }
+
+    @Test
+    public void failedInitializationClosesConnections() throws 
InterruptedException {
+        String user = "seatunnel_denied";
+        jedis.aclSetUser(user, "reset", "on", ">named-password", "+select");
+        try {
+            assertConnectionCountUnchanged(() -> connectionParameters(user, 
"wrong").buildJedis());
+            RedisParameters parameters = connectionParameters(user, 
"named-password");
+            parameters.setDbNum(-1);
+            assertConnectionCountUnchanged(parameters::buildJedis);
+            parameters.setDbNum(0);
+            assertConnectionCountUnchanged(parameters::buildRedisClient);
+            jedis.aclSetUser(user, "-select");
+            assertConnectionCountUnchanged(parameters::buildJedis);
+        } finally {
+            jedis.aclDelUser(user);
+        }
+    }
+
+    private void assertConnectionCountUnchanged(Runnable connect) throws 
InterruptedException {
+        int connections = jedis.clientList().split("\n").length;
+        Assertions.assertThrows(JedisDataException.class, connect::run);
+        awaitCondition(() -> connections == 
jedis.clientList().split("\n").length);
+    }
+
+    @Test
+    public void namedClusterAuthentication() throws InterruptedException, 
UnknownHostException {
+        // A real cluster-enabled server owns all slots; no external cluster 
is required.
+        try (GenericContainer<?> cluster =
+                new GenericContainer<>("redis:7")
+                        .withExposedPorts(6379)
+                        .withCommand(
+                                "redis-server",
+                                "--cluster-enabled",
+                                "yes",
+                                "--save",
+                                "",
+                                "--appendonly",
+                                "no")) {
+            cluster.start();
+            try (Jedis admin = new Jedis(cluster.getHost(), 
cluster.getFirstMappedPort())) {
+                admin.configSet(
+                        "cluster-announce-ip",
+                        
InetAddress.getByName(cluster.getHost()).getHostAddress());
+                admin.configSet("cluster-announce-port", 
cluster.getFirstMappedPort().toString());
+                admin.clusterAddSlots(IntStream.range(0, 16384).toArray());
+                awaitCondition(() -> 
admin.clusterInfo().contains("cluster_state:ok"));
+                admin.aclSetUser(
+                        "seatunnel_cluster",
+                        "reset",
+                        "on",
+                        ">named-password",
+                        "~auth:*",
+                        "+cluster|slots",
+                        "+info",
+                        "+get",
+                        "+set",
+                        "+acl|whoami");
+                RedisParameters parameters =
+                        connectionParameters("seatunnel_cluster", 
"named-password");
+                parameters.setMode(RedisBaseOptions.RedisMode.CLUSTER);
+                parameters.setRedisNodes(
+                        Collections.singletonList(
+                                cluster.getHost() + ":" + 
cluster.getFirstMappedPort()));
+                List<String> aclBefore = admin.aclList();
+                try (JedisWrapper connection = (JedisWrapper) 
parameters.buildJedis()) {
+                    for (String node : connection.getClusterNodes().keySet()) {
+                        Assertions.assertEquals(
+                                "seatunnel_cluster", 
connection.getJedis(node).aclWhoAmI());
+                    }
+                    connection.set("auth:cluster", "value");
+                    Assertions.assertEquals("value", 
connection.get("auth:cluster"));
+                }
+                RedisClient client = parameters.buildRedisClient();
+                client.close();
+                Assertions.assertEquals(aclBefore, admin.aclList());
+
+                parameters.setUser("");
+                parameters.setAuth("");
+                assertClusterIdentity(parameters, "default");
+                admin.aclSetUser("default", "resetpass", ">named-password");
+                parameters.setAuth("named-password");
+                assertClusterIdentity(parameters, "default");
+                parameters.setUser("seatunnel_cluster");
+                assertClusterIdentity(parameters, "seatunnel_cluster");
+
+                parameters.setAuth("wrong-password");
+                int beforeAuthenticationFailure = 
admin.clientList().split("\n").length;
+                Assertions.assertThrows(JedisDataException.class, 
parameters::buildJedis);
+                awaitCondition(
+                        () -> beforeAuthenticationFailure == 
admin.clientList().split("\n").length);
+                parameters.setAuth("named-password");
+                admin.aclSetUser("seatunnel_cluster", "-info");
+                int connections = admin.clientList().split("\n").length;
+                Assertions.assertThrows(RuntimeException.class, 
parameters::buildRedisClient);
+                awaitCondition(() -> connections == 
admin.clientList().split("\n").length);
+
+                admin.aclSetUser("seatunnel_cluster", "nopass");
+                for (String auth : new String[] {null, "", "  "}) {
+                    parameters.setAuth(auth);
+                    assertClusterIdentity(parameters, "seatunnel_cluster");
+                }
+            }
+        }
+    }
+
+    private void assertClusterIdentity(RedisParameters parameters, String 
user) {
+        try (JedisWrapper connection = (JedisWrapper) parameters.buildJedis()) 
{
+            for (String node : connection.getClusterNodes().keySet()) {
+                Assertions.assertEquals(user, 
connection.getJedis(node).aclWhoAmI());
+            }
+        }
+    }
+
+    private void awaitCondition(BooleanSupplier condition) throws 
InterruptedException {
+        for (int attempt = 0; attempt < 100; attempt++) {
+            if (condition.getAsBoolean()) {
+                return;
+            }
+            Thread.sleep(50);
+        }
+        Assertions.assertTrue(condition.getAsBoolean());
+    }
+
+    private RedisParameters connectionParameters(String user, String auth) {
+        Map<String, Object> config = new HashMap<>();
+        config.put("host", redisContainer.getHost());
+        config.put("port", redisContainer.getFirstMappedPort());
+        if (user != null) {
+            config.put("user", user);
+        }
+        if (auth != null) {
+            config.put("auth", auth);
+        }
+        RedisParameters parameters = new RedisParameters();
+        parameters.buildConnectionConfig(ReadonlyConfig.fromMap(config));
+        return parameters;
+    }
+
     @Override
     public RedisContainerInfo getRedisContainerInfo() {
         return new RedisContainerInfo("redis-e2e", 6379, "SeaTunnel", 
"redis:7");
diff --git 
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
 
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
index 67fb4a8989..f1b6d3850e 100644
--- 
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
@@ -29,6 +29,8 @@ import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
 import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
 import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisBaseOptions;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
 import org.apache.seatunnel.connectors.seatunnel.redis.sink.RedisSink;
 import org.apache.seatunnel.connectors.seatunnel.redis.sink.RedisSinkFactory;
 import 
org.apache.seatunnel.connectors.seatunnel.redis.source.RedisSourceFactory;
@@ -38,6 +40,11 @@ import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
 import org.junit.jupiter.params.provider.Arguments;
 import org.junit.jupiter.params.provider.MethodSource;
+import org.mockito.MockedConstruction;
+import org.mockito.Mockito;
+
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.exceptions.JedisDataException;
 
 import java.io.IOException;
 import java.util.ArrayList;
@@ -53,6 +60,73 @@ class RedisFactoryTest {
     private static final OptionRule SOURCE_RULE = new 
RedisSourceFactory().optionRule();
     private static final OptionRule SINK_RULE = new 
RedisSinkFactory().optionRule();
 
+    @Test
+    void namedAuthenticationDoesNotAdministerUsers() {
+        try (MockedConstruction<Jedis> construction = 
Mockito.mockConstruction(Jedis.class)) {
+            RedisParameters parameters = singleConnectionParameters();
+            parameters.setUser("named-user");
+            parameters.setAuth("named-password");
+            parameters.setDbNum(3);
+            Jedis connection = parameters.buildJedis();
+            Assertions.assertSame(construction.constructed().get(0), 
connection);
+            Mockito.verify(connection).auth("named-user", "named-password");
+            Mockito.verify(connection).select(3);
+            Mockito.verifyNoMoreInteractions(connection);
+            connection.close();
+        }
+    }
+
+    @Test
+    void authenticationFailureClosesConnectionAndPreservesCause() {
+        JedisDataException failure = new JedisDataException("Authentication 
failed");
+        RuntimeException closeFailure = new RuntimeException("Close failed");
+        try (MockedConstruction<Jedis> construction =
+                Mockito.mockConstruction(
+                        Jedis.class,
+                        (connection, context) -> {
+                            Mockito.when(connection.auth("named-user", 
"wrong")).thenThrow(failure);
+                            
Mockito.doThrow(closeFailure).when(connection).close();
+                        })) {
+            RedisParameters parameters = singleConnectionParameters();
+            parameters.setUser("named-user");
+            parameters.setAuth("wrong");
+            Assertions.assertSame(
+                    failure,
+                    Assertions.assertThrows(JedisDataException.class, 
parameters::buildJedis));
+            Mockito.verify(construction.constructed().get(0)).close();
+            Assertions.assertArrayEquals(new Throwable[] {closeFailure}, 
failure.getSuppressed());
+        }
+    }
+
+    @Test
+    void versionInitializationFailureClosesConnection() {
+        for (String info : Arrays.asList("redis_version:not-a-version", 
"no-version")) {
+            RedisParameters parameters = 
Mockito.spy(singleConnectionParameters());
+            Jedis connection = Mockito.mock(Jedis.class);
+            Mockito.doReturn(connection).when(parameters).buildJedis();
+            Mockito.when(connection.info()).thenReturn(info);
+            Assertions.assertThrows(RuntimeException.class, 
parameters::buildRedisClient);
+            Mockito.verify(connection).close();
+        }
+        RedisParameters parameters = Mockito.spy(singleConnectionParameters());
+        Jedis connection = Mockito.mock(Jedis.class);
+        Mockito.doReturn(connection).when(parameters).buildJedis();
+        JedisDataException denied = new JedisDataException("INFO denied");
+        Mockito.when(connection.info()).thenThrow(denied);
+        Assertions.assertSame(
+                denied,
+                Assertions.assertThrows(JedisDataException.class, 
parameters::buildRedisClient));
+        Mockito.verify(connection).close();
+    }
+
+    private RedisParameters singleConnectionParameters() {
+        RedisParameters parameters = new RedisParameters();
+        parameters.setHost("localhost");
+        parameters.setPort(6379);
+        parameters.setMode(RedisBaseOptions.RedisMode.SINGLE);
+        return parameters;
+    }
+
     @Test
     void optionRule() {
         Assertions.assertNotNull(SOURCE_RULE);
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java
index 7e065837ff..33c6516539 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java
@@ -17,12 +17,58 @@
 package org.apache.seatunnel.e2e.connector.redis;
 
 import 
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisContainerInfo;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
 
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.TestTemplate;
 import org.junit.jupiter.api.parallel.ResourceLock;
+import org.testcontainers.containers.Container;
+
+import java.io.IOException;
+import java.util.Collections;
+import java.util.List;
 
 @ResourceLock("redis-standalone-e2e")
 public class Redis7IT extends RedisTestCaseTemplateIT {
 
+    @TestTemplate
+    public void testNamedUserSourceAndSink(TestContainer container)
+            throws IOException, InterruptedException {
+        jedis.aclSetUser(
+                "seatunnel_reader",
+                "reset",
+                "on",
+                ">reader-password",
+                "~acl:source:*",
+                "+select",
+                "+info",
+                "+scan",
+                "+type",
+                "+get",
+                "+mget");
+        jedis.aclSetUser(
+                "seatunnel_writer",
+                "reset",
+                "on",
+                ">writer-password",
+                "~acl:result",
+                "+select",
+                "+info",
+                "+lpush");
+        List<String> aclBefore = jedis.aclList();
+        jedis.set("acl:source:1", "{\"value\":\"named-user\"}");
+        try {
+            Container.ExecResult result = 
container.executeJob("/redis-named-user.conf");
+            Assertions.assertEquals(0, result.getExitCode());
+            Assertions.assertEquals(
+                    Collections.singletonList("named-user"), 
jedis.lrange("acl:result", 0, -1));
+            Assertions.assertEquals(aclBefore, jedis.aclList());
+        } finally {
+            jedis.del("acl:source:1", "acl:result");
+            jedis.aclDelUser("seatunnel_reader", "seatunnel_writer");
+        }
+    }
+
     @Override
     public RedisContainerInfo getRedisContainerInfo() {
         return new RedisContainerInfo("redis-e2e", 6379, "SeaTunnel", 
"redis:7");
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java
index bb1650f93d..e6aabdae18 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java
@@ -89,7 +89,7 @@ public abstract class RedisTestCaseTemplateIT extends 
TestSuiteBase implements T
 
     private GenericContainer<?> redisContainer;
 
-    private Jedis jedis;
+    protected Jedis jedis;
 
     @BeforeAll
     @Override
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf
new file mode 100644
index 0000000000..2f88da84e2
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf
@@ -0,0 +1,50 @@
+#
+# 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.
+#
+
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  Redis {
+    host = "redis-e2e"
+    port = 6379
+    user = "seatunnel_reader"
+    auth = "reader-password"
+    keys = "acl:source:*"
+    data_type = string
+    format = json
+    schema = {
+      fields {
+        value = string
+      }
+    }
+  }
+}
+
+sink {
+  Redis {
+    host = "redis-e2e"
+    port = 6379
+    user = "seatunnel_writer"
+    auth = "writer-password"
+    key = "acl:result"
+    data_type = list
+    value_field = "value"
+  }
+}

Reply via email to