How to use Aggregator EIP using Spring and Apache Camel?

Jasvinder S Saggu
1 min readJul 9, 2021

--

pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 https://maven.apache.org/xsd/maven-4.0.0.xsd">
<modelVersion>4.0.0</modelVersion>
<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>2.4.7</version>
<relativePath/> <!-- lookup parent from repository -->
</parent>
<groupId>com.jss</groupId>
<artifactId>camel</artifactId>
<version>0.0.1-SNAPSHOT</version>
<name>camel</name>
<description>Apache Camel with Spring Boot</description>

<properties>
<java.version>1.8</java.version>
<junit.platform.version>1.6.0</junit.platform.version>
<camel.version>3.7.3</camel.version>
</properties>

<dependencies>
<dependency>
<groupId>org.apache.camel.springboot</groupId>
<artifactId>camel-spring-boot-starter</artifactId>
<version>${camel.version}</version>
</dependency>

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>
</dependencies>
</project>

AggregatorRoute.java

package com.jss.camel.components.routes.aggregation;

import org.apache.camel.LoggingLevel;
import org.apache.camel.Message;
import org.apache.camel.builder.RouteBuilder;
import org.springframework.stereotype.Component;

import java.util.Date;
import java.util.Random;

@Component
public class AggregatorRoute extends RouteBuilder {
final String CORRELATION_ID = "correlationId";

@Override
public void configure() throws Exception {

Random random = new Random();

from("timer:insurance?period=200")
.process(exchange ->
{
Message message = exchange.getMessage();
message.setHeader(CORRELATION_ID, random.nextInt(4));
message.setBody(new Date() + "");
})
.aggregate(header(CORRELATION_ID), new MyAggregationStrategy())
.completionSize(5)
.log(LoggingLevel.ERROR, "${header." + CORRELATION_ID + "} ${body}")
;
}
}

MyAggregationStrategy.java

package com.jss.camel.components.routes.aggregation;

import org.apache.camel.Exchange;

import java.util.Objects;

public class MyAggregationStrategy implements org.apache.camel.AggregationStrategy {
@Override
public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
if (Objects.isNull(oldExchange)) {
return newExchange;
}

String oldBody = oldExchange.getIn().getBody(String.class);
String newBody = newExchange.getIn().getBody(String.class);

String aggBody = oldBody + "->" + newBody;

oldExchange.getIn().setBody(aggBody);

return oldExchange;
}
}

--

--