From d5a77d139967fea26cb43f5c81af1f2865e85332 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 19 Dec 2024 11:55:55 -0500 Subject: [PATCH] GH-112: Fix `JdbcConsumerConfiguration` for auto-wire ambiguity Fixes: https://github.com/spring-cloud/spring-functions-catalog/issues/112 --- .../jdbc/JdbcConsumerConfiguration.java | 20 +++++++------------ 1 file changed, 7 insertions(+), 13 deletions(-) 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 dc7451d4..a0ac7256 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 @@ -51,7 +51,6 @@ 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; import org.springframework.integration.store.SimpleMessageStore; import org.springframework.integration.support.MutableMessage; import org.springframework.jdbc.core.namedparam.MapSqlParameterSource; @@ -120,8 +119,8 @@ public class JdbcConsumerConfiguration { } @Bean - IntegrationFlow jdbcConsumerFlow(@Qualifier("aggregator") MessageHandler aggregator, - JdbcMessageHandler jdbcMessageHandler) { + IntegrationFlow jdbcConsumerFlow(@Qualifier("jdbcConsumerAggregator") MessageHandler aggregator, + @Qualifier("jdbcConsumerMessageHandler") JdbcMessageHandler jdbcMessageHandler) { return (flow) -> { if (this.properties.getBatchSize() > 1 || this.properties.getIdleTimeout() > 0) { @@ -140,13 +139,16 @@ public class JdbcConsumerConfiguration { } @Bean - FactoryBean aggregator(MessageGroupStore messageGroupStore) { + FactoryBean jdbcConsumerAggregator() { AggregatorFactoryBean aggregatorFactoryBean = new AggregatorFactoryBean(); aggregatorFactoryBean.setCorrelationStrategy((message) -> message.getPayload().getClass().getName()); aggregatorFactoryBean.setReleaseStrategy(new MessageCountReleaseStrategy(this.properties.getBatchSize())); if (this.properties.getIdleTimeout() >= 0) { aggregatorFactoryBean.setGroupTimeoutExpression(new ValueExpression<>(this.properties.getIdleTimeout())); } + SimpleMessageStore messageGroupStore = new SimpleMessageStore(); + messageGroupStore.setTimeoutOnIdle(true); + messageGroupStore.setCopyOnGet(false); aggregatorFactoryBean.setMessageStore(messageGroupStore); aggregatorFactoryBean.setProcessorBean(new DefaultAggregatingMessageGroupProcessor()); aggregatorFactoryBean.setExpireGroupsUponCompletion(true); @@ -155,15 +157,7 @@ public class JdbcConsumerConfiguration { } @Bean - MessageGroupStore messageGroupStore() { - SimpleMessageStore messageGroupStore = new SimpleMessageStore(); - messageGroupStore.setTimeoutOnIdle(true); - messageGroupStore.setCopyOnGet(false); - return messageGroupStore; - } - - @Bean - public JdbcMessageHandler jdbcMessageHandler(DataSource dataSource, + public JdbcMessageHandler jdbcConsumerMessageHandler(DataSource dataSource, @Qualifier(IntegrationContextUtils.INTEGRATION_EVALUATION_CONTEXT_BEAN_NAME) EvaluationContext evaluationContext) { final MultiValueMap columnExpressionVariations = new LinkedMultiValueMap<>();