Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -65,7 +69,7 @@
* {@link RedisClusterConnection} implementation on top of {@link RedisClusterClient}.
* <p>
* Uses the native {@link RedisClusterClient} api where possible and falls back to direct node communication using
* {@link Jedis} where needed.
* {@link UnifiedJedis} where needed.
* <p>
* Pipelines and transactions are not supported in cluster mode. This class is not Thread-safe and instances should not
* be shared across threads.
Expand Down Expand Up @@ -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<Object> commandCallback = jedis -> jedis
.sendCommand(JedisClientUtils.getCommand(command), args);
JedisClusterCommandCallback<Object> commandCallback = client -> client
.executeCommand(new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(args));

return this.clusterCommandExecutor.executeCommandOnArbitraryNode(commandCallback).getValue();
}
Expand All @@ -268,8 +272,8 @@ public <T> T execute(@NonNull String command, byte @NonNull [] key, @NonNull Col

RedisClusterNode keyMaster = this.topologyProvider.getTopology().getKeyServingMasterNode(key);

JedisClusterCommandCallback<T> commandCallback = jedis -> (T) jedis
.sendCommand(JedisClientUtils.getCommand(command), commandArgs);
JedisClusterCommandCallback<T> commandCallback = client -> (T) client
.executeCommand(new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(commandArgs));

return this.clusterCommandExecutor.executeCommandOnSingleNode(commandCallback, keyMaster).getValue();
}
Expand Down Expand Up @@ -318,8 +322,8 @@ public <T> List<T> execute(@NonNull String command, @NonNull Collection<byte @No
Assert.notNull(keys, "Key must not be null");
Assert.notNull(args, "Args must not be null");

JedisMultiKeyClusterCommandCallback<T> commandCallback = (jedis,
key) -> (T) jedis.sendCommand(JedisClientUtils.getCommand(command), getCommandArguments(key, args));
JedisMultiKeyClusterCommandCallback<T> commandCallback = (jedis, key) -> (T) jedis.executeCommand(
new CommandArguments(JedisClientUtils.getCommand(command)).addObjects(getCommandArguments(key, args)));

return this.clusterCommandExecutor.executeMultiKeyCommand(commandCallback, keys).resultsAsList();
}
Expand Down Expand Up @@ -450,15 +454,15 @@ public byte[] echo(byte @NonNull [] message) {
@Override
public String ping() {

JedisClusterCommandCallback<String> command = Jedis::ping;
JedisClusterCommandCallback<String> command = UnifiedJedis::ping;

return !this.clusterCommandExecutor.executeCommandOnAllNodes(command).resultsAsList().isEmpty() ? "PONG" : null;
}

@Override
public String ping(@NonNull RedisClusterNode node) {

JedisClusterCommandCallback<String> command = Jedis::ping;
JedisClusterCommandCallback<String> command = UnifiedJedis::ping;

return this.clusterCommandExecutor.executeCommandOnSingleNode(command, node).getValue();
}
Expand All @@ -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<String> 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<Object> 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);
Expand All @@ -490,9 +499,11 @@ public void clusterSetSlot(@NonNull RedisClusterNode node, int slot, @NonNull Ad
public List<byte[]> clusterGetKeysInSlot(int slot, @NonNull Integer count) {

RedisClusterNode node = clusterGetNodeForSlot(slot);
String slotId = String.valueOf(slot);

JedisClusterCommandCallback<List<byte[]>> command = jedis -> JedisConverters.stringListToByteList()
.convert(jedis.clusterGetKeysInSlot(slot, nullSafeIntValue(count)));
JedisClusterCommandCallback<List<byte[]>> command = client -> (List<byte[]>) client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("GETKEYSINSLOT").add(slotId)
.add(String.valueOf(nullSafeIntValue(count))));

NodeResult<List<byte[]>> result = this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);

Expand All @@ -506,7 +517,11 @@ private int nullSafeIntValue(@Nullable Integer value) {
@Override
public void clusterAddSlots(@NonNull RedisClusterNode node, int @NonNull... slots) {

JedisClusterCommandCallback<String> command = jedis -> jedis.clusterAddSlots(slots);
String[] args = Stream.concat(Stream.of("ADDSLOTS"), Arrays.stream(slots).mapToObj(String::valueOf))
.toArray(String[]::new);

JedisClusterCommandCallback<Object> command = client -> client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).addObjects(args));

