You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用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.


Apache Camel Implementation

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_Id header is auto-generated by Camel when parsing CSV with headers, matching your Id column.
  • Adjust the aggregation strategy if you need to handle partial joins (e.g., left/right joins) or filter unmatched rows.

Spring Integration Implementation

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 Id field as a message header to use for correlation.
  • The groupTimeout ensures 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 Id values 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.29 07:57:56