如何用Apache Camel Java DSL实现类似SQL Join的双数据源合并?
Great question! Both Apache Camel and Spring Integration are perfect fits for this kind of data joining scenario, and yes—routing-based aggregation is exactly the approach you’ll want to use here. Let’s walk through how to implement this with both frameworks, complete with Java DSL examples.
Camel provides built-in CSV parsing support and flexible aggregation strategies to handle the "join by Id" logic. Here's a complete example that reads two CSV files, merges rows with matching Id values, and outputs the result:
import org.apache.camel.builder.RouteBuilder; import org.apache.camel.dataformat.csv.CsvDataFormat; import org.apache.camel.processor.aggregate.AggregationStrategy; import org.apache.camel.Exchange; import java.util.HashMap; import java.util.Map; public class CsvJoinRoute extends RouteBuilder { @Override public void configure() throws Exception { // Configure CSV format to map rows to key-value pairs (using header row as keys) CsvDataFormat csvFormat = new CsvDataFormat(); csvFormat.setUseMaps(true); csvFormat.setHeader(true); // Custom aggregation strategy to merge rows with the same Id AggregationStrategy idBasedJoinStrategy = new AggregationStrategy() { @Override public Exchange aggregate(Exchange existingExchange, Exchange newExchange) { // Handle first entry in the group if (existingExchange == null) { return newExchange; } // Extract data from both exchanges Map<String, Object> existingRow = existingExchange.getIn().getBody(Map.class); Map<String, Object> newRow = newExchange.getIn().getBody(Map.class); // Merge rows (preserves all fields from both sources) Map<String, Object> mergedRow = new HashMap<>(existingRow); mergedRow.putAll(newRow); // Update the existing exchange with merged data existingExchange.getIn().setBody(mergedRow); return existingExchange; } }; // Route for first CSV file from("file:input?fileName=first.csv") .unmarshal(csvFormat) // Parse CSV into a list of maps .split(body()) // Process each row individually .aggregate(header("CsvRecord_Id"), idBasedJoinStrategy) // Group by Id field .completionFromBatchConsumer() // Wait until entire file is processed .marshal(csvFormat) // Convert merged maps back to CSV .to("file:output?fileName=joined_result.csv"); // Route for second CSV file (shares the same aggregation strategy) from("file:input?fileName=second.csv") .unmarshal(csvFormat) .split(body()) .aggregate(header("CsvRecord_Id"), idBasedJoinStrategy) .completionFromBatchConsumer() .marshal(csvFormat) .to("file:output?fileName=joined_result.csv"); } }
Key Notes for Camel:
- The
CsvRecord_Idheader is auto-generated by Camel when parsing CSV with headers, matching yourIdcolumn. - Adjust the aggregation strategy if you need to handle partial joins (e.g., left/right joins) or filter unmatched rows.
Spring Integration uses aggregators, CSV transformers, and file adapters to achieve the same result. Here's a Java DSL configuration example:
import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.file.dsl.Files; import org.springframework.integration.file.support.FileExistsMode; import org.springframework.integration.csv.transformer.CsvToMapTransformer; import org.springframework.integration.aggregator.HeaderAttributeCorrelationStrategy; import org.springframework.messaging.Message; import java.io.File; import java.util.HashMap; import java.util.List; import java.util.Map; @Configuration public class CsvJoinIntegrationConfig { // Flow to read and parse the first CSV file @Bean public IntegrationFlow firstCsvProcessingFlow() { return IntegrationFlows.from(Files.inboundAdapter(new File("input")) .patternFilter("first.csv")) .transform(new CsvToMapTransformer(true)) // Parse CSV rows to maps (uses header row) .split() // Split list of maps into individual row messages .enrichHeaders(h -> h.headerExpression("Id", "payload.get('Id')")) // Add Id as message header .channel("csvJoinChannel") .get(); } // Flow to read and parse the second CSV file @Bean public IntegrationFlow secondCsvProcessingFlow() { return IntegrationFlows.from(Files.inboundAdapter(new File("input")) .patternFilter("second.csv")) .transform(new CsvToMapTransformer(true)) .split() .enrichHeaders(h -> h.headerExpression("Id", "payload.get('Id')")) .channel("csvJoinChannel") .get(); } // Flow to aggregate rows by Id and write the result @Bean public IntegrationFlow joinAndWriteResultFlow() { return IntegrationFlows.from("csvJoinChannel") .aggregate(aggregator -> aggregator .correlationStrategy(new HeaderAttributeCorrelationStrategy("Id")) // Group by Id header .messageGroupProcessor((List<Message<?>> messages) -> { // Merge all rows in the group into one map Map<String, Object> mergedRow = new HashMap<>(); messages.forEach(msg -> mergedRow.putAll((Map<String, Object>) msg.getPayload())); return mergedRow; }) .groupTimeout(5000) // Wait 5s for all matching rows to arrive .sendPartialResultOnExpiry(true)) // Send results even if some groups are incomplete .transform(payload -> { // Convert merged map back to CSV row (handles headers automatically) Map<String, Object> row = (Map<String, Object>) payload; return String.join(",", row.values().stream().map(Object::toString).toList()); }) .handle(Files.outboundAdapter(new File("output")) .fileExistsMode(FileExistsMode.APPEND) .fileNameGenerator(msg -> "joined_result.csv")) .get(); } }
Key Notes for Spring Integration:
- We explicitly add the
Idfield as a message header to use for correlation. - The
groupTimeoutensures the aggregator doesn’t wait forever for matching rows (adjust based on your data processing timeline).
General Tips
- Both frameworks let you extend the aggregation logic to handle edge cases like unmatched
Idvalues or custom field merging rules. - For large datasets, consider adding batch processing or persistence for aggregation groups to avoid memory issues.
内容的提问来源于stack exchange,提问作者Johan Fredin

