diff --git a/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandler.java b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandler.java new file mode 100644 index 0000000000..a424dc57ee --- /dev/null +++ b/spring-integration-redis/src/main/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandler.java @@ -0,0 +1,153 @@ +/* + * Copyright 2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.redis.outbound; + +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory; +import org.springframework.data.redis.connection.stream.Record; +import org.springframework.data.redis.connection.stream.StreamRecords; +import org.springframework.data.redis.core.ReactiveRedisTemplate; +import org.springframework.data.redis.core.ReactiveStreamOperations; +import org.springframework.data.redis.hash.HashMapper; +import org.springframework.data.redis.serializer.RedisSerializationContext; +import org.springframework.expression.EvaluationContext; +import org.springframework.expression.Expression; +import org.springframework.expression.common.LiteralExpression; +import org.springframework.integration.expression.ExpressionUtils; +import org.springframework.integration.handler.AbstractReactiveMessageHandler; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.util.Assert; + +import reactor.core.publisher.Mono; + +/** + * Implementation of {@link org.springframework.messaging.ReactiveMessageHandler} which writes + * Message payload or Message itself (see {@link #extractPayload}) into a Redis stream using Reactive Stream operations. + * + * @author Attoumane Ahamadi + * @author Artem Bilan + * + * @since 5.4 + */ +public class ReactiveRedisStreamMessageHandler extends AbstractReactiveMessageHandler { + + private final Expression streamKeyExpression; + + private final ReactiveRedisConnectionFactory connectionFactory; + + private EvaluationContext evaluationContext; + + private boolean extractPayload = true; + + private ReactiveStreamOperations reactiveStreamOperations; + + private RedisSerializationContext serializationContext = RedisSerializationContext.string(); + + @Nullable + private HashMapper hashMapper; + + /** + * Create an instance based on provided {@link ReactiveRedisConnectionFactory} and key for stream. + * @param connectionFactory the {@link ReactiveRedisConnectionFactory} to use + * @param streamKey the key for stream + */ + public ReactiveRedisStreamMessageHandler(ReactiveRedisConnectionFactory connectionFactory, String streamKey) { + this(connectionFactory, new LiteralExpression(streamKey)); + } + + /** + * Create an instance based on provided {@link ReactiveRedisConnectionFactory} and expression for stream key. + * @param connectionFactory the {@link ReactiveRedisConnectionFactory} to use + * @param streamKeyExpression the SpEL expression to evaluate a key for stream + */ + public ReactiveRedisStreamMessageHandler(ReactiveRedisConnectionFactory connectionFactory, + Expression streamKeyExpression) { + + Assert.notNull(streamKeyExpression, "'streamKeyExpression' must not be null"); + Assert.notNull(connectionFactory, "'connectionFactory' must not be null"); + this.streamKeyExpression = streamKeyExpression; + this.connectionFactory = connectionFactory; + } + + public void setSerializationContext(RedisSerializationContext serializationContext) { + Assert.notNull(serializationContext, "'serializationContext' must not be null"); + this.serializationContext = serializationContext; + } + + /** + * (Optional) Set the {@link HashMapper} used to create {@link #reactiveStreamOperations}. + * The default {@link HashMapper} is defined from the provided {@link RedisSerializationContext} + * @param hashMapper the wanted hashMapper + * */ + public void setHashMapper(@Nullable HashMapper hashMapper) { + this.hashMapper = hashMapper; + } + + /** + * Set to {@code true} to extract the payload; otherwise + * the entire message is sent. Default {@code true}. + * @param extractPayload false to not extract. + */ + public void setExtractPayload(boolean extractPayload) { + this.extractPayload = extractPayload; + } + + @Override + public String getComponentType() { + return "redis:stream-outbound-channel-adapter"; + } + + @Override + @SuppressWarnings("unchecked") + protected void onInit() { + super.onInit(); + + this.evaluationContext = ExpressionUtils.createStandardEvaluationContext(getBeanFactory()); + + ReactiveRedisTemplate template = + new ReactiveRedisTemplate<>(this.connectionFactory, this.serializationContext); + this.reactiveStreamOperations = + this.hashMapper == null + ? template.opsForStream() + : template.opsForStream( + (HashMapper) this.hashMapper); + } + + @Override + protected Mono handleMessageInternal(Message message) { + return Mono + .fromSupplier(() -> { + String streamKey = this.streamKeyExpression.getValue(this.evaluationContext, message, String.class); + Assert.notNull(streamKey, "'streamKey' must not be null"); + return streamKey; + }) + .flatMap((streamKey) -> { + Object value = message; + if (this.extractPayload) { + value = message.getPayload(); + } + + Record record = + StreamRecords.objectBacked(value) + .withStreamKey(streamKey); + + return this.reactiveStreamOperations.add(record); + }) + .then(); + } + +} diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandlerTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandlerTests.java new file mode 100644 index 0000000000..91eb67e89f --- /dev/null +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/outbound/ReactiveRedisStreamMessageHandlerTests.java @@ -0,0 +1,170 @@ +/* + * Copyright 2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.redis.outbound; + +import static org.assertj.core.api.Assertions.assertThat; + +import org.junit.Before; +import org.junit.Test; +import org.junit.runner.RunWith; + +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory; +import org.springframework.data.redis.connection.stream.ObjectRecord; +import org.springframework.data.redis.connection.stream.StreamOffset; +import org.springframework.data.redis.core.ReactiveRedisTemplate; +import org.springframework.data.redis.serializer.RedisSerializationContext; +import org.springframework.data.redis.serializer.StringRedisSerializer; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.handler.ReactiveMessageHandlerAdapter; +import org.springframework.integration.redis.rules.RedisAvailable; +import org.springframework.integration.redis.rules.RedisAvailableRule; +import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.util.Address; +import org.springframework.integration.redis.util.Person; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.support.GenericMessage; +import org.springframework.test.annotation.DirtiesContext; +import org.springframework.test.context.junit4.SpringRunner; + +/** + * @author Attoumane Ahamadi + * @author Artem Bilan + * + * @since 5.4 + */ +@RunWith(SpringRunner.class) +@DirtiesContext +public class ReactiveRedisStreamMessageHandlerTests extends RedisAvailableTests { + + private static final String STREAM_KEY = "myStream"; + + @Autowired + @Qualifier("streamChannel") + private MessageChannel messageChannel; + + @Autowired + private ReactiveRedisConnectionFactory redisConnectionFactory; + + @Autowired + private ReactiveMessageHandlerAdapter handlerAdapter; + + @Autowired + private ReactiveRedisStreamMessageHandler streamMessageHandler; + + @Before + public void deleteStreamKey() { + ReactiveRedisTemplate template = new ReactiveRedisTemplate<>(this.redisConnectionFactory, + RedisSerializationContext.string()); + template.delete(STREAM_KEY).block(); + } + + + @Test + @RedisAvailable + public void integrationStreamOutboundTest() { + String messagePayload = "Hello stream message"; + + messageChannel.send(new GenericMessage<>(messagePayload)); + + RedisSerializationContext serializationContext = redisSerializationContext(); + + ReactiveRedisTemplate template = + new ReactiveRedisTemplate<>(redisConnectionFactory, serializationContext); + + ObjectRecord record = + template.opsForStream() + .read(String.class, StreamOffset.fromStart(STREAM_KEY)) + .blockFirst(); + + assertThat(record.getStream()).isEqualTo(STREAM_KEY); + + assertThat(record.getValue()).isEqualTo(messagePayload); + } + + @Test + @RedisAvailable + public void explicitSerializationContextWithModelTest() { + Address address = new Address("Rennes, France"); + Person person = new Person(address, "Attoumane"); + + Message message = new GenericMessage<>(person); + + RedisSerializationContext serializationContext = redisSerializationContext(); + + streamMessageHandler.setSerializationContext(serializationContext); + streamMessageHandler.afterPropertiesSet(); + + handlerAdapter.handleMessage(message); + + ReactiveRedisTemplate template = + new ReactiveRedisTemplate<>(redisConnectionFactory, serializationContext); + + ObjectRecord record = + template.opsForStream() + .read(Person.class, StreamOffset.fromStart(STREAM_KEY)) + .blockFirst(); + + assertThat(record.getStream()).isEqualTo(STREAM_KEY); + assertThat(record.getValue().getName()).isEqualTo("Attoumane"); + assertThat(record.getValue().getAddress().getAddress()).isEqualTo("Rennes, France"); + } + + + private RedisSerializationContext redisSerializationContext() { + return RedisSerializationContext.fromSerializer(StringRedisSerializer.UTF_8); + } + + + @Configuration + public static class ReactiveRedisStreamMessageHandlerTestsContext { + + @Bean + public MessageChannel streamChannel(ReactiveMessageHandlerAdapter messageHandlerAdapter) { + DirectChannel directChannel = new DirectChannel(); + directChannel.subscribe(messageHandlerAdapter); + directChannel.setMaxSubscribers(1); + return directChannel; + } + + + @Bean + public ReactiveRedisStreamMessageHandler streamMessageHandler( + ReactiveRedisConnectionFactory connectionFactory) { + + return new ReactiveRedisStreamMessageHandler(connectionFactory, STREAM_KEY); + } + + @Bean + public ReactiveMessageHandlerAdapter reactiveMessageHandlerAdapter( + ReactiveRedisStreamMessageHandler streamMessageHandler) { + + return new ReactiveMessageHandlerAdapter(streamMessageHandler); + } + + @Bean + public ReactiveRedisConnectionFactory reactiveRedisConnectionFactory() { + return RedisAvailableRule.connectionFactory; + } + + } + +} diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageStoreTests.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageStoreTests.java index c30298f117..193b04e1ad 100644 --- a/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageStoreTests.java +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/store/RedisMessageStoreTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2007-2019 the original author or authors. + * Copyright 2007-2020 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -18,7 +18,6 @@ package org.springframework.integration.redis.store; import static org.assertj.core.api.Assertions.assertThat; -import java.io.Serializable; import java.util.ArrayList; import java.util.List; import java.util.Properties; @@ -35,6 +34,8 @@ import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.history.MessageHistory; import org.springframework.integration.redis.rules.RedisAvailable; import org.springframework.integration.redis.rules.RedisAvailableTests; +import org.springframework.integration.redis.util.Address; +import org.springframework.integration.redis.util.Person; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; @@ -43,6 +44,7 @@ import org.springframework.messaging.support.GenericMessage; /** * @author Oleg Zhurakousky * @author Gary Russell + * @author Artem Bilan * */ public class RedisMessageStoreTests extends RedisAvailableTests { @@ -57,7 +59,7 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testGetNonExistingMessage() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); Message message = store.getMessage(UUID.randomUUID()); assertThat(message).isNull(); @@ -66,7 +68,7 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testGetMessageCountWhenEmpty() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); assertThat(store.getMessageCount()).isEqualTo(0); } @@ -76,8 +78,8 @@ public class RedisMessageStoreTests extends RedisAvailableTests { public void testAddStringMessage() { RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - Message stringMessage = new GenericMessage("Hello Redis"); - Message storedMessage = store.addMessage(stringMessage); + Message stringMessage = new GenericMessage<>("Hello Redis"); + Message storedMessage = store.addMessage(stringMessage); assertThat(storedMessage).isNotSameAs(stringMessage); assertThat(storedMessage.getPayload()).isEqualTo("Hello Redis"); } @@ -85,14 +87,14 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testAddSerializableObjectMessage() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); Address address = new Address(); address.setAddress("1600 Pennsylvania Av, Washington, DC"); Person person = new Person(address, "Barak Obama"); - Message objectMessage = new GenericMessage(person); - Message storedMessage = store.addMessage(objectMessage); + Message objectMessage = new GenericMessage<>(person); + Message storedMessage = store.addMessage(objectMessage); assertThat(storedMessage).isNotSameAs(objectMessage); assertThat(storedMessage.getPayload().getName()).isEqualTo("Barak Obama"); } @@ -100,10 +102,10 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test(expected = IllegalArgumentException.class) @RedisAvailable public void testAddNonSerializableObjectMessage() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - Message objectMessage = new GenericMessage(new Foo()); + Message objectMessage = new GenericMessage<>(new Foo()); store.addMessage(objectMessage); } @@ -111,9 +113,9 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testAddAndGetStringMessage() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - Message stringMessage = new GenericMessage("Hello Redis"); + Message stringMessage = new GenericMessage<>("Hello Redis"); store.addMessage(stringMessage); Message retrievedMessage = (Message) store.getMessage(stringMessage.getHeaders().getId()); assertThat(retrievedMessage).isNotNull(); @@ -124,9 +126,9 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testAddAndGetWithPrefix() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf, "foo"); - Message stringMessage = new GenericMessage("Hello Redis"); + Message stringMessage = new GenericMessage<>("Hello Redis"); store.addMessage(stringMessage); Message retrievedMessage = (Message) store.getMessage(stringMessage.getHeaders().getId()); assertThat(retrievedMessage).isNotNull(); @@ -142,9 +144,9 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable public void testAddAndRemoveStringMessage() { - RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); + RedisConnectionFactory jcf = getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - Message stringMessage = new GenericMessage("Hello Redis"); + Message stringMessage = new GenericMessage<>("Hello Redis"); store.addMessage(stringMessage); Message retrievedMessage = (Message) store.removeMessage(stringMessage.getHeaders().getId()); assertThat(retrievedMessage).isNotNull(); @@ -154,11 +156,11 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable - public void testWithMessageHistory() throws Exception { + public void testWithMessageHistory() { RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); RedisMessageStore store = new RedisMessageStore(jcf); - Message message = new GenericMessage("Hello"); + Message message = new GenericMessage<>("Hello"); DirectChannel fooChannel = new DirectChannel(); fooChannel.setBeanName("fooChannel"); DirectChannel barChannel = new DirectChannel(); @@ -178,11 +180,11 @@ public class RedisMessageStoreTests extends RedisAvailableTests { @Test @RedisAvailable - public void testAddAndRemoveMessagesFromMessageGroup() throws Exception { + public void testAddAndRemoveMessagesFromMessageGroup() { RedisConnectionFactory jcf = this.getConnectionFactoryForTest(); RedisMessageStore messageStore = new RedisMessageStore(jcf); String groupId = "X"; - List> messages = new ArrayList>(); + List> messages = new ArrayList<>(); for (int i = 0; i < 25; i++) { Message message = MessageBuilder.withPayload("foo").setCorrelationId(groupId).build(); messageStore.addMessagesToGroup(groupId, message); @@ -194,51 +196,6 @@ public class RedisMessageStoreTests extends RedisAvailableTests { messageStore.removeMessageGroup("X"); } - @SuppressWarnings("serial") - public static class Person implements Serializable { - - private Address address; - - private String name; - - public Person(Address address, String name) { - this.address = address; - this.name = name; - } - - public Address getAddress() { - return address; - } - - public void setAddress(Address address) { - this.address = address; - } - - public String getName() { - return name; - } - - public void setName(String name) { - this.name = name; - } - - } - - @SuppressWarnings("serial") - public static class Address implements Serializable { - - private String address; - - public String getAddress() { - return address; - } - - public void setAddress(String address) { - this.address = address; - } - - } - public static class Foo { } diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Address.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Address.java new file mode 100644 index 0000000000..6551477139 --- /dev/null +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Address.java @@ -0,0 +1,40 @@ +/* + * Copyright 2013-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.redis.util; + +import java.io.Serializable; + +@SuppressWarnings("serial") +public class Address implements Serializable { + + private String address; + + public String getAddress() { + return address; + } + + public void setAddress(String address) { + this.address = address; + } + + public Address() { + } + + public Address(String address) { + this.address = address; + } +} diff --git a/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Person.java b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Person.java new file mode 100644 index 0000000000..982aac5c4e --- /dev/null +++ b/spring-integration-redis/src/test/java/org/springframework/integration/redis/util/Person.java @@ -0,0 +1,52 @@ +/* + * Copyright 2013-2020 the original author or authors. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.springframework.integration.redis.util; + +import java.io.Serializable; + + +@SuppressWarnings("serial") +public class Person implements Serializable { + + private Address address; + + private String name; + + public Person(Address address, String name) { + this.address = address; + this.name = name; + } + + public Person() { + } + + public Address getAddress() { + return address; + } + + public void setAddress(Address address) { + this.address = address; + } + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } +} diff --git a/src/reference/asciidoc/redis.adoc b/src/reference/asciidoc/redis.adoc index e7861acc71..023922fd31 100644 --- a/src/reference/asciidoc/redis.adoc +++ b/src/reference/asciidoc/redis.adoc @@ -788,3 +788,8 @@ However, the resources protected by such a lock may have been compromised, so su You should set the expiry at a large enough value to prevent this condition, but set it low enough that the lock can be recovered after a server failure in a reasonable amount of time. Starting with version 5.0, the `RedisLockRegistry` implements `ExpirableLockRegistry`, which removes locks last acquired more than `age` ago and that are not currently locked. + +[[redis-stream-outbound]] +=== Redis Stream Outbound Channel Adapter + +TBD diff --git a/src/reference/asciidoc/whats-new.adoc b/src/reference/asciidoc/whats-new.adoc index ae3de2f627..be4ec644be 100644 --- a/src/reference/asciidoc/whats-new.adoc +++ b/src/reference/asciidoc/whats-new.adoc @@ -23,7 +23,12 @@ See <<./kafka.adoc#kafka,Spring for Apache Kafka Support>> for more information. ==== R2DBC Channel Adapters The Channel Adapters for R2DBC database interaction have been introduced. -See <<./r2dbc.adoc#r2dbc,R2DBC Support>> for more information. +See <<./r2dbc.adoc#r2dbc,R2DBC Support>> for more information. + +==== Redis Stream Support + +The Channel Adapters for Redis Stream support have been introduced. +See <<./redis.adoc#redis-stream-outbound,Redis Stream Outbound Channel Adapter>> for more information. [[x5.4-general]] === General Changes