如何通过Apache Camel+Spring Batch处理文件并传递JAXBElement至Camel路由
解决方案:Apache Camel + Spring Batch 集成,将处理后的数据回传Camel路由
针对你提出的两个需求,我整理了一套完整的实现方案,替代原来的JMS Writer,让Spring Batch处理后的每行数据直接回到Camel路由中:
一、整体流程实现:Camel读文件 → Spring Batch处理 → 数据回传Camel
1. Camel路由配置(触发Batch Job)
首先,我们配置Camel路由读取文件,并将文件路径作为参数传递给Spring Batch Job:
import org.apache.camel.builder.RouteBuilder; import org.springframework.stereotype.Component; @Component public class FileToBatchRoute extends RouteBuilder { @Override public void configure() throws Exception { from("file:/path/to/input/directory?noop=true") // noop=true避免重复处理文件 .log("Starting processing file: ${file:name}") .setHeader("CamelSpringBatchJobName", constant("fileProcessingJob")) .setBody(simple("${file:absolute.path}")) // 将文件绝对路径作为Job参数传递 .to("spring-batch:fileProcessingJob"); // 调用Spring Batch Job } }
2. Spring Batch Job配置(处理文件并回传Camel)
核心是自定义ItemWriter,用Camel的ProducerTemplate将处理后的数据发送回Camel路由,替代原来的JMSTemplate:
import org.apache.camel.CamelContext; import org.apache.camel.ProducerTemplate; import org.springframework.batch.core.Job; import org.springframework.batch.core.Step; import org.springframework.batch.core.configuration.annotation.EnableBatchProcessing; import org.springframework.batch.core.configuration.annotation.JobBuilderFactory; import org.springframework.batch.core.configuration.annotation.StepBuilderFactory; import org.springframework.batch.core.launch.support.RunIdIncrementer; import org.springframework.batch.item.ItemProcessor; import org.springframework.batch.item.ItemReader; import org.springframework.batch.item.ItemWriter; import org.springframework.batch.item.file.FlatFileItemReader; import org.springframework.batch.item.file.mapping.PassThroughLineMapper; import org.springframework.batch.item.file.builder.FlatFileItemReaderBuilder; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.core.io.FileSystemResource; import javax.xml.bind.JAXBContext; import javax.xml.bind.JAXBElement; import javax.xml.bind.Unmarshaller; import java.io.StringReader; @Configuration @EnableBatchProcessing public class BatchProcessingConfig { @Autowired private JobBuilderFactory jobBuilderFactory; @Autowired private StepBuilderFactory stepBuilderFactory; @Autowired private CamelContext camelContext; // 自定义ItemWriter:将JAXBElement发送到Camel路由 @Bean public ItemWriter<JAXBElement<?>> camelItemWriter() { ProducerTemplate producerTemplate = camelContext.createProducerTemplate(); // 复用线程安全的ProducerTemplate return items -> { for (JAXBElement<?> item : items) { // 发送到Camel的后续处理端点 producerTemplate.sendBody("direct:postBatchProcessing", item); } }; } // ItemReader:读取文件的每一行内容 @Bean public ItemReader<String> fileItemReader(@Value("#{jobParameters['CamelSpringBatchJobBody']}") String filePath) { return new FlatFileItemReaderBuilder<String>() .name("fileLineReader") .resource(new FileSystemResource(filePath)) .lineMapper(new PassThroughLineMapper()) // 直接读取每行字符串 .build(); } // ItemProcessor:将每行字符串转换为JAXBElement @Bean public ItemProcessor<String, JAXBElement<?>> jaxbElementProcessor() { return line -> { // 替换成你的实际JAXB模型类 Class<?> modelClass = YourJaxbModel.class; JAXBContext jaxbContext = JAXBContext.newInstance(modelClass); Unmarshaller unmarshaller = jaxbContext.createUnmarshaller(); Object model = unmarshaller.unmarshal(new StringReader(line)); // 根据你的命名空间和模型调整QName return new JAXBElement<>(new QName("http://your.namespace", modelClass.getSimpleName()), modelClass, model); }; } // 定义Batch Step @Bean public Step fileProcessingStep(ItemReader<String> fileItemReader, ItemProcessor<String, JAXBElement<?>> jaxbElementProcessor, ItemWriter<JAXBElement<?>> camelItemWriter) { return stepBuilderFactory.get("fileProcessingStep") .<String, JAXBElement<?>>chunk(10) // 按10条数据为一个批次处理,可按需调整 .reader(fileItemReader) .processor(jaxbElementProcessor) .writer(camelItemWriter) .build(); } // 定义Batch Job @Bean public Job fileProcessingJob(Step fileProcessingStep) { return jobBuilderFactory.get("fileProcessingJob") .incrementer(new RunIdIncrementer()) // 保证每次Job运行ID唯一 .flow(fileProcessingStep) .end() .build(); } }
3. Camel后续路由(接收Spring Batch回传的数据)
配置Camel路由接收Spring Batch发送的JAXBElement,进行后续业务处理:
import org.apache.camel.builder.RouteBuilder; import org.springframework.stereotype.Component; @Component public class PostBatchProcessingRoute extends RouteBuilder { @Override public void configure() throws Exception { from("direct:postBatchProcessing") .log("Received JAXBElement from Spring Batch: ${body}") // 替换为你的实际业务逻辑,比如调用REST接口、存储到数据库等 .to("mock:finalDestination"); } }
二、关键注意事项
- 依赖引入:确保项目中引入
camel-spring-batch和spring-batch-core相关依赖,保证Camel与Spring Batch正常集成。 - ProducerTemplate复用:
ProducerTemplate是线程安全的,建议在camelItemWriter中复用同一个实例,减少频繁创建的开销。 - JAXBElement传递:JAXBElement本身支持序列化,Camel可直接传递,无需额外配置。
- Chunk大小调整:根据文件大小和系统性能,调整Step中的
chunk参数,平衡处理效率与内存占用。
内容的提问来源于stack exchange,提问作者Vijai
相关产品推荐
相关产品推荐

