From 78e5b89d1b66cfa6a8f13fa31f7847c6af8acb28 Mon Sep 17 00:00:00 2001 From: Mark Paluch Date: Thu, 11 May 2017 14:15:47 +0200 Subject: [PATCH] DATAREDIS-647 - Correct SDIFF & SINTER behavior. We now consider the last key in SDIFF command execution on Redis Cluster to correctly compute the set difference. For reactive connections we collect the result sets entirely before applying further computation. Previously, the zip function could return a previous view of the result that caused too many/too few results. Original Pull Request: #250 --- .../jedis/JedisClusterSetCommands.java | 2 +- .../lettuce/LettuceClusterSetCommands.java | 2 +- .../LettuceReactiveClusterSetCommands.java | 15 +++++++--- .../jedis/JedisClusterConnectionTests.java | 30 ++++++------------- .../LettuceClusterConnectionTests.java | 10 +++---- .../LettuceReactiveSetCommandsTests.java | 12 ++++++-- 6 files changed, 36 insertions(+), 35 deletions(-) diff --git a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterSetCommands.java b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterSetCommands.java index 06b1f2b51..8b9ef2ecd 100644 --- a/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterSetCommands.java +++ b/src/main/java/org/springframework/data/redis/connection/jedis/JedisClusterSetCommands.java @@ -276,7 +276,7 @@ class JedisClusterSetCommands implements RedisSetCommands { } byte[] source = keys[0]; - byte[][] others = Arrays.copyOfRange(keys, 1, keys.length - 1); + byte[][] others = Arrays.copyOfRange(keys, 1, keys.length); ByteArraySet values = new ByteArraySet(sMembers(source)); Collection> resultList = connection.getClusterCommandExecutor() diff --git a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceClusterSetCommands.java b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceClusterSetCommands.java index df5810f1d..e3e500bdb 100644 --- a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceClusterSetCommands.java +++ b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceClusterSetCommands.java @@ -177,7 +177,7 @@ class LettuceClusterSetCommands extends LettuceSetCommands { } byte[] source = keys[0]; - byte[][] others = Arrays.copyOfRange(keys, 1, keys.length - 1); + byte[][] others = Arrays.copyOfRange(keys, 1, keys.length); ByteArraySet values = new ByteArraySet(sMembers(source)); Collection> nodeResult = connection.getClusterCommandExecutor() diff --git a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveClusterSetCommands.java b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveClusterSetCommands.java index 0f597b29e..6bda5093b 100644 --- a/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveClusterSetCommands.java +++ b/src/main/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveClusterSetCommands.java @@ -115,9 +115,12 @@ class LettuceReactiveClusterSetCommands extends LettuceReactiveSetCommands imple intersectingSets.add(cmd.smembers(command.getKeys().get(i)).distinct().collectList()); } - Flux> result = Flux.zip(sourceSet, Flux.merge(intersectingSets), (source, intersecting) -> { + Flux> result = Flux.zip(sourceSet, Flux.merge(intersectingSets).collectList(), + (source, intersectings) -> { - source.retainAll(intersecting); + for (List intersecting : intersectings) { + source.retainAll(intersecting); + } return source; }); @@ -172,9 +175,13 @@ class LettuceReactiveClusterSetCommands extends LettuceReactiveSetCommands imple intersectingSets.add(cmd.smembers(command.getKeys().get(i)).distinct().collectList()); } - Flux> result = Flux.zip(sourceSet, Flux.merge(intersectingSets), (source, intersecting) -> { + Flux> result = Flux.zip(sourceSet, Flux.merge(intersectingSets).collectList(), + (source, intersectings) -> { + + for (List intersecting : intersectings) { + source.removeAll(intersecting); + } - source.removeAll(intersecting); return source; }); diff --git a/src/test/java/org/springframework/data/redis/connection/jedis/JedisClusterConnectionTests.java b/src/test/java/org/springframework/data/redis/connection/jedis/JedisClusterConnectionTests.java index 58212aed1..6966befda 100644 --- a/src/test/java/org/springframework/data/redis/connection/jedis/JedisClusterConnectionTests.java +++ b/src/test/java/org/springframework/data/redis/connection/jedis/JedisClusterConnectionTests.java @@ -15,30 +15,20 @@ */ package org.springframework.data.redis.connection.jedis; -import static org.hamcrest.CoreMatchers.*; -import static org.hamcrest.collection.IsCollectionWithSize.*; -import static org.hamcrest.collection.IsIterableContainingInOrder.*; -import static org.hamcrest.core.Is.is; -import static org.hamcrest.number.IsCloseTo.*; +import static org.hamcrest.Matchers.*; import static org.junit.Assert.*; import static org.springframework.data.redis.connection.ClusterTestVariables.*; import static org.springframework.data.redis.connection.RedisGeoCommands.DistanceUnit.*; import static org.springframework.data.redis.connection.RedisGeoCommands.GeoRadiusCommandArgs.*; import static org.springframework.data.redis.core.ScanOptions.*; +import redis.clients.jedis.HostAndPort; +import redis.clients.jedis.JedisCluster; +import redis.clients.jedis.JedisPool; + import java.io.IOException; import java.nio.charset.Charset; -import java.util.Arrays; -import java.util.Collection; -import java.util.Collections; -import java.util.HashMap; -import java.util.HashSet; -import java.util.LinkedHashMap; -import java.util.List; -import java.util.ListIterator; -import java.util.Map; -import java.util.Properties; -import java.util.Set; +import java.util.*; import java.util.concurrent.TimeUnit; import org.junit.After; @@ -71,10 +61,6 @@ import org.springframework.data.redis.test.util.MinimumRedisVersionRule; import org.springframework.data.redis.test.util.RedisClusterRule; import org.springframework.test.annotation.IfProfileValue; -import redis.clients.jedis.HostAndPort; -import redis.clients.jedis.JedisCluster; -import redis.clients.jedis.JedisPool; - /** * @author Christoph Strobl * @author Mark Paluch @@ -1072,13 +1058,15 @@ public class JedisClusterConnectionTests implements ClusterConnectionTests { assertThat(clusterConnection.sDiff(SAME_SLOT_KEY_1_BYTES, SAME_SLOT_KEY_2_BYTES), hasItems(VALUE_1_BYTES)); } - @Test // DATAREDIS-315 + @Test // DATAREDIS-315, DATAREDIS-647 public void sDiffShouldWorkWhenKeysNotMapToSameSlot() { nativeConnection.sadd(KEY_1_BYTES, VALUE_1_BYTES, VALUE_2_BYTES); nativeConnection.sadd(KEY_2_BYTES, VALUE_2_BYTES, VALUE_3_BYTES); + nativeConnection.sadd(KEY_3_BYTES, VALUE_1_BYTES, VALUE_3_BYTES); assertThat(clusterConnection.sDiff(KEY_1_BYTES, KEY_2_BYTES), hasItems(VALUE_1_BYTES)); + assertThat(clusterConnection.sDiff(KEY_1_BYTES, KEY_2_BYTES, KEY_3_BYTES), is(empty())); } @Test // DATAREDIS-315 diff --git a/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceClusterConnectionTests.java b/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceClusterConnectionTests.java index 5f7fc5c99..1e4860e4a 100644 --- a/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceClusterConnectionTests.java +++ b/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceClusterConnectionTests.java @@ -15,11 +15,7 @@ */ package org.springframework.data.redis.connection.lettuce; -import static org.hamcrest.CoreMatchers.*; -import static org.hamcrest.collection.IsCollectionWithSize.*; -import static org.hamcrest.collection.IsIterableContainingInOrder.*; -import static org.hamcrest.core.Is.is; -import static org.hamcrest.number.IsCloseTo.*; +import static org.hamcrest.Matchers.*; import static org.junit.Assert.*; import static org.springframework.data.redis.connection.ClusterTestVariables.*; import static org.springframework.data.redis.connection.RedisGeoCommands.DistanceUnit.*; @@ -1078,13 +1074,15 @@ public class LettuceClusterConnectionTests implements ClusterConnectionTests { assertThat(clusterConnection.sDiff(SAME_SLOT_KEY_1_BYTES, SAME_SLOT_KEY_2_BYTES), hasItems(VALUE_1_BYTES)); } - @Test // DATAREDIS-315 + @Test // DATAREDIS-315, DATAREDIS-647 public void sDiffShouldWorkWhenKeysNotMapToSameSlot() { nativeConnection.sadd(KEY_1, VALUE_1, VALUE_2); nativeConnection.sadd(KEY_2, VALUE_2, VALUE_3); + nativeConnection.sadd(KEY_3, VALUE_1, VALUE_3); assertThat(clusterConnection.sDiff(KEY_1_BYTES, KEY_2_BYTES), hasItems(VALUE_1_BYTES)); + assertThat(clusterConnection.sDiff(KEY_1_BYTES, KEY_2_BYTES, KEY_3_BYTES), is(empty())); } @Test // DATAREDIS-315 diff --git a/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveSetCommandsTests.java b/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveSetCommandsTests.java index 1d8ca5380..84bb5ce3b 100644 --- a/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveSetCommandsTests.java +++ b/src/test/java/org/springframework/data/redis/connection/lettuce/LettuceReactiveSetCommandsTests.java @@ -129,15 +129,19 @@ public class LettuceReactiveSetCommandsTests extends LettuceReactiveCommandsTest assertThat(connection.setCommands().sIsMember(KEY_1_BBUFFER, VALUE_3_BBUFFER).block(), is(false)); } - @Test // DATAREDIS-525 + @Test // DATAREDIS-525, DATAREDIS-647 public void sInterShouldIntersectSetsCorrectly() { nativeCommands.sadd(KEY_1, VALUE_1, VALUE_2); nativeCommands.sadd(KEY_2, VALUE_2, VALUE_3); + nativeCommands.sadd(KEY_3, VALUE_1, VALUE_3); StepVerifier.create(connection.setCommands().sInter(Arrays.asList(KEY_1_BBUFFER, KEY_2_BBUFFER))) // .expectNext(VALUE_2_BBUFFER) // .verifyComplete(); + + StepVerifier.create(connection.setCommands().sInter(Arrays.asList(KEY_1_BBUFFER, KEY_2_BBUFFER, KEY_3_BBUFFER))) // + .verifyComplete(); } @Test // DATAREDIS-525 @@ -172,15 +176,19 @@ public class LettuceReactiveSetCommandsTests extends LettuceReactiveCommandsTest is(3L)); } - @Test // DATAREDIS-525 + @Test // DATAREDIS-525, DATAREDIS-647 public void sDiffShouldBeExcecutedCorrectly() { nativeCommands.sadd(KEY_1, VALUE_1, VALUE_2); nativeCommands.sadd(KEY_2, VALUE_2, VALUE_3); + nativeCommands.sadd(KEY_3, VALUE_2, VALUE_1); StepVerifier.create(connection.setCommands().sDiff(Arrays.asList(KEY_1_BBUFFER, KEY_2_BBUFFER))) // .expectNext(VALUE_1_BBUFFER) // .verifyComplete(); + + StepVerifier.create(connection.setCommands().sDiff(Arrays.asList(KEY_1_BBUFFER, KEY_2_BBUFFER, KEY_3_BBUFFER))) // + .verifyComplete(); } @Test // DATAREDIS-525