From 0f7607e515e5c1a2fd4646cc34065af3287d5b32 Mon Sep 17 00:00:00 2001 From: Christoph Strobl Date: Tue, 10 Jan 2017 13:33:47 +0100 Subject: [PATCH] DATAREDIS-589 - Move secondary index cleanup to MappingExpirationListener. We now no longer rely on ApplicationEvents captured in the RedisKeyValueAdapter for performing cleanup operations for expired keys, but do this along with the phantom key removal. This removes a flaw when initializing a non repository related KeyspaceEventListener publishing events that actually are unrelated to the Adapter. Additionally upgraded test infrastructure to utilize Redis 3.2.6 with disabled protected-mode and enabled keyspace-events. Original pull request: #232. --- Makefile | 8 ++- .../data/redis/core/RedisKeyExpiredEvent.java | 16 +++++- .../data/redis/core/RedisKeyValueAdapter.java | 44 +++++++--------- .../data/redis/core/RedisKeyspaceEvent.java | 26 +++++++++- .../redis/core/RedisKeyValueAdapterTests.java | 52 +++++++++++++++++-- 5 files changed, 114 insertions(+), 32 deletions(-) diff --git a/Makefile b/Makefile index ace3b52b7..aa3fd8e71 100644 --- a/Makefile +++ b/Makefile @@ -12,7 +12,7 @@ # See the License for the specific language governing permissions and # limitations under the License. -REDIS_VERSION:=3.2.0 +REDIS_VERSION:=3.2.6 SPRING_PROFILE?=ci ####### @@ -25,6 +25,8 @@ work/redis-%.conf: echo port $* >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ + echo notify-keyspace-events Ex >> $@ echo pidfile $(shell pwd)/work/redis-$*.pid >> $@ echo logfile $(shell pwd)/work/redis-$*.log >> $@ echo save \"\" >> $@ @@ -36,6 +38,8 @@ work/redis-6379.conf: echo port 6379 >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ + echo notify-keyspace-events Ex >> $@ echo pidfile $(shell pwd)/work/redis-6379.pid >> $@ echo logfile $(shell pwd)/work/redis-6379.log >> $@ echo save \"\" >> $@ @@ -57,6 +61,7 @@ work/sentinel-%.conf: echo port $* >> $@ echo daemonize yes >> $@ + echo protected-mode no >> $@ echo bind 0.0.0.0 >> $@ echo pidfile $(shell pwd)/work/sentinel-$*.pid >> $@ echo logfile $(shell pwd)/work/sentinel-$*.log >> $@ @@ -80,6 +85,7 @@ work/cluster-%.conf: @mkdir -p $(@D) echo port $* >> $@ + echo protected-mode no >> $@ echo cluster-enabled yes >> $@ echo cluster-config-file $(shell pwd)/work/nodes-$*.conf >> $@ echo cluster-node-timeout 5 >> $@ diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java b/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java index 054860c29..2795a9f7b 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyExpiredEvent.java @@ -48,12 +48,24 @@ public class RedisKeyExpiredEvent extends RedisKeyspaceEvent { /** * Creates new {@link RedisKeyExpiredEvent} - * + * * @param key * @param value */ public RedisKeyExpiredEvent(byte[] key, Object value) { - super(key); + this(null, key, value); + } + + /** + * Creates new {@link RedisKeyExpiredEvent} + * + * @pamam channel + * @param key + * @param value + * @since 1.8 + */ + public RedisKeyExpiredEvent(String channel, byte[] key, Object value) { + super(channel, key); args = ByteUtils.split(key, ':'); this.value = value; diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java b/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java index ee136b3cc..e36162405 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyValueAdapter.java @@ -64,6 +64,7 @@ import org.springframework.data.redis.util.ByteUtils; import org.springframework.data.util.CloseableIterator; import org.springframework.util.Assert; import org.springframework.util.ObjectUtils; +import org.springframework.util.StringUtils; /** * Redis specific {@link KeyValueAdapter} implementation. Uses binary codec to read/write data from/to Redis. Objects @@ -367,7 +368,6 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter List keys = new ArrayList(ids); - if (keys.isEmpty() || keys.size() < offset) { return Collections.emptyList(); } @@ -701,28 +701,7 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter */ @Override public void onApplicationEvent(RedisKeyspaceEvent event) { - - LOGGER.debug("Received %s .", event); - - if (event instanceof RedisKeyExpiredEvent) { - - final RedisKeyExpiredEvent expiredEvent = (RedisKeyExpiredEvent) event; - - redisOps.execute(new RedisCallback() { - - @Override - public Void doInRedis(RedisConnection connection) throws DataAccessException { - - LOGGER.debug("Cleaning up expired key '%s' data structures in keyspace '%s'.", expiredEvent.getSource(), - expiredEvent.getKeyspace()); - - connection.sRem(toBytes(expiredEvent.getKeyspace()), expiredEvent.getId()); - new IndexWriter(connection, converter).removeKeyFromIndexes(expiredEvent.getKeyspace(), expiredEvent.getId()); - return null; - } - }); - - } + // just a customization hook } /* @@ -814,12 +793,29 @@ public class RedisKeyValueAdapter extends AbstractKeyValueAdapter if (!org.springframework.util.CollectionUtils.isEmpty(hash)) { connection.del(phantomKey); } + return hash; } }); Object value = converter.read(Object.class, new RedisData(hash)); - publishEvent(new RedisKeyExpiredEvent(key, value)); + + String channel = !ObjectUtils.isEmpty(message.getChannel()) + ? converter.getConversionService().convert(message.getChannel(), String.class) : null; + + final RedisKeyExpiredEvent event = new RedisKeyExpiredEvent(channel, key, value); + + ops.execute(new RedisCallback() { + @Override + public Void doInRedis(RedisConnection connection) throws DataAccessException { + + connection.sRem(converter.getConversionService().convert(event.getKeyspace(), byte[].class), event.getId()); + new IndexWriter(connection, converter).removeKeyFromIndexes(event.getKeyspace(), event.getId()); + return null; + } + }); + + publishEvent(event); } private boolean isKeyExpirationMessage(Message message) { diff --git a/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java b/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java index 503825903..fee229eb0 100644 --- a/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java +++ b/src/main/java/org/springframework/data/redis/core/RedisKeyspaceEvent.java @@ -25,13 +25,28 @@ import org.springframework.context.ApplicationEvent; */ public class RedisKeyspaceEvent extends ApplicationEvent { + private final String channel; + /** * Creates new {@link RedisKeyspaceEvent}. - * + * * @param key The key that expired. Must not be {@literal null}. */ public RedisKeyspaceEvent(byte[] key) { + this(null, key); + } + + /** + * Creates new {@link RedisKeyspaceEvent}. + * + * @param channel The source channel aka subscription topic. Can be {@literal null}. + * @param key The key that expired. Must not be {@literal null}. + * @since 1.8 + */ + public RedisKeyspaceEvent(String channel, byte[] key) { + super(key); + this.channel = channel; } /* @@ -42,4 +57,13 @@ public class RedisKeyspaceEvent extends ApplicationEvent { return (byte[]) super.getSource(); } + /** + * + * @return can be {@literal null}. + * @since 1.8 + */ + public String getChannel() { + return this.channel; + } + } diff --git a/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java b/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java index 54a921e6b..ae065a466 100644 --- a/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java +++ b/src/test/java/org/springframework/data/redis/core/RedisKeyValueAdapterTests.java @@ -24,6 +24,7 @@ import java.util.Date; import java.util.LinkedHashMap; import java.util.List; import java.util.Map; +import java.util.concurrent.TimeUnit; import org.junit.After; import org.junit.AfterClass; @@ -43,6 +44,7 @@ import org.springframework.data.redis.connection.RedisConnection; import org.springframework.data.redis.connection.RedisConnectionFactory; import org.springframework.data.redis.connection.jedis.JedisConnectionFactory; import org.springframework.data.redis.connection.lettuce.LettuceConnectionFactory; +import org.springframework.data.redis.core.RedisKeyValueAdapter.EnableKeyspaceEvents; import org.springframework.data.redis.core.convert.Bucket; import org.springframework.data.redis.core.convert.KeyspaceConfiguration; import org.springframework.data.redis.core.convert.MappingConfiguration; @@ -92,6 +94,7 @@ public class RedisKeyValueAdapterTests { mappingContext.afterPropertiesSet(); adapter = new RedisKeyValueAdapter(template, mappingContext); + adapter.setEnableKeyspaceEvents(EnableKeyspaceEvents.ON_STARTUP); adapter.afterPropertiesSet(); template.execute(new RedisCallback() { @@ -239,26 +242,67 @@ public class RedisKeyValueAdapterTests { assertThat(template.opsForSet().members("persons:firstname:rand"), not(hasItem("1"))); } - @Test // DATAREDIS-425 - public void keyExpiredEventShouldRemoveHelperStructures() { + /** + * @see DATAREDIS-425 + */ + @Test + public void keyExpiredEventShouldRemoveHelperStructures() throws InterruptedException { Map map = new LinkedHashMap(); map.put("_class", Person.class.getName()); map.put("firstname", "rand"); map.put("address.country", "Andor"); + template.opsForHash().putAll("persons:1", map); + template.expire("persons:1", 1, TimeUnit.SECONDS); + template.opsForSet().add("persons", "1"); template.opsForSet().add("persons:firstname:rand", "1"); template.opsForSet().add("persons:1:idx", "persons:firstname:rand"); - adapter.onApplicationEvent(new RedisKeyExpiredEvent("persons:1".getBytes(Bucket.CHARSET))); + int iterationCount = 0; + while (template.hasKey("persons:1") && iterationCount++ < 3) { // ci might be a little slow + Thread.sleep(2000); + } + assertThat(template.hasKey("persons:1"), is(false)); assertThat(template.hasKey("persons:firstname:rand"), is(false)); assertThat(template.hasKey("persons:1:idx"), is(false)); assertThat(template.opsForSet().members("persons"), not(hasItem("1"))); } - @Test // DATAREDIS-512 + /** + * @see DATAREDIS-589 + */ + @Test + public void keyExpiredEventWithoutKeyspaceShouldBeIgnored() throws InterruptedException { + + Map map = new LinkedHashMap(); + map.put("_class", Person.class.getName()); + map.put("firstname", "rand"); + map.put("address.country", "Andor"); + + template.opsForHash().putAll("persons:1", map); + template.opsForHash().putAll("1", map); + + template.expire("1", 1, TimeUnit.SECONDS); + + template.opsForSet().add("persons", "1"); + template.opsForSet().add("persons:firstname:rand", "1"); + template.opsForSet().add("persons:1:idx", "persons:firstname:rand"); + + Thread.sleep(2000); + + assertThat(template.hasKey("persons:1"), is(true)); + assertThat(template.hasKey("persons:firstname:rand"), is(true)); + assertThat(template.hasKey("persons:1:idx"), is(true)); + assertThat(template.opsForSet().members("persons"), hasItem("1")); + } + + /** + * @see DATAREDIS-512 + */ + @Test public void putWritesIndexDataCorrectly() { Person rand = new Person();