diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterConnection.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterConnection.java
index 15a83f3841..8a2526cc4f 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterConnection.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterConnection.java
@@ -15,18 +15,21 @@
*/
package org.springframework.data.redis.connection.jedis;
+import redis.clients.jedis.CommandArguments;
import redis.clients.jedis.Connection;
import redis.clients.jedis.ConnectionPool;
import redis.clients.jedis.HostAndPort;
import redis.clients.jedis.Jedis;
import redis.clients.jedis.JedisCluster;
import redis.clients.jedis.JedisClusterInfoCache;
+import redis.clients.jedis.Protocol;
import redis.clients.jedis.RedisClusterClient;
import redis.clients.jedis.UnifiedJedis;
import redis.clients.jedis.providers.ClusterConnectionProvider;
import java.time.Duration;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collection;
import java.util.Collections;
import java.util.LinkedHashMap;
@@ -35,6 +38,7 @@
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
+import java.util.stream.Stream;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
@@ -65,7 +69,7 @@
* {@link RedisClusterConnection} implementation on top of {@link RedisClusterClient}.
*
* Uses the native {@link RedisClusterClient} api where possible and falls back to direct node communication using
- * {@link Jedis} where needed.
+ * {@link UnifiedJedis} where needed.
*
* Pipelines and transactions are not supported in cluster mode. This class is not Thread-safe and instances should not
* be shared across threads.
@@ -250,8 +254,8 @@ public Object execute(@NonNull String command, byte @NonNull [] @NonNull... args
Assert.notNull(command, "Command must not be null");
Assert.notNull(args, "Args must not be null");
- JedisClusterCommandCallback commandCallback = jedis -> jedis
- .sendCommand(JedisClientUtils.getCommand(command), args);
+ JedisClusterCommandCallback commandCallback = client -> client
+ .executeCommand(new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(args));
return this.clusterCommandExecutor.executeCommandOnArbitraryNode(commandCallback).getValue();
}
@@ -268,8 +272,8 @@ public T execute(@NonNull String command, byte @NonNull [] key, @NonNull Col
RedisClusterNode keyMaster = this.topologyProvider.getTopology().getKeyServingMasterNode(key);
- JedisClusterCommandCallback commandCallback = jedis -> (T) jedis
- .sendCommand(JedisClientUtils.getCommand(command), commandArgs);
+ JedisClusterCommandCallback commandCallback = client -> (T) client
+ .executeCommand(new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(commandArgs));
return this.clusterCommandExecutor.executeCommandOnSingleNode(commandCallback, keyMaster).getValue();
}
@@ -318,8 +322,8 @@ public List execute(@NonNull String command, @NonNull Collection commandCallback = (jedis,
- key) -> (T) jedis.sendCommand(JedisClientUtils.getCommand(command), getCommandArguments(key, args));
+ JedisMultiKeyClusterCommandCallback commandCallback = (jedis, key) -> (T) jedis.executeCommand(
+ new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(getCommandArguments(key, args)));
return this.clusterCommandExecutor.executeMultiKeyCommand(commandCallback, keys).resultsAsList();
}
@@ -450,7 +454,7 @@ public byte[] echo(byte @NonNull [] message) {
@Override
public String ping() {
- JedisClusterCommandCallback command = Jedis::ping;
+ JedisClusterCommandCallback command = UnifiedJedis::ping;
return !this.clusterCommandExecutor.executeCommandOnAllNodes(command).resultsAsList().isEmpty() ? "PONG" : null;
}
@@ -458,7 +462,7 @@ public String ping() {
@Override
public String ping(@NonNull RedisClusterNode node) {
- JedisClusterCommandCallback command = Jedis::ping;
+ JedisClusterCommandCallback command = UnifiedJedis::ping;
return this.clusterCommandExecutor.executeCommandOnSingleNode(command, node).getValue();
}
@@ -475,12 +479,17 @@ public void clusterSetSlot(@NonNull RedisClusterNode node, int slot, @NonNull Ad
RedisClusterNode nodeToUse = this.topologyProvider.getTopology().lookup(node);
String nodeId = nodeToUse.getId();
-
- JedisClusterCommandCallback command = jedis -> switch (mode) {
- case IMPORTING -> jedis.clusterSetSlotImporting(slot, nodeId);
- case MIGRATING -> jedis.clusterSetSlotMigrating(slot, nodeId);
- case STABLE -> jedis.clusterSetSlotStable(slot);
- case NODE -> jedis.clusterSetSlotNode(slot, nodeId);
+ String slotId = String.valueOf(slot);
+
+ JedisClusterCommandCallback command = client -> switch (mode) {
+ case IMPORTING -> client.executeCommand(
+ new CommandArguments(Protocol.Command.CLUSTER).add("SETSLOT").add(slotId).add("IMPORTING").add(nodeId));
+ case MIGRATING -> client.executeCommand(
+ new CommandArguments(Protocol.Command.CLUSTER).add("SETSLOT").add(slotId).add("MIGRATING").add(nodeId));
+ case STABLE ->
+ client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("SETSLOT").add(slotId).add("STABLE"));
+ case NODE -> client.executeCommand(
+ new CommandArguments(Protocol.Command.CLUSTER).add("SETSLOT").add(slotId).add("NODE").add(nodeId));
};
this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
@@ -490,9 +499,11 @@ public void clusterSetSlot(@NonNull RedisClusterNode node, int slot, @NonNull Ad
public List clusterGetKeysInSlot(int slot, @NonNull Integer count) {
RedisClusterNode node = clusterGetNodeForSlot(slot);
+ String slotId = String.valueOf(slot);
- JedisClusterCommandCallback> command = jedis -> JedisConverters.stringListToByteList()
- .convert(jedis.clusterGetKeysInSlot(slot, nullSafeIntValue(count)));
+ JedisClusterCommandCallback> command = client -> (List) client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("GETKEYSINSLOT").add(slotId)
+ .add(String.valueOf(nullSafeIntValue(count))));
NodeResult> result = this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
@@ -506,7 +517,11 @@ private int nullSafeIntValue(@Nullable Integer value) {
@Override
public void clusterAddSlots(@NonNull RedisClusterNode node, int @NonNull... slots) {
- JedisClusterCommandCallback command = jedis -> jedis.clusterAddSlots(slots);
+ String[] args = Stream.concat(Stream.of("ADDSLOTS"), Arrays.stream(slots).mapToObj(String::valueOf))
+ .toArray(String[]::new);
+
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).addObjects(args));
this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
}
@@ -523,8 +538,10 @@ public void clusterAddSlots(@NonNull RedisClusterNode node, @NonNull SlotRange r
public Long clusterCountKeysInSlot(int slot) {
RedisClusterNode node = clusterGetNodeForSlot(slot);
+ String slotId = String.valueOf(slot);
- JedisClusterCommandCallback command = jedis -> jedis.clusterCountKeysInSlot(slot);
+ JedisClusterCommandCallback command = client -> (Long) client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("COUNTKEYSINSLOT").add(slotId));
return this.clusterCommandExecutor.executeCommandOnSingleNode(command, node).getValue();
}
@@ -532,7 +549,10 @@ public Long clusterCountKeysInSlot(int slot) {
@Override
public void clusterDeleteSlots(@NonNull RedisClusterNode node, int @NonNull... slots) {
- JedisClusterCommandCallback command = jedis -> jedis.clusterDelSlots(slots);
+ String[] args = Stream.concat(Stream.of("DELSLOTS"), Arrays.stream(slots).mapToObj(String::valueOf))
+ .toArray(String[]::new);
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).addObjects(args));
this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
}
@@ -553,7 +573,8 @@ public void clusterForget(@NonNull RedisClusterNode node) {
nodes.remove(nodeToRemove);
- JedisClusterCommandCallback command = jedis -> jedis.clusterForget(node.getId());
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("FORGET").add(node.getId()));
this.clusterCommandExecutor.executeCommandAsyncOnNodes(command, nodes);
}
@@ -566,8 +587,9 @@ public void clusterMeet(@NonNull RedisClusterNode node) {
Assert.hasText(node.getHost(), "Node to meet cluster must have a host");
Assert.isTrue(node.getPort() > 0, "Node to meet cluster must have a port greater 0");
- JedisClusterCommandCallback command = jedis -> jedis.clusterMeet(node.getRequiredHost(),
- node.getRequiredPort());
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("MEET").add(node.getRequiredHost())
+ .add(String.valueOf(node.getRequiredPort())));
this.clusterCommandExecutor.executeCommandOnAllNodes(command);
}
@@ -577,7 +599,8 @@ public void clusterReplicate(@NonNull RedisClusterNode master, @NonNull RedisClu
RedisClusterNode masterNode = this.topologyProvider.getTopology().lookup(master);
- JedisClusterCommandCallback command = jedis -> jedis.clusterReplicate(masterNode.getId());
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("REPLICATE").add(masterNode.getId()));
this.clusterCommandExecutor.executeCommandOnSingleNode(command, replica);
}
@@ -585,8 +608,8 @@ public void clusterReplicate(@NonNull RedisClusterNode master, @NonNull RedisClu
@Override
public Integer clusterGetSlotForKey(byte @NonNull [] key) {
- JedisClusterCommandCallback command = jedis -> Long
- .valueOf(jedis.clusterKeySlot(JedisConverters.toString(key))).intValue();
+ JedisClusterCommandCallback command = client -> ((Long) client.executeCommand(
+ new CommandArguments(Protocol.Command.CLUSTER).add("KEYSLOT").add(key))).intValue();
return this.clusterCommandExecutor.executeCommandOnArbitraryNode(command).getValue();
}
@@ -620,7 +643,8 @@ public Set clusterGetReplicas(@NonNull RedisClusterNode master
RedisClusterNode nodeToUse = this.topologyProvider.getTopology().lookup(master);
- JedisClusterCommandCallback> command = jedis -> jedis.clusterSlaves(nodeToUse.getId());
+ JedisClusterCommandCallback> command = client -> JedisConverters.toStrings((List) client
+ .executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("SLAVES").add(nodeToUse.getId())));
List clusterNodes = this.clusterCommandExecutor.executeCommandOnSingleNode(command, master).getValue();
@@ -630,8 +654,15 @@ public Set clusterGetReplicas(@NonNull RedisClusterNode master
@Override
public Map> clusterGetMasterReplicaMap() {
- JedisClusterCommandCallback> command = jedis -> JedisConverters
- .toSetOfRedisClusterNodes(jedis.clusterSlaves(jedis.clusterMyId()));
+ JedisClusterCommandCallback> command = client -> {
+
+ String myId = JedisConverters
+ .toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("MYID")));
+ List replicas = JedisConverters.toStrings(
+ (List) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("SLAVES").add(myId)));
+
+ return JedisConverters.toSetOfRedisClusterNodes(replicas);
+ };
Set activeMasterNodes = this.topologyProvider.getTopology().getActiveMasterNodes();
@@ -643,14 +674,14 @@ public Map> clusterGetMasterRepli
for (NodeResult> nodeResult : nodeResults) {
result.put(nodeResult.getNode(), nodeResult.getValue());
}
-
return result;
}
@Override
public ClusterInfo clusterGetClusterInfo() {
- JedisClusterCommandCallback command = Jedis::clusterInfo;
+ JedisClusterCommandCallback command = client -> JedisConverters
+ .toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("INFO")));
String source = this.clusterCommandExecutor.executeCommandOnArbitraryNode(command).getValue();
@@ -712,22 +743,22 @@ public void rewriteConfig() {
}
/**
- * {@link Jedis} specific {@link ClusterCommandCallback}.
+ * {@link UnifiedJedis} specific {@link ClusterCommandCallback}.
*
* @author Christoph Strobl
* @param
* @since 1.7
*/
- protected interface JedisClusterCommandCallback extends ClusterCommandCallback {}
+ protected interface JedisClusterCommandCallback extends ClusterCommandCallback {}
/**
- * {@link Jedis} specific {@link MultiKeyClusterCommandCallback}.
+ * {@link UnifiedJedis} specific {@link MultiKeyClusterCommandCallback}.
*
* @author Christoph Strobl
* @param
* @since 1.7
*/
- protected interface JedisMultiKeyClusterCommandCallback extends MultiKeyClusterCommandCallback {}
+ protected interface JedisMultiKeyClusterCommandCallback extends MultiKeyClusterCommandCallback {}
/**
* Jedis specific implementation of {@link ClusterNodeResourceProvider}.
@@ -763,19 +794,19 @@ static class JedisClusterNodeResourceProvider implements ClusterNodeResourceProv
@Override
@SuppressWarnings("unchecked")
- public Jedis getResourceForSpecificNode(RedisClusterNode node) {
+ public UnifiedJedis getResourceForSpecificNode(RedisClusterNode node) {
Assert.notNull(node, "Cannot get Pool for 'null' node");
ConnectionPool pool = getResourcePoolForSpecificNode(node);
if (pool != null) {
- return new Jedis(pool.getResource());
+ return new UnifiedJedisAdapter(new Jedis(pool.getResource()));
}
Connection connection = getConnectionForSpecificNode(node);
if (connection != null) {
- return new Jedis(connection);
+ return new UnifiedJedisAdapter(new Jedis(connection));
}
throw new DataAccessResourceFailureException("Node %s is unknown to cluster".formatted(node));
@@ -812,7 +843,7 @@ public Jedis getResourceForSpecificNode(RedisClusterNode node) {
@Override
public void returnResourceForSpecificNode(@NonNull RedisClusterNode node, @NonNull Object client) {
- ((Jedis) client).close();
+ ((UnifiedJedisAdapter) client).getJedis().close();
}
}
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterJsonCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterJsonCommands.java
index 183e8e1ab0..4cab6cd5a9 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterJsonCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterJsonCommands.java
@@ -21,6 +21,7 @@
import org.springframework.data.redis.connection.ClusterSlotHashUtil;
import org.springframework.data.redis.connection.RedisJsonCommands;
+import org.springframework.data.redis.connection.jedis.JedisClusterConnection.JedisMultiKeyClusterCommandCallback;
import org.springframework.data.redis.connection.json.JsonPath;
import org.springframework.util.Assert;
@@ -50,9 +51,10 @@ public List jsonMGet(JsonPath path, byte[]... keys) {
return super.jsonMGet(path, keys);
}
- List> results = connection.getClusterCommandExecutor().executeMultiKeyCommand((client, key) -> {
- return super.jsonMGet(path, key);
- }, Arrays.asList(keys)).resultsAsListSortBy(keys).stream().toList();
+ List> results = connection.getClusterCommandExecutor()
+ .executeMultiKeyCommand((JedisMultiKeyClusterCommandCallback>) (client,
+ key) -> toJsonBytes(client.jsonMGet(getPath(path), JedisConverters.toString(key))), Arrays.asList(keys))
+ .resultsAsListSortBy(keys).stream().toList();
List result = new ArrayList<>();
for (List list : results) {
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterScriptingCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterScriptingCommands.java
index 9d3c5edd05..d391db2ac4 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterScriptingCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterScriptingCommands.java
@@ -15,8 +15,6 @@
*/
package org.springframework.data.redis.connection.jedis;
-import redis.clients.jedis.Jedis;
-
import java.util.List;
import org.jspecify.annotations.NonNull;
@@ -25,6 +23,7 @@
import org.springframework.data.redis.connection.ClusterCommandExecutor;
import org.springframework.data.redis.connection.RedisScriptingCommands;
import org.springframework.util.Assert;
+import redis.clients.jedis.UnifiedJedis;
/**
* Cluster {@link RedisScriptingCommands} implementation for Jedis.
@@ -52,8 +51,8 @@ class JedisClusterScriptingCommands extends JedisScriptingCommands {
public void scriptFlush() {
try {
- connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback) Jedis::scriptFlush);
+ connection.getClusterCommandExecutor().executeCommandOnAllNodes(
+ (JedisClusterConnection.JedisClusterCommandCallback) UnifiedJedis::scriptFlush);
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
@@ -63,8 +62,8 @@ public void scriptFlush() {
public void scriptKill() {
try {
- connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback) Jedis::scriptKill);
+ connection.getClusterCommandExecutor().executeCommandOnAllNodes(
+ (JedisClusterConnection.JedisClusterCommandCallback) UnifiedJedis::scriptKill);
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
@@ -76,11 +75,11 @@ public String scriptLoad(byte @NonNull [] script) {
Assert.notNull(script, "Script must not be null");
try {
- ClusterCommandExecutor.MultiNodeResult multiNodeResult = connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes(
- (JedisClusterConnection.JedisClusterCommandCallback) client -> client.scriptLoad(script));
+ ClusterCommandExecutor.MultiNodeResult multiNodeResult = connection.getClusterCommandExecutor()
+ .executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback) client -> client
+ .scriptLoad(JedisConverters.toString(script)));
- return JedisConverters.toString(multiNodeResult.getFirstNonNullNotEmptyOrDefault(new byte[0]));
+ return multiNodeResult.getFirstNonNullNotEmptyOrDefault("");
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterServerCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterServerCommands.java
index a0402cf74a..e0926b962b 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterServerCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterServerCommands.java
@@ -15,7 +15,10 @@
*/
package org.springframework.data.redis.connection.jedis;
-import redis.clients.jedis.Jedis;
+import redis.clients.jedis.CommandArguments;
+import redis.clients.jedis.Protocol;
+import redis.clients.jedis.UnifiedJedis;
+import redis.clients.jedis.params.MigrateParams;
import java.util.ArrayList;
import java.util.Collection;
@@ -57,30 +60,36 @@ class JedisClusterServerCommands implements RedisClusterServerCommands {
@Override
public void bgReWriteAof(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::bgrewriteaof, node);
+ executeCommandOnSingleNode(client -> client.executeCommand(new CommandArguments(Protocol.Command.BGREWRITEAOF)),
+ node);
}
@Override
public void bgReWriteAof() {
connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::bgrewriteaof);
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.BGREWRITEAOF)));
}
@Override
public void bgSave() {
connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::bgsave);
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.BGSAVE)));
}
@Override
public void bgSave(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::bgsave, node);
+ executeCommandOnSingleNode(client -> client.executeCommand(new CommandArguments(Protocol.Command.BGSAVE)), node);
}
@Override
public Long lastSave() {
- List result = new ArrayList<>(executeCommandOnAllNodes(Jedis::lastsave).resultsAsList());
+ JedisClusterCommandCallback command = client -> (Long) client
+ .executeCommand(new CommandArguments(Protocol.Command.LASTSAVE));
+
+ List result = new ArrayList<>(executeCommandOnAllNodes(command).resultsAsList());
if (CollectionUtils.isEmpty(result)) {
return null;
@@ -92,23 +101,27 @@ public Long lastSave() {
@Override
public Long lastSave(@NonNull RedisClusterNode node) {
- return executeCommandOnSingleNode(Jedis::lastsave, node).getValue();
+
+ JedisClusterCommandCallback command = client -> (Long) client
+ .executeCommand(new CommandArguments(Protocol.Command.LASTSAVE));
+
+ return executeCommandOnSingleNode(command, node).getValue();
}
@Override
public void save() {
- executeCommandOnAllNodes(Jedis::save);
+ executeCommandOnAllNodes(client -> client.executeCommand(new CommandArguments(Protocol.Command.SAVE)));
}
@Override
public void save(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::save, node);
+ executeCommandOnSingleNode(client -> client.executeCommand(new CommandArguments(Protocol.Command.SAVE)), node);
}
@Override
public Long dbSize() {
- Collection dbSizes = executeCommandOnAllNodes(Jedis::dbSize).resultsAsList();
+ Collection dbSizes = executeCommandOnAllNodes(UnifiedJedis::dbSize).resultsAsList();
if (CollectionUtils.isEmpty(dbSizes)) {
return 0L;
@@ -123,49 +136,54 @@ public Long dbSize() {
@Override
public Long dbSize(@NonNull RedisClusterNode node) {
- return executeCommandOnSingleNode(Jedis::dbSize, node).getValue();
+ return executeCommandOnSingleNode(UnifiedJedis::dbSize, node).getValue();
}
@Override
public void flushDb() {
- executeCommandOnAllNodes(Jedis::flushDB);
+ executeCommandOnAllNodes(UnifiedJedis::flushDB);
}
@Override
public void flushDb(@NonNull FlushOption option) {
- executeCommandOnAllNodes(it -> it.flushDB(JedisConverters.toFlushMode(option)));
+ executeCommandOnAllNodes(client -> client.executeCommand(
+ new CommandArguments(Protocol.Command.FLUSHDB).add(JedisConverters.toFlushMode(option).name())));
}
@Override
public void flushDb(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::flushDB, node);
+ executeCommandOnSingleNode(UnifiedJedis::flushDB, node);
}
@Override
public void flushDb(@NonNull RedisClusterNode node, @NonNull FlushOption option) {
- executeCommandOnSingleNode(it -> it.flushDB(JedisConverters.toFlushMode(option)), node);
+ executeCommandOnSingleNode(client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.FLUSHDB).add(JedisConverters.toFlushMode(option).name())),
+ node);
}
@Override
public void flushAll() {
connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::flushAll);
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) UnifiedJedis::flushAll);
}
@Override
public void flushAll(@NonNull FlushOption option) {
- connection.getClusterCommandExecutor().executeCommandOnAllNodes(
- (JedisClusterCommandCallback) it -> it.flushAll(JedisConverters.toFlushMode(option)));
+ connection.getClusterCommandExecutor()
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client.executeCommand(
+ new CommandArguments(Protocol.Command.FLUSHALL).add(JedisConverters.toFlushMode(option).name())));
}
@Override
public void flushAll(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::flushAll, node);
+ executeCommandOnSingleNode(UnifiedJedis::flushAll, node);
}
@Override
public void flushAll(@NonNull RedisClusterNode node, @NonNull FlushOption option) {
- executeCommandOnSingleNode(it -> it.flushAll(JedisConverters.toFlushMode(option)), node);
+ executeCommandOnSingleNode(client -> client.executeCommand(
+ new CommandArguments(Protocol.Command.FLUSHALL).add(JedisConverters.toFlushMode(option).name())), node);
}
@Override
@@ -189,7 +207,7 @@ public Properties info() {
@Override
public Properties info(@NonNull RedisClusterNode node) {
- return JedisConverters.toProperties(executeCommandOnSingleNode(Jedis::info, node).getValue());
+ return JedisConverters.toProperties(executeCommandOnSingleNode(UnifiedJedis::info, node).getValue());
}
@Override
@@ -223,18 +241,14 @@ public Properties info(@NonNull RedisClusterNode node, @NonNull String section)
@Override
public void shutdown() {
- connection.getClusterCommandExecutor().executeCommandOnAllNodes((JedisClusterCommandCallback) jedis -> {
- jedis.shutdown();
- return null;
- });
+ connection.getClusterCommandExecutor()
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.SHUTDOWN)));
}
@Override
public void shutdown(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(jedis -> {
- jedis.shutdown();
- return null;
- }, node);
+ executeCommandOnSingleNode(client -> client.executeCommand(new CommandArguments(Protocol.Command.SHUTDOWN)), node);
}
@Override
@@ -308,41 +322,49 @@ public void setConfig(@NonNull RedisClusterNode node, @NonNull String param, @No
@Override
public void resetConfigStats() {
connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::configResetStat);
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("RESETSTAT")));
}
@Override
public void rewriteConfig() {
connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::configRewrite);
+ .executeCommandOnAllNodes((JedisClusterCommandCallback) client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("REWRITE")));
}
@Override
public void resetConfigStats(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::configResetStat, node);
+ executeCommandOnSingleNode(
+ client -> client.executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("RESETSTAT")), node);
}
@Override
public void rewriteConfig(@NonNull RedisClusterNode node) {
- executeCommandOnSingleNode(Jedis::configRewrite, node);
+ executeCommandOnSingleNode(
+ client -> client.executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("REWRITE")), node);
}
@Override
+ @SuppressWarnings("unchecked")
public Long time(@NonNull TimeUnit timeUnit) {
+ JedisClusterCommandCallback> command = client -> JedisConverters
+ .toStrings((List) client.executeCommand(new CommandArguments(Protocol.Command.TIME)));
+
return convertListOfStringToTime(
- connection.getClusterCommandExecutor()
- .executeCommandOnArbitraryNode((JedisClusterCommandCallback>) Jedis::time).getValue(),
- timeUnit);
+ connection.getClusterCommandExecutor().executeCommandOnArbitraryNode(command).getValue(), timeUnit);
}
@Override
+ @SuppressWarnings("unchecked")
public Long time(@NonNull RedisClusterNode node, @NonNull TimeUnit timeUnit) {
+ JedisClusterCommandCallback> command = client -> JedisConverters
+ .toStrings((List) client.executeCommand(new CommandArguments(Protocol.Command.TIME)));
+
return convertListOfStringToTime(
- connection.getClusterCommandExecutor()
- .executeCommandOnSingleNode((JedisClusterCommandCallback>) Jedis::time, node).getValue(),
- timeUnit);
+ connection.getClusterCommandExecutor().executeCommandOnSingleNode(command, node).getValue(), timeUnit);
}
@Override
@@ -351,7 +373,8 @@ public void killClient(@NonNull String host, int port) {
Assert.hasText(host, "Host for 'CLIENT KILL' must not be 'null' or 'empty'");
String hostAndPort = "%s:%d".formatted(host, port);
- JedisClusterCommandCallback command = client -> client.clientKill(hostAndPort);
+ JedisClusterCommandCallback command = client -> client
+ .executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("KILL").add(hostAndPort));
connection.getClusterCommandExecutor().executeCommandOnAllNodes(command);
}
@@ -369,8 +392,10 @@ public String getClientName() {
@Override
public List<@NonNull RedisClientInfo> getClientList() {
- Collection map = connection.getClusterCommandExecutor()
- .executeCommandOnAllNodes((JedisClusterCommandCallback) Jedis::clientList).resultsAsList();
+ JedisClusterCommandCallback command = client -> JedisConverters
+ .toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("LIST")));
+
+ Collection map = connection.getClusterCommandExecutor().executeCommandOnAllNodes(command).resultsAsList();
ArrayList result = new ArrayList<>();
for (String infos : map) {
@@ -382,8 +407,10 @@ public String getClientName() {
@Override
public List<@NonNull RedisClientInfo> getClientList(@NonNull RedisClusterNode node) {
- return JedisConverters
- .toListOfRedisClientInformation(executeCommandOnSingleNode(Jedis::clientList, node).getValue());
+ JedisClusterCommandCallback command = client -> JedisConverters
+ .toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("LIST")));
+
+ return JedisConverters.toListOfRedisClientInformation(executeCommandOnSingleNode(command, node).getValue());
}
@Override
@@ -414,8 +441,15 @@ public void migrate(byte @NonNull [] key, @NonNull RedisNode target, int dbIndex
RedisClusterNode node = connection.getTopologyProvider().getTopology().lookup(target.getRequiredHost(),
target.getRequiredPort());
+ MigrateParams params = new MigrateParams();
+ if (option == MigrateOption.COPY) {
+ params.copy();
+ } else if (option == MigrateOption.REPLACE) {
+ params.replace();
+ }
+
executeCommandOnSingleNode(
- client -> client.migrate(target.getRequiredHost(), target.getRequiredPort(), key, dbIndex, timeoutToUse), node);
+ client -> client.migrate(target.getRequiredHost(), target.getRequiredPort(), timeoutToUse, params, key), node);
}
private Long convertListOfStringToTime(List<@NonNull String> serverTimeInformation, TimeUnit timeUnit) {
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisConnection.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisConnection.java
index d939baf84c..f3f8cee2ec 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisConnection.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisConnection.java
@@ -342,7 +342,7 @@ public Object execute(@NonNull String command, byte @NonNull []... args) {
return null;
}
- return it.sendCommand(protocolCommand, args);
+ return it.executeCommand(new CommandArguments(protocolCommand).addObjects(args));
});
}
@@ -470,7 +470,7 @@ public byte[] echo(byte @NonNull [] message) {
Assert.notNull(message, "Message must not be null");
- return invoke().from(jedis -> jedis.sendCommand(Protocol.Command.ECHO, message))
+ return invoke().from(client -> client.executeCommand(new CommandArguments(Protocol.Command.ECHO).add(message)))
.get(response -> (byte[]) response);
}
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisJsonCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisJsonCommands.java
index 9275973a5f..4d39d677a7 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisJsonCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisJsonCommands.java
@@ -23,6 +23,7 @@
import java.util.List;
import java.util.stream.Stream;
+import org.json.JSONArray;
import org.jspecify.annotations.NonNull;
import org.jspecify.annotations.NullUnmarked;
@@ -163,7 +164,7 @@ public List jsonMGet(@NonNull JsonPath path, byte @NonNull [] @NonNull..
return connection.invoke()
.from(UnifiedJedis::jsonMGet, RedisJsonPipelineCommands::jsonMGet, getPath(path), stringKeys)
- .get(jsonArrList -> jsonArrList.stream().map(arr -> arr != null ? arr.toString().getBytes(StandardCharsets.UTF_8) : null).toList());
+ .get(JedisJsonCommands::toJsonBytes);
}
@Override
@@ -226,4 +227,9 @@ static Path2 getPath(JsonPath path) {
return Path2.of(path.asString());
}
+ static List toJsonBytes(List jsonArrList) {
+ return jsonArrList.stream().map(arr -> arr != null ? arr.toString().getBytes(StandardCharsets.UTF_8) : null)
+ .toList();
+ }
+
}
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisKeyCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisKeyCommands.java
index b637418bcd..c8563f9b51 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisKeyCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisKeyCommands.java
@@ -15,6 +15,7 @@
*/
package org.springframework.data.redis.connection.jedis;
+import redis.clients.jedis.CommandArguments;
import redis.clients.jedis.Protocol;
import redis.clients.jedis.args.ExpiryOption;
import redis.clients.jedis.commands.JedisBinaryCommands;
@@ -327,7 +328,8 @@ public Boolean move(byte @NonNull [] key, int dbIndex) {
Assert.notNull(key, "Key must not be null");
- return connection.invoke().from(j -> j.sendCommand(Protocol.Command.MOVE, key, Protocol.toByteArray(dbIndex)))
+ return connection.invoke().from(
+ j -> j.executeCommand(new CommandArguments(Protocol.Command.MOVE).add(key).add(Protocol.toByteArray(dbIndex))))
.get(response -> JedisConverters.longToBoolean().convert(((Long) response)));
}
diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisServerCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisServerCommands.java
index 0c40c5dc76..3b77dd565a 100644
--- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisServerCommands.java
+++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisServerCommands.java
@@ -15,6 +15,7 @@
*/
package org.springframework.data.redis.connection.jedis;
+import redis.clients.jedis.CommandArguments;
import redis.clients.jedis.Protocol;
import redis.clients.jedis.UnifiedJedis;
import redis.clients.jedis.params.MigrateParams;
@@ -50,23 +51,23 @@ class JedisServerCommands implements RedisServerCommands {
@Override
public void bgReWriteAof() {
- connection.invoke().just(j -> j.sendCommand(Protocol.Command.BGREWRITEAOF));
+ connection.invoke().just(j -> j.executeCommand(new CommandArguments(Protocol.Command.BGREWRITEAOF)));
}
@Override
public void bgSave() {
- connection.invoke().just(j -> j.sendCommand(Protocol.Command.BGSAVE));
+ connection.invoke().just(j -> j.executeCommand(new CommandArguments(Protocol.Command.BGSAVE)));
}
@Override
public Long lastSave() {
- return connection.invoke().from(j -> j.sendCommand(Protocol.Command.LASTSAVE))
+ return connection.invoke().from(j -> j.executeCommand(new CommandArguments(Protocol.Command.LASTSAVE)))
.get(response -> (Long) response);
}
@Override
public void save() {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.SAVE));
+ connection.invokeStatus().just(j -> j.executeCommand(new CommandArguments(Protocol.Command.SAVE)));
}
@Override
@@ -81,7 +82,8 @@ public void flushDb() {
@Override
public void flushDb(@NonNull FlushOption option) {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.FLUSHDB, JedisConverters.toFlushMode(option).name()));
+ connection.invokeStatus().just(j -> j.executeCommand(
+ new CommandArguments(Protocol.Command.FLUSHDB).add(JedisConverters.toFlushMode(option).name())));
}
@Override
@@ -91,7 +93,8 @@ public void flushAll() {
@Override
public void flushAll(@NonNull FlushOption option) {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.FLUSHALL, JedisConverters.toFlushMode(option).name()));
+ connection.invokeStatus().just(j -> j.executeCommand(
+ new CommandArguments(Protocol.Command.FLUSHALL).add(JedisConverters.toFlushMode(option).name())));
}
@Override
@@ -109,7 +112,7 @@ public Properties info(@NonNull String section) {
@Override
public void shutdown() {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.SHUTDOWN));
+ connection.invokeStatus().just(j -> j.executeCommand(new CommandArguments(Protocol.Command.SHUTDOWN)));
}
@Override
@@ -121,7 +124,8 @@ public void shutdown(@Nullable ShutdownOption option) {
}
String saveOption = (option == ShutdownOption.NOSAVE) ? "NOSAVE" : "SAVE";
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.SHUTDOWN, saveOption));
+ connection.invokeStatus()
+ .just(j -> j.executeCommand(new CommandArguments(Protocol.Command.SHUTDOWN).add(saveOption)));
}
@Override
@@ -130,7 +134,8 @@ public Properties getConfig(@NonNull String pattern) {
Assert.notNull(pattern, "Pattern must not be null");
- return connection.invoke().from(j -> j.sendCommand(Protocol.Command.CONFIG, "GET", pattern))
+ return connection.invoke()
+ .from(j -> j.executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("GET").add(pattern)))
.get(response -> {
List list = (List) response;
Properties props = new Properties();
@@ -166,12 +171,13 @@ public void setConfig(@NonNull String param, @NonNull String value) {
@Override
public void resetConfigStats() {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.CONFIG, "RESETSTAT"));
+ connection.invokeStatus()
+ .just(j -> j.executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("RESETSTAT")));
}
@Override
public void rewriteConfig() {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.CONFIG, "REWRITE"));
+ connection.invokeStatus().just(j -> j.executeCommand(new CommandArguments(Protocol.Command.CONFIG).add("REWRITE")));
}
@Override
@@ -180,7 +186,7 @@ public Long time(@NonNull TimeUnit timeUnit) {
Assert.notNull(timeUnit, "TimeUnit must not be null");
- return connection.invoke().from(j -> j.sendCommand(Protocol.Command.TIME))
+ return connection.invoke().from(j -> j.executeCommand(new CommandArguments(Protocol.Command.TIME)))
.get(response -> {
List list = (List) response;
List timeList = new ArrayList<>();
@@ -196,7 +202,8 @@ public void killClient(@NonNull String host, int port) {
Assert.hasText(host, "Host for 'CLIENT KILL' must not be 'null' or 'empty'");
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.CLIENT, "KILL", "%s:%s".formatted(host, port)));
+ connection.invokeStatus().just(j -> j
+ .executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("KILL").add("%s:%s".formatted(host, port))));
}
@Override
@@ -204,18 +211,21 @@ public void setClientName(byte @NonNull [] name) {
Assert.notNull(name, "Name must not be null");
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.CLIENT, "SETNAME".getBytes(), name));
+ connection.invokeStatus()
+ .just(j -> j.executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("SETNAME".getBytes()).add(name)));
}
@Override
public String getClientName() {
- return connection.invokeStatus().from(j -> j.sendCommand(Protocol.Command.CLIENT, "GETNAME"))
+ return connection.invokeStatus()
+ .from(j -> j.executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("GETNAME")))
.get(response -> new String((byte[]) response));
}
@Override
public List<@NonNull RedisClientInfo> getClientList() {
- return connection.invokeStatus().from(j -> j.sendCommand(Protocol.Command.CLIENT, "LIST"))
+ return connection.invokeStatus()
+ .from(j -> j.executeCommand(new CommandArguments(Protocol.Command.CLIENT).add("LIST")))
.get(response -> JedisConverters.toListOfRedisClientInformation(new String((byte[]) response)));
}
@@ -224,12 +234,14 @@ public void replicaOf(@NonNull String host, int port) {
Assert.hasText(host, "Host must not be null for 'REPLICAOF' command");
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.REPLICAOF, host, String.valueOf(port)));
+ connection.invokeStatus().just(
+ j -> j.executeCommand(new CommandArguments(Protocol.Command.REPLICAOF).add(host).add(String.valueOf(port))));
}
@Override
public void replicaOfNoOne() {
- connection.invokeStatus().just(j -> j.sendCommand(Protocol.Command.REPLICAOF, "NO", "ONE"));
+ connection.invokeStatus()
+ .just(j -> j.executeCommand(new CommandArguments(Protocol.Command.REPLICAOF).add("NO").add("ONE")));
}