diff --git a/function-dependencies/pom.xml b/function-dependencies/pom.xml
index cceac754..232cfe6c 100644
--- a/function-dependencies/pom.xml
+++ b/function-dependencies/pom.xml
@@ -185,6 +185,11 @@
websocket-consumer
${project.version}
+
+ org.springframework.cloud.fn
+ aggregator-function
+ ${project.version}
+
org.springframework.cloud.fn
filter-function
diff --git a/function/aggregator-function/README.adoc b/function/aggregator-function/README.adoc
new file mode 100644
index 00000000..2af08ebf
--- /dev/null
+++ b/function/aggregator-function/README.adoc
@@ -0,0 +1,25 @@
+# Aggregator Function
+
+This module provides an aggregation function that can be reused and composed in other applications.
+
+## Beans for injection
+
+You can import the `AggregatorFunctionConfiguration` in a Spring Boot application and then inject the following bean.
+
+`aggregatorFunction`
+
+You can use `aggregatorFunction` as a qualifier when injecting.
+
+Once injected, you can use the `apply` method of the `Function` to invoke it and get the result.
+
+## Configuration Options
+
+For more information on the various options available, please see link:src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionProperties.java[AggregatorFunctionProperties.java]
+
+## Tests
+
+See this link:src/test/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionApplicationTests.java[test suite] for examples of how this function is used.
+
+## Other usage
+
+See this link:../../../applications/processor/aggregator-processor/README.adoc[README] where this function is used to create a Spring Cloud Stream application.
diff --git a/function/aggregator-function/pom.xml b/function/aggregator-function/pom.xml
new file mode 100644
index 00000000..2577d59f
--- /dev/null
+++ b/function/aggregator-function/pom.xml
@@ -0,0 +1,127 @@
+
+
+ 4.0.0
+ aggregator-function
+ 1.0.0-SNAPSHOT
+ aggregator-function
+ Spring Native Function for Aggregator
+
+
+ org.springframework.cloud.fn
+ spring-functions-parent
+ 1.0.0-SNAPSHOT
+ ../../spring-functions-parent
+
+
+
+
+ org.springframework.cloud.fn
+ config-common
+ ${project.version}
+
+
+ org.springframework.boot
+ spring-boot-starter-integration
+
+
+
+ org.springframework.boot
+ spring-boot-configuration-processor
+ provided
+
+
+
+
+ org.springframework.integration
+ spring-integration-mongodb
+
+
+ org.springframework.boot
+ spring-boot-starter-data-mongodb
+ runtime
+
+
+ de.flapdoodle.embed
+ de.flapdoodle.embed.mongo
+ test
+
+
+
+
+ org.springframework.integration
+ spring-integration-redis
+
+
+ org.springframework.boot
+ spring-boot-starter-data-redis
+ runtime
+
+
+
+
+ org.springframework.integration
+ spring-integration-gemfire
+
+
+ org.springframework.geode
+ spring-geode-starter
+ 1.3.2.RELEASE
+
+
+
+
+ org.springframework.integration
+ spring-integration-jdbc
+
+
+ org.springframework.boot
+ spring-boot-starter-jdbc
+ runtime
+
+
+ org.hsqldb
+ hsqldb
+ runtime
+
+
+ com.h2database
+ h2
+ runtime
+
+
+ org.mariadb.jdbc
+ mariadb-java-client
+ runtime
+
+
+ org.postgresql
+ postgresql
+ runtime
+
+
+
+
+ org.springframework.boot
+ spring-boot-starter-test
+ test
+
+
+ org.junit.vintage
+ junit-vintage-engine
+
+
+
+
+ io.projectreactor
+ reactor-test
+ test
+
+
+ org.springframework.integration
+ spring-integration-test
+ test
+
+
+
+
diff --git a/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java
new file mode 100644
index 00000000..166693de
--- /dev/null
+++ b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionConfiguration.java
@@ -0,0 +1,140 @@
+/*
+ * Copyright 2020-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.cloud.fn.aggregator;
+
+import java.util.function.Function;
+
+import reactor.core.publisher.Flux;
+
+import org.springframework.beans.factory.BeanFactory;
+import org.springframework.beans.factory.BeanFactoryAware;
+import org.springframework.beans.factory.ObjectProvider;
+import org.springframework.beans.factory.annotation.Autowired;
+import org.springframework.beans.factory.annotation.Qualifier;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.boot.context.properties.EnableConfigurationProperties;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Configuration;
+import org.springframework.context.annotation.Import;
+import org.springframework.integration.aggregator.CorrelationStrategy;
+import org.springframework.integration.aggregator.DefaultAggregatingMessageGroupProcessor;
+import org.springframework.integration.aggregator.ExpressionEvaluatingCorrelationStrategy;
+import org.springframework.integration.aggregator.ExpressionEvaluatingMessageGroupProcessor;
+import org.springframework.integration.aggregator.ExpressionEvaluatingReleaseStrategy;
+import org.springframework.integration.aggregator.MessageGroupProcessor;
+import org.springframework.integration.aggregator.ReleaseStrategy;
+import org.springframework.integration.annotation.ServiceActivator;
+import org.springframework.integration.channel.FluxMessageChannel;
+import org.springframework.integration.config.AggregatorFactoryBean;
+import org.springframework.integration.store.MessageGroupStore;
+import org.springframework.messaging.Message;
+import org.springframework.messaging.MessageChannel;
+
+/**
+ * @author Artem Bilan
+ */
+@Configuration(proxyBeanMethods = false)
+@EnableConfigurationProperties(AggregatorFunctionProperties.class)
+public class AggregatorFunctionConfiguration {
+
+ @Autowired
+ private AggregatorFunctionProperties properties;
+
+ @Autowired
+ private BeanFactory beanFactory;
+
+ @Bean
+ public Function>, Flux>> aggregatorFunction(FluxMessageChannel inputChannel,
+ FluxMessageChannel outputChannel) {
+
+ return input -> Flux.from(outputChannel)
+ .doOnSubscribe((sub) -> inputChannel.subscribeTo(input));
+ }
+
+ @Bean
+ public FluxMessageChannel inputChannel() {
+ return new FluxMessageChannel();
+ }
+
+ @Bean
+ public FluxMessageChannel outputChannel() {
+ return new FluxMessageChannel();
+ }
+
+ @Bean
+ @ServiceActivator(inputChannel = "inputChannel")
+ public AggregatorFactoryBean aggregator(
+ ObjectProvider correlationStrategy,
+ ObjectProvider releaseStrategy,
+ ObjectProvider messageGroupProcessor,
+ ObjectProvider messageStore,
+ @Qualifier("outputChannel") MessageChannel outputChannel) {
+
+ AggregatorFactoryBean aggregator = new AggregatorFactoryBean();
+ aggregator.setExpireGroupsUponCompletion(true);
+ aggregator.setSendPartialResultOnExpiry(true);
+ aggregator.setGroupTimeoutExpression(this.properties.getGroupTimeout());
+
+ aggregator.setCorrelationStrategy(correlationStrategy.getIfAvailable());
+ aggregator.setReleaseStrategy(releaseStrategy.getIfAvailable());
+
+ MessageGroupProcessor groupProcessor = messageGroupProcessor.getIfAvailable();
+
+ if (groupProcessor == null) {
+ groupProcessor = new DefaultAggregatingMessageGroupProcessor();
+ ((BeanFactoryAware) groupProcessor).setBeanFactory(this.beanFactory);
+ }
+ aggregator.setProcessorBean(groupProcessor);
+
+ aggregator.setMessageStore(messageStore.getIfAvailable());
+ aggregator.setOutputChannel(outputChannel);
+
+ return aggregator;
+ }
+
+ @Bean
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX, name = "correlation")
+ @ConditionalOnMissingBean
+ public CorrelationStrategy correlationStrategy() {
+ return new ExpressionEvaluatingCorrelationStrategy(this.properties.getCorrelation());
+ }
+
+ @Bean
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX, name = "release")
+ @ConditionalOnMissingBean
+ public ReleaseStrategy releaseStrategy() {
+ return new ExpressionEvaluatingReleaseStrategy(this.properties.getRelease());
+ }
+
+ @Bean
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX, name = "aggregation")
+ @ConditionalOnMissingBean
+ public MessageGroupProcessor messageGroupProcessor() {
+ return new ExpressionEvaluatingMessageGroupProcessor(this.properties.getAggregation().getExpressionString());
+ }
+
+
+ @Configuration
+ @ConditionalOnMissingBean(MessageGroupStore.class)
+ @Import({ MessageStoreConfiguration.Mongo.class, MessageStoreConfiguration.Redis.class,
+ MessageStoreConfiguration.Gemfire.class, MessageStoreConfiguration.Jdbc.class })
+ protected static class MessageStoreAutoConfiguration {
+
+ }
+
+}
diff --git a/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionProperties.java b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionProperties.java
new file mode 100644
index 00000000..05b4003c
--- /dev/null
+++ b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/AggregatorFunctionProperties.java
@@ -0,0 +1,124 @@
+/*
+ * Copyright 2020-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.cloud.fn.aggregator;
+
+import org.springframework.boot.context.properties.ConfigurationProperties;
+import org.springframework.expression.Expression;
+
+/**
+ * Configuration properties for the Aggregator function.
+ *
+ * @author Artem Bilan
+ */
+@ConfigurationProperties("aggregator")
+public class AggregatorFunctionProperties {
+
+ static final String PREFIX = "aggregator";
+
+ /**
+ * SpEL expression for correlation key. Default to correlationId header.
+ */
+ private Expression correlation;
+
+ /**
+ * SpEL expression for release strategy. Default is based on the sequenceSize header.
+ */
+ private Expression release;
+
+ /**
+ * SpEL expression for aggregation strategy. Default is collection of payloads.
+ */
+ private Expression aggregation;
+
+ /**
+ * SpEL expression for timeout to expiring uncompleted groups.
+ */
+ private Expression groupTimeout;
+
+ /**
+ * Message store type.
+ */
+ private String messageStoreType = MessageStoreType.SIMPLE;
+
+ /**
+ * Persistence message store entity: table prefix in RDBMS, collection name in MongoDb, etc.
+ */
+ private String messageStoreEntity;
+
+ public Expression getCorrelation() {
+ return this.correlation;
+ }
+
+ public void setCorrelation(Expression correlation) {
+ this.correlation = correlation;
+ }
+
+ public Expression getRelease() {
+ return this.release;
+ }
+
+ public void setRelease(Expression release) {
+ this.release = release;
+ }
+
+ public Expression getAggregation() {
+ return this.aggregation;
+ }
+
+ public void setAggregation(Expression aggregation) {
+ this.aggregation = aggregation;
+ }
+
+ public Expression getGroupTimeout() {
+ return this.groupTimeout;
+ }
+
+ public void setGroupTimeout(Expression groupTimeout) {
+ this.groupTimeout = groupTimeout;
+ }
+
+ public String getMessageStoreEntity() {
+ return this.messageStoreEntity;
+ }
+
+ public void setMessageStoreEntity(String messageStoreEntity) {
+ this.messageStoreEntity = messageStoreEntity;
+ }
+
+ public String getMessageStoreType() {
+ return this.messageStoreType;
+ }
+
+ public void setMessageStoreType(String messageStoreType) {
+ this.messageStoreType = messageStoreType;
+ }
+
+ static final class MessageStoreType {
+
+ static final String SIMPLE = "simple";
+
+ static final String JDBC = "jdbc";
+
+ static final String MONGODB = "mongodb";
+
+ static final String REDIS = "redis";
+
+ static final String GEMFIRE = "gemfire";
+
+ }
+
+}
diff --git a/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/ExcludeStoresAutoConfigurationEnvironmentPostProcessor.java b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/ExcludeStoresAutoConfigurationEnvironmentPostProcessor.java
new file mode 100644
index 00000000..9e2fd4dc
--- /dev/null
+++ b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/ExcludeStoresAutoConfigurationEnvironmentPostProcessor.java
@@ -0,0 +1,64 @@
+/*
+ * Copyright 2020-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.cloud.fn.aggregator;
+
+import java.util.Properties;
+
+import org.springframework.boot.SpringApplication;
+import org.springframework.boot.autoconfigure.data.mongo.MongoDataAutoConfiguration;
+import org.springframework.boot.autoconfigure.data.mongo.MongoRepositoriesAutoConfiguration;
+import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
+import org.springframework.boot.autoconfigure.data.redis.RedisRepositoriesAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceTransactionManagerAutoConfiguration;
+import org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration;
+import org.springframework.boot.autoconfigure.mongo.embedded.EmbeddedMongoAutoConfiguration;
+import org.springframework.boot.env.EnvironmentPostProcessor;
+import org.springframework.core.env.ConfigurableEnvironment;
+import org.springframework.core.env.MutablePropertySources;
+import org.springframework.core.env.PropertiesPropertySource;
+import org.springframework.geode.boot.autoconfigure.ClientCacheAutoConfiguration;
+
+/**
+ * An {@link EnvironmentPostProcessor} to add {@code spring.autoconfigure.exclude} property
+ * since we can't use {@code application.properties} from the library perspective.
+ *
+ * @author Artem Bilan
+ */
+public class ExcludeStoresAutoConfigurationEnvironmentPostProcessor implements EnvironmentPostProcessor {
+
+ @Override
+ public void postProcessEnvironment(ConfigurableEnvironment environment, SpringApplication application) {
+ MutablePropertySources propertySources = environment.getPropertySources();
+ Properties properties = new Properties();
+
+ properties.setProperty("spring.autoconfigure.exclude",
+ DataSourceAutoConfiguration.class.getName() + ", " +
+ DataSourceTransactionManagerAutoConfiguration.class.getName() + ", " +
+ MongoAutoConfiguration.class.getName() + ", " +
+ MongoDataAutoConfiguration.class.getName() + ", " +
+ MongoRepositoriesAutoConfiguration.class.getName() + ", " +
+ EmbeddedMongoAutoConfiguration.class.getName() + ", " +
+ ClientCacheAutoConfiguration.class.getName() + ", " +
+ RedisAutoConfiguration.class.getName() + ", " +
+ RedisRepositoriesAutoConfiguration.class.getName());
+
+ propertySources.addLast(
+ new PropertiesPropertySource("aggregator.exclude.stores.auto-configuration", properties));
+ }
+
+}
diff --git a/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/MessageStoreConfiguration.java b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/MessageStoreConfiguration.java
new file mode 100644
index 00000000..9ad7a43c
--- /dev/null
+++ b/function/aggregator-function/src/main/java/org/springframework/cloud/fn/aggregator/MessageStoreConfiguration.java
@@ -0,0 +1,148 @@
+/*
+ * Copyright 2020-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.cloud.fn.aggregator;
+
+import java.util.Arrays;
+
+import org.apache.geode.cache.GemFireCache;
+import org.apache.geode.cache.Region;
+
+import org.springframework.boot.autoconfigure.condition.ConditionalOnClass;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
+import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
+import org.springframework.boot.autoconfigure.data.mongo.MongoDataAutoConfiguration;
+import org.springframework.boot.autoconfigure.data.redis.RedisAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceAutoConfiguration;
+import org.springframework.boot.autoconfigure.jdbc.DataSourceTransactionManagerAutoConfiguration;
+import org.springframework.boot.autoconfigure.mongo.MongoAutoConfiguration;
+import org.springframework.boot.autoconfigure.mongo.embedded.EmbeddedMongoAutoConfiguration;
+import org.springframework.context.annotation.Bean;
+import org.springframework.context.annotation.Import;
+import org.springframework.context.annotation.Primary;
+import org.springframework.data.gemfire.client.ClientRegionFactoryBean;
+import org.springframework.data.gemfire.config.annotation.EnablePdx;
+import org.springframework.data.mongodb.core.MongoTemplate;
+import org.springframework.data.mongodb.core.convert.MongoCustomConversions;
+import org.springframework.data.redis.core.RedisTemplate;
+import org.springframework.geode.boot.autoconfigure.ClientCacheAutoConfiguration;
+import org.springframework.integration.gemfire.store.GemfireMessageStore;
+import org.springframework.integration.jdbc.store.JdbcMessageStore;
+import org.springframework.integration.mongodb.store.ConfigurableMongoDbMessageStore;
+import org.springframework.integration.mongodb.support.BinaryToMessageConverter;
+import org.springframework.integration.mongodb.support.MessageToBinaryConverter;
+import org.springframework.integration.redis.store.RedisMessageStore;
+import org.springframework.integration.store.MessageGroupStore;
+import org.springframework.jdbc.core.JdbcTemplate;
+import org.springframework.util.StringUtils;
+
+
+/**
+ * A helper class containing configuration classes for particular technologies
+ * to expose an appropriate {@link org.springframework.integration.store.MessageStore} bean
+ * via matched configuration properties.
+ *
+ * @author Artem Bilan
+ */
+class MessageStoreConfiguration {
+
+ @ConditionalOnClass(ConfigurableMongoDbMessageStore.class)
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX,
+ name = "message-store-type",
+ havingValue = AggregatorFunctionProperties.MessageStoreType.MONGODB)
+ @Import({ MongoAutoConfiguration.class,
+ MongoDataAutoConfiguration.class,
+ EmbeddedMongoAutoConfiguration.class })
+ static class Mongo {
+
+ @Bean
+ public MessageGroupStore messageStore(MongoTemplate mongoTemplate, AggregatorFunctionProperties properties) {
+ if (StringUtils.hasText(properties.getMessageStoreEntity())) {
+ return new ConfigurableMongoDbMessageStore(mongoTemplate, properties.getMessageStoreEntity());
+ }
+ else {
+ return new ConfigurableMongoDbMessageStore(mongoTemplate);
+ }
+ }
+
+ @Bean
+ @Primary
+ public MongoCustomConversions mongoDbCustomConversions() {
+ return new MongoCustomConversions(Arrays.asList(
+ new MessageToBinaryConverter(), new BinaryToMessageConverter()));
+ }
+
+ }
+
+ @ConditionalOnClass(RedisMessageStore.class)
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX,
+ name = "message-store-type",
+ havingValue = AggregatorFunctionProperties.MessageStoreType.REDIS)
+ @Import(RedisAutoConfiguration.class)
+ static class Redis {
+
+ @Bean
+ public MessageGroupStore messageStore(RedisTemplate, ?> redisTemplate) {
+ return new RedisMessageStore(redisTemplate.getConnectionFactory());
+ }
+
+ }
+
+ @ConditionalOnClass(GemfireMessageStore.class)
+ @ConditionalOnProperty(prefix = AggregatorFunctionProperties.PREFIX,
+ name = "message-store-type",
+ havingValue = AggregatorFunctionProperties.MessageStoreType.GEMFIRE)
+ @Import(ClientCacheAutoConfiguration.class)
+ @EnablePdx
+ static class Gemfire {
+
+ @Bean
+ @ConditionalOnMissingBean
+ public ClientRegionFactoryBean, ?> gemfireRegion(GemFireCache cache, AggregatorFunctionProperties properties) {
+ ClientRegionFactoryBean, ?> clientRegionFactoryBean = new ClientRegionFactoryBean<>();
+ clientRegionFactoryBean.setCache(cache);
+ clientRegionFactoryBean.setName(properties.getMessageStoreEntity());
+ return clientRegionFactoryBean;
+ }
+
+ @Bean
+ public MessageGroupStore messageStore(Region