From 9a1fb207d548d4705ee34c863c6ba97fbb469998 Mon Sep 17 00:00:00 2001 From: Spencer Gibb Date: Wed, 1 May 2019 17:39:58 -0400 Subject: [PATCH] Updates to use ReactiveStringRedisTemplate --- .../config/GatewayRedisAutoConfiguration.java | 20 ++-------------- .../filter/ratelimit/RedisRateLimiter.java | 11 +++++---- .../GatewayRSocketAutoConfiguration.java | 3 ++- .../core/GatewayRSocketIntegrationTests.java | 4 ++-- .../gateway/rsocket/test/PingPongApp.java | 24 +++++++++---------- .../sample/GatewaySampleApplicationTests.java | 4 ++-- 6 files changed, 25 insertions(+), 41 deletions(-) diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayRedisAutoConfiguration.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayRedisAutoConfiguration.java index 70b30e40..ab7630b9 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayRedisAutoConfiguration.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/config/GatewayRedisAutoConfiguration.java @@ -29,14 +29,11 @@ import org.springframework.cloud.gateway.filter.ratelimit.RedisRateLimiter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.io.ClassPathResource; -import org.springframework.data.redis.connection.ReactiveRedisConnectionFactory; import org.springframework.data.redis.core.ReactiveRedisTemplate; +import org.springframework.data.redis.core.ReactiveStringRedisTemplate; import org.springframework.data.redis.core.RedisTemplate; import org.springframework.data.redis.core.script.DefaultRedisScript; import org.springframework.data.redis.core.script.RedisScript; -import org.springframework.data.redis.serializer.RedisSerializationContext; -import org.springframework.data.redis.serializer.RedisSerializer; -import org.springframework.data.redis.serializer.StringRedisSerializer; import org.springframework.scripting.support.ResourceScriptSource; import org.springframework.validation.Validator; import org.springframework.web.reactive.DispatcherHandler; @@ -58,22 +55,9 @@ class GatewayRedisAutoConfiguration { return redisScript; } - @Bean - // TODO: replace with ReactiveStringRedisTemplate in future - public ReactiveRedisTemplate stringReactiveRedisTemplate( - ReactiveRedisConnectionFactory reactiveRedisConnectionFactory) { - RedisSerializer serializer = new StringRedisSerializer(); - RedisSerializationContext serializationContext = RedisSerializationContext - .newSerializationContext().key(serializer) - .value(serializer).hashKey(serializer).hashValue(serializer).build(); - return new ReactiveRedisTemplate<>(reactiveRedisConnectionFactory, - serializationContext); - } - @Bean @ConditionalOnMissingBean - public RedisRateLimiter redisRateLimiter( - ReactiveRedisTemplate redisTemplate, + public RedisRateLimiter redisRateLimiter(ReactiveStringRedisTemplate redisTemplate, @Qualifier(RedisRateLimiter.REDIS_SCRIPT_NAME) RedisScript> redisScript, Validator validator) { return new RedisRateLimiter(redisTemplate, redisScript, validator); diff --git a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/ratelimit/RedisRateLimiter.java b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/ratelimit/RedisRateLimiter.java index 4e5b7aee..1b9c0635 100644 --- a/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/ratelimit/RedisRateLimiter.java +++ b/spring-cloud-gateway-core/src/main/java/org/springframework/cloud/gateway/filter/ratelimit/RedisRateLimiter.java @@ -37,7 +37,7 @@ import org.springframework.boot.context.properties.ConfigurationProperties; import org.springframework.cloud.gateway.route.RouteDefinitionRouteLocator; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; -import org.springframework.data.redis.core.ReactiveRedisTemplate; +import org.springframework.data.redis.core.ReactiveStringRedisTemplate; import org.springframework.data.redis.core.script.RedisScript; import org.springframework.validation.Validator; import org.springframework.validation.annotation.Validated; @@ -92,7 +92,7 @@ public class RedisRateLimiter extends AbstractRateLimiter redisTemplate; + private ReactiveStringRedisTemplate redisTemplate; private RedisScript> script; @@ -119,7 +119,7 @@ public class RedisRateLimiter extends AbstractRateLimiter redisTemplate, + public RedisRateLimiter(ReactiveStringRedisTemplate redisTemplate, RedisScript> script, Validator validator) { super(Config.class, CONFIGURATION_PROPERTY_NAME, validator); this.redisTemplate = redisTemplate; @@ -182,8 +182,9 @@ public class RedisRateLimiter extends AbstractRateLimiter 0) { this.setValidator(context.getBean(Validator.class)); diff --git a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java index e17c7740..54a57c63 100644 --- a/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java +++ b/spring-cloud-gateway-rsocket/src/main/java/org/springframework/cloud/gateway/rsocket/autoconfigure/GatewayRSocketAutoConfiguration.java @@ -42,7 +42,8 @@ import org.springframework.core.env.Environment; * @author Spencer Gibb */ @Configuration -@ConditionalOnProperty(name = "spring.cloud.gateway.rsocket.enabled", matchIfMissing = true) +@ConditionalOnProperty(name = "spring.cloud.gateway.rsocket.enabled", + matchIfMissing = true) @EnableConfigurationProperties @ConditionalOnClass(RSocket.class) public class GatewayRSocketAutoConfiguration { diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java index 526a2039..02978c3c 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/core/GatewayRSocketIntegrationTests.java @@ -37,8 +37,8 @@ import org.springframework.util.SocketUtils; import static org.assertj.core.api.Assertions.assertThat; @RunWith(SpringRunner.class) -@SpringBootTest(classes = PingPongApp.class, properties = { - "ping.take=5" }, webEnvironment = WebEnvironment.RANDOM_PORT) +@SpringBootTest(classes = PingPongApp.class, properties = { "ping.take=5" }, + webEnvironment = WebEnvironment.RANDOM_PORT) public class GatewayRSocketIntegrationTests { private static int port; diff --git a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java index 5c7537f7..9df8e4ba 100644 --- a/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java +++ b/spring-cloud-gateway-rsocket/src/test/java/org/springframework/cloud/gateway/rsocket/test/PingPongApp.java @@ -146,21 +146,19 @@ public class PingPongApp { } Publisher doPing(Integer take, RSocket socket) { - Flux pong = socket.requestChannel( - Flux.interval(Duration.ofSeconds(1)).map(i -> { - ByteBuf data = ByteBufUtil.writeUtf8( - ByteBufAllocator.DEFAULT, "ping" + id); - ByteBuf routingMetadata = Metadata.from("pong") - .encode(); + Flux pong = socket + .requestChannel(Flux.interval(Duration.ofSeconds(1)).map(i -> { + ByteBuf data = ByteBufUtil.writeUtf8(ByteBufAllocator.DEFAULT, + "ping" + id); + ByteBuf routingMetadata = Metadata.from("pong").encode(); return DefaultPayload.create(data, routingMetadata); // onBackpressue is needed in case pong is not available yet - }).onBackpressureDrop(payload -> log.debug( - "Dropped payload " + payload.getDataUtf8())) - ).map(Payload::getDataUtf8).doOnNext(str -> { - int received = pongsReceived.incrementAndGet(); - log.info("received " + str + "(" + received + ") in Ping" - + id); - }).doFinally(signal -> socket.dispose()); + }).onBackpressureDrop(payload -> log + .debug("Dropped payload " + payload.getDataUtf8()))) + .map(Payload::getDataUtf8).doOnNext(str -> { + int received = pongsReceived.incrementAndGet(); + log.info("received " + str + "(" + received + ") in Ping" + id); + }).doFinally(signal -> socket.dispose()); if (take != null) { return pong.take(take); } diff --git a/spring-cloud-gateway-sample/src/test/java/org/springframework/cloud/gateway/sample/GatewaySampleApplicationTests.java b/spring-cloud-gateway-sample/src/test/java/org/springframework/cloud/gateway/sample/GatewaySampleApplicationTests.java index 50df3069..50d89c36 100644 --- a/spring-cloud-gateway-sample/src/test/java/org/springframework/cloud/gateway/sample/GatewaySampleApplicationTests.java +++ b/spring-cloud-gateway-sample/src/test/java/org/springframework/cloud/gateway/sample/GatewaySampleApplicationTests.java @@ -51,8 +51,8 @@ import static org.springframework.boot.test.context.SpringBootTest.WebEnvironmen * @author Spencer Gibb */ @RunWith(SpringRunner.class) -@SpringBootTest(classes = { - GatewaySampleApplicationTests.TestConfig.class }, webEnvironment = RANDOM_PORT, properties = "management.server.port=${test.port}") +@SpringBootTest(classes = { GatewaySampleApplicationTests.TestConfig.class }, + webEnvironment = RANDOM_PORT, properties = "management.server.port=${test.port}") public class GatewaySampleApplicationTests { protected static int managementPort;