diff --git a/consumer/spring-cassandra-consumer/src/main/java/org/springframework/cloud/fn/consumer/cassandra/CassandraConsumerConfiguration.java b/consumer/spring-cassandra-consumer/src/main/java/org/springframework/cloud/fn/consumer/cassandra/CassandraConsumerConfiguration.java index ae7e8c3f..34d9a6e8 100644 --- a/consumer/spring-cassandra-consumer/src/main/java/org/springframework/cloud/fn/consumer/cassandra/CassandraConsumerConfiguration.java +++ b/consumer/spring-cassandra-consumer/src/main/java/org/springframework/cloud/fn/consumer/cassandra/CassandraConsumerConfiguration.java @@ -47,9 +47,9 @@ import org.springframework.data.cassandra.core.UpdateOptions; import org.springframework.data.cassandra.core.WriteResult; import org.springframework.data.cassandra.core.cql.WriteOptions; import org.springframework.integration.JavaUtils; -import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.cassandra.outbound.CassandraMessageHandler; import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.gateway.AnnotationGatewayProxyFactoryBean; import org.springframework.integration.support.json.Jackson2JsonObjectMapper; import org.springframework.integration.transformer.AbstractPayloadTransformer; import org.springframework.messaging.MessageHandler; @@ -71,7 +71,7 @@ public class CassandraConsumerConfiguration { private CassandraConsumerProperties cassandraSinkProperties; @Bean - public Consumer cassandraConsumer(CassandraConsumerFunction cassandraConsumerFunction) { + Consumer cassandraConsumer(CassandraConsumerFunction cassandraConsumerFunction) { return (payload) -> cassandraConsumerFunction.apply(payload).block(); } @@ -121,6 +121,13 @@ public class CassandraConsumerConfiguration { return cassandraMessageHandler; } + @Bean + AnnotationGatewayProxyFactoryBean cassandraConsumerFunction() { + var gatewayProxyFactoryBean = new AnnotationGatewayProxyFactoryBean<>(CassandraConsumerFunction.class); + gatewayProxyFactoryBean.setDefaultRequestChannelName("cassandraConsumerFlow.input"); + return gatewayProxyFactoryBean; + } + private static boolean isUuid(String uuid) { if (uuid.length() == 36) { String[] parts = uuid.split("-"); @@ -198,7 +205,6 @@ public class CassandraConsumerConfiguration { } - @MessagingGateway(name = "cassandraConsumerFunction", defaultRequestChannel = "cassandraConsumerFlow.input") interface CassandraConsumerFunction extends Function> { } diff --git a/consumer/spring-elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java b/consumer/spring-elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java index c1592b9f..37aa0b13 100644 --- a/consumer/spring-elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java +++ b/consumer/spring-elasticsearch-consumer/src/main/java/org/springframework/cloud/fn/consumer/elasticsearch/ElasticsearchConsumerConfiguration.java @@ -45,10 +45,10 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.context.annotation.Bean; import org.springframework.integration.aggregator.AbstractAggregatingMessageGroupProcessor; import org.springframework.integration.aggregator.MessageCountReleaseStrategy; -import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.config.AggregatorFactoryBean; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.expression.ValueExpression; +import org.springframework.integration.gateway.AnnotationGatewayProxyFactoryBean; import org.springframework.integration.store.MessageGroup; import org.springframework.integration.store.MessageGroupStore; import org.springframework.integration.store.SimpleMessageStore; @@ -138,6 +138,14 @@ public class ElasticsearchConsumerConfiguration { }; } + @Bean + @SuppressWarnings({ "unchecked", "rawtypes" }) + AnnotationGatewayProxyFactoryBean>> elasticsearchConsumer() { + var gatewayProxyFactoryBean = new AnnotationGatewayProxyFactoryBean<>(Consumer.class); + gatewayProxyFactoryBean.setDefaultRequestChannelName("elasticsearchConsumerFlow.input"); + return (AnnotationGatewayProxyFactoryBean) gatewayProxyFactoryBean; + } + @Bean public MessageHandler indexingHandler(ElasticsearchClient elasticsearchClient, ElasticsearchConsumerProperties consumerProperties) { @@ -303,9 +311,4 @@ public class ElasticsearchConsumerConfiguration { } - @MessagingGateway(name = "elasticsearchConsumer", defaultRequestChannel = "elasticsearchConsumerFlow.input") - private interface MessageConsumer extends Consumer> { - - } - } diff --git a/consumer/spring-jdbc-consumer/src/main/java/org/springframework/cloud/fn/consumer/jdbc/JdbcConsumerConfiguration.java b/consumer/spring-jdbc-consumer/src/main/java/org/springframework/cloud/fn/consumer/jdbc/JdbcConsumerConfiguration.java index 57fc3b0e..dc7451d4 100644 --- a/consumer/spring-jdbc-consumer/src/main/java/org/springframework/cloud/fn/consumer/jdbc/JdbcConsumerConfiguration.java +++ b/consumer/spring-jdbc-consumer/src/main/java/org/springframework/cloud/fn/consumer/jdbc/JdbcConsumerConfiguration.java @@ -44,11 +44,11 @@ import org.springframework.expression.spel.SpelParseException; import org.springframework.expression.spel.standard.SpelExpressionParser; import org.springframework.integration.aggregator.DefaultAggregatingMessageGroupProcessor; import org.springframework.integration.aggregator.MessageCountReleaseStrategy; -import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.config.AggregatorFactoryBean; import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.expression.ValueExpression; +import org.springframework.integration.gateway.AnnotationGatewayProxyFactoryBean; import org.springframework.integration.jdbc.JdbcMessageHandler; import org.springframework.integration.jdbc.SqlParameterSourceFactory; import org.springframework.integration.store.MessageGroupStore; @@ -131,6 +131,14 @@ public class JdbcConsumerConfiguration { }; } + @Bean + @SuppressWarnings({ "unchecked", "rawtypes" }) + AnnotationGatewayProxyFactoryBean>> jdbcConsumer() { + var gatewayProxyFactoryBean = new AnnotationGatewayProxyFactoryBean<>(Consumer.class); + gatewayProxyFactoryBean.setDefaultRequestChannelName("jdbcConsumerFlow.input"); + return (AnnotationGatewayProxyFactoryBean) gatewayProxyFactoryBean; + } + @Bean FactoryBean aggregator(MessageGroupStore messageGroupStore) { AggregatorFactoryBean aggregatorFactoryBean = new AggregatorFactoryBean(); @@ -232,11 +240,6 @@ public class JdbcConsumerConfiguration { return dataSourceInitializer; } - @MessagingGateway(name = "jdbcConsumer", defaultRequestChannel = "jdbcConsumerFlow.input") - public interface MessageConsumer extends Consumer> { - - } - private record ParameterFactory(MultiValueMap columnExpressions, EvaluationContext context) implements SqlParameterSourceFactory { diff --git a/consumer/spring-log-consumer/src/main/java/org/springframework/cloud/fn/consumer/log/LogConsumerConfiguration.java b/consumer/spring-log-consumer/src/main/java/org/springframework/cloud/fn/consumer/log/LogConsumerConfiguration.java index f98993b4..97cd61f7 100644 --- a/consumer/spring-log-consumer/src/main/java/org/springframework/cloud/fn/consumer/log/LogConsumerConfiguration.java +++ b/consumer/spring-log-consumer/src/main/java/org/springframework/cloud/fn/consumer/log/LogConsumerConfiguration.java @@ -21,8 +21,8 @@ import java.util.function.Consumer; import org.springframework.boot.autoconfigure.AutoConfiguration; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.context.annotation.Bean; -import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.integration.gateway.AnnotationGatewayProxyFactoryBean; import org.springframework.messaging.Message; /** @@ -46,9 +46,12 @@ public class LogConsumerConfiguration { .nullChannel(); } - @MessagingGateway(name = "logConsumer", defaultRequestChannel = "logConsumerFlow.input") - private interface MessageConsumer extends Consumer> { - + @Bean + @SuppressWarnings({ "unchecked", "rawtypes" }) + AnnotationGatewayProxyFactoryBean>> logConsumer() { + var gatewayProxyFactoryBean = new AnnotationGatewayProxyFactoryBean<>(Consumer.class); + gatewayProxyFactoryBean.setDefaultRequestChannelName("logConsumerFlow.input"); + return (AnnotationGatewayProxyFactoryBean) gatewayProxyFactoryBean; } } diff --git a/consumer/spring-sftp-consumer/src/main/java/org/springframework/cloud/fn/consumer/sftp/SftpConsumerConfiguration.java b/consumer/spring-sftp-consumer/src/main/java/org/springframework/cloud/fn/consumer/sftp/SftpConsumerConfiguration.java index 1ac1aee7..0f59a906 100644 --- a/consumer/spring-sftp-consumer/src/main/java/org/springframework/cloud/fn/consumer/sftp/SftpConsumerConfiguration.java +++ b/consumer/spring-sftp-consumer/src/main/java/org/springframework/cloud/fn/consumer/sftp/SftpConsumerConfiguration.java @@ -25,9 +25,9 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.fn.common.config.ComponentCustomizer; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Import; -import org.springframework.integration.annotation.MessagingGateway; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.file.remote.session.SessionFactory; +import org.springframework.integration.gateway.AnnotationGatewayProxyFactoryBean; import org.springframework.integration.sftp.dsl.Sftp; import org.springframework.integration.sftp.dsl.SftpMessageHandlerSpec; import org.springframework.integration.sftp.session.SftpRemoteFileTemplate; @@ -69,9 +69,12 @@ public class SftpConsumerConfiguration { return (flow) -> flow.handle(handlerSpec); } - @MessagingGateway(name = "sftpConsumer", defaultRequestChannel = "ftpOutboundFlow.input") - public interface MessageConsumer extends Consumer> { - + @Bean + @SuppressWarnings({ "unchecked", "rawtypes" }) + AnnotationGatewayProxyFactoryBean>> sftpConsumer() { + var gatewayProxyFactoryBean = new AnnotationGatewayProxyFactoryBean<>(Consumer.class); + gatewayProxyFactoryBean.setDefaultRequestChannelName("ftpOutboundFlow.input"); + return (AnnotationGatewayProxyFactoryBean) gatewayProxyFactoryBean; } }