this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
}
Expand All @@ -523,16 +538,21 @@ public void clusterAddSlots(@NonNull RedisClusterNode node, @NonNull SlotRange r
public Long clusterCountKeysInSlot(int slot) {

RedisClusterNode node = clusterGetNodeForSlot(slot);
String slotId = String.valueOf(slot);

JedisClusterCommandCallback<Long> command = jedis -> jedis.clusterCountKeysInSlot(slot);
JedisClusterCommandCallback<Long> command = client -> (Long) client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("COUNTKEYSINSLOT").add(slotId));

return this.clusterCommandExecutor.executeCommandOnSingleNode(command, node).getValue();
}

@Override
public void clusterDeleteSlots(@NonNull RedisClusterNode node, int @NonNull... slots) {

JedisClusterCommandCallback<String> command = jedis -> jedis.clusterDelSlots(slots);
String[] args = Stream.concat(Stream.of("DELSLOTS"), Arrays.stream(slots).mapToObj(String::valueOf))
.toArray(String[]::new);
JedisClusterCommandCallback<Object> command = client -> client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).addObjects(args));

this.clusterCommandExecutor.executeCommandOnSingleNode(command, node);
}
Expand All @@ -553,7 +573,8 @@ public void clusterForget(@NonNull RedisClusterNode node) {

nodes.remove(nodeToRemove);

JedisClusterCommandCallback<String> command = jedis -> jedis.clusterForget(node.getId());
JedisClusterCommandCallback<Object> command = client -> client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("FORGET").add(node.getId()));

this.clusterCommandExecutor.executeCommandAsyncOnNodes(command, nodes);
}
Expand All @@ -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<String> command = jedis -> jedis.clusterMeet(node.getRequiredHost(),
node.getRequiredPort());
JedisClusterCommandCallback<Object> command = client -> client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("MEET").add(node.getRequiredHost())
.add(String.valueOf(node.getRequiredPort())));

this.clusterCommandExecutor.executeCommandOnAllNodes(command);
}
Expand All @@ -577,16 +599,17 @@ public void clusterReplicate(@NonNull RedisClusterNode master, @NonNull RedisClu

RedisClusterNode masterNode = this.topologyProvider.getTopology().lookup(master);

JedisClusterCommandCallback<String> command = jedis -> jedis.clusterReplicate(masterNode.getId());
JedisClusterCommandCallback<Object> command = client -> client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("REPLICATE").add(masterNode.getId()));

this.clusterCommandExecutor.executeCommandOnSingleNode(command, replica);
}

@Override
public Integer clusterGetSlotForKey(byte @NonNull [] key) {

JedisClusterCommandCallback<Integer> command = jedis -> Long
.valueOf(jedis.clusterKeySlot(JedisConverters.toString(key))).intValue();
JedisClusterCommandCallback<Integer> command = client -> ((Long) client.executeCommand(
new CommandArguments(Protocol.Command.CLUSTER).add("KEYSLOT").add(key))).intValue();

return this.clusterCommandExecutor.executeCommandOnArbitraryNode(command).getValue();
}
Expand Down Expand Up @@ -620,7 +643,8 @@ public Set<RedisClusterNode> clusterGetReplicas(@NonNull RedisClusterNode master

RedisClusterNode nodeToUse = this.topologyProvider.getTopology().lookup(master);

JedisClusterCommandCallback<List<String>> command = jedis -> jedis.clusterSlaves(nodeToUse.getId());
JedisClusterCommandCallback<List<String>> command = client -> JedisConverters.toStrings((List<byte[]>) client
.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("SLAVES").add(nodeToUse.getId())));

List<String> clusterNodes = this.clusterCommandExecutor.executeCommandOnSingleNode(command, master).getValue();

Expand All @@ -630,8 +654,15 @@ public Set<RedisClusterNode> clusterGetReplicas(@NonNull RedisClusterNode master
@Override
public Map<RedisClusterNode, Collection<RedisClusterNode>> clusterGetMasterReplicaMap() {

JedisClusterCommandCallback<Collection<RedisClusterNode>> command = jedis -> JedisConverters
.toSetOfRedisClusterNodes(jedis.clusterSlaves(jedis.clusterMyId()));
JedisClusterCommandCallback<Collection<RedisClusterNode>> command = client -> {

String myId = JedisConverters
.toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("MYID")));
List<String> replicas = JedisConverters.toStrings(
(List<byte[]>) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("SLAVES").add(myId)));

return JedisConverters.toSetOfRedisClusterNodes(replicas);
};

Set<RedisClusterNode> activeMasterNodes = this.topologyProvider.getTopology().getActiveMasterNodes();

