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"))); }