Expand All @@ -643,14 +674,14 @@ public Map<RedisClusterNode, Collection<RedisClusterNode>> clusterGetMasterRepli
for (NodeResult<Collection<RedisClusterNode>> nodeResult : nodeResults) {
result.put(nodeResult.getNode(), nodeResult.getValue());
}

return result;
}

@Override
public ClusterInfo clusterGetClusterInfo() {

JedisClusterCommandCallback<String> command = Jedis::clusterInfo;
JedisClusterCommandCallback<String> command = client -> JedisConverters
.toString((byte[]) client.executeCommand(new CommandArguments(Protocol.Command.CLUSTER).add("INFO")));

String source = this.clusterCommandExecutor.executeCommandOnArbitraryNode(command).getValue();

Expand Down Expand Up @@ -712,22 +743,22 @@ public void rewriteConfig() {
}

/**
* {@link Jedis} specific {@link ClusterCommandCallback}.
* {@link UnifiedJedis} specific {@link ClusterCommandCallback}.
*
* @author Christoph Strobl
* @param <T>
* @since 1.7
*/
protected interface JedisClusterCommandCallback<T> extends ClusterCommandCallback<Jedis, T> {}
protected interface JedisClusterCommandCallback<T> extends ClusterCommandCallback<UnifiedJedis, T> {}

/**
* {@link Jedis} specific {@link MultiKeyClusterCommandCallback}.
* {@link UnifiedJedis} specific {@link MultiKeyClusterCommandCallback}.
*
* @author Christoph Strobl
* @param <T>
* @since 1.7
*/
protected interface JedisMultiKeyClusterCommandCallback<T> extends MultiKeyClusterCommandCallback<Jedis, T> {}
protected interface JedisMultiKeyClusterCommandCallback<T> extends MultiKeyClusterCommandCallback<UnifiedJedis, T> {}

/**
* Jedis specific implementation of {@link ClusterNodeResourceProvider}.
Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -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();
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

Expand Down Expand Up @@ -50,9 +51,10 @@ public List<byte[]> jsonMGet(JsonPath path, byte[]... keys) {
return super.jsonMGet(path, keys);
}

List<List<byte[]>> results = connection.getClusterCommandExecutor().executeMultiKeyCommand((client, key) -> {
return super.jsonMGet(path, key);
}, Arrays.asList(keys)).resultsAsListSortBy(keys).stream().toList();
List<List<byte[]>> results = connection.getClusterCommandExecutor()
.executeMultiKeyCommand((JedisMultiKeyClusterCommandCallback<List<byte[]>>) (client,
key) -> toJsonBytes(client.jsonMGet(getPath(path), JedisConverters.toString(key))), Arrays.asList(keys))
.resultsAsListSortBy(keys).stream().toList();

List<byte[]> result = new ArrayList<>();
for (List<byte[]> list : results) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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.
Expand Down Expand Up @@ -52,8 +51,8 @@ class JedisClusterScriptingCommands extends JedisScriptingCommands {
public void scriptFlush() {

try {
connection.getClusterCommandExecutor()
.executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback<String>) Jedis::scriptFlush);
connection.getClusterCommandExecutor().executeCommandOnAllNodes(
(JedisClusterConnection.JedisClusterCommandCallback<String>) UnifiedJedis::scriptFlush);
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
Expand All @@ -63,8 +62,8 @@ public void scriptFlush() {
public void scriptKill() {

try {
connection.getClusterCommandExecutor()
.executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback<String>) Jedis::scriptKill);
connection.getClusterCommandExecutor().executeCommandOnAllNodes(
(JedisClusterConnection.JedisClusterCommandCallback<String>) UnifiedJedis::scriptKill);
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
Expand All @@ -76,11 +75,11 @@ public String scriptLoad(byte @NonNull [] script) {
Assert.notNull(script, "Script must not be null");

try {
ClusterCommandExecutor.MultiNodeResult<byte[]> multiNodeResult = connection.getClusterCommandExecutor()
.executeCommandOnAllNodes(
(JedisClusterConnection.JedisClusterCommandCallback<byte[]>) client -> client.scriptLoad(script));
ClusterCommandExecutor.MultiNodeResult<String> multiNodeResult = connection.getClusterCommandExecutor()
.executeCommandOnAllNodes((JedisClusterConnection.JedisClusterCommandCallback<String>) client -> client
.scriptLoad(JedisConverters.toString(script)));

return JedisConverters.toString(multiNodeResult.getFirstNonNullNotEmptyOrDefault(new byte[0]));
return multiNodeResult.getFirstNonNullNotEmptyOrDefault("");
} catch (Exception ex) {
throw connection.convertJedisAccessException(ex);
}
Expand Down
Loading
Loading