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

Apache Camel:基于pollEnrich结果聚合并保留Header的路由实现问题

解决Apache Camel路由聚合PollEnrich结果并保留原消息Header的问题

我来帮你调整这个Camel路由,实现基于pollEnrich拉取的文件进行聚合,同时完整保留原消息的Header信息。你的核心需求是用原消息里的pad和file参数生成对应路径的Tar包,下面是优化后的路由方案:

修改后的完整路由代码

from("broker1:files.queue")
    .log("Starting processing for file prefix: ${header.file} in directory: ${header.pad}")
    // 把原消息的关键Header存到Exchange属性,避免后续操作覆盖
    .setProperty("originalPad", simple("${header.pad}"))
    .setProperty("originalFilePrefix", simple("${header.file}"))
    // 拉取指定目录下所有匹配前缀的文件,maxMessagesPerPoll=-1表示拉取全部符合条件的文件
    .pollEnrich()
        .simple("file:${header.pad}?antInclude=${header.file}.*&maxMessagesPerPoll=-1&readLock=none")
        .timeout(5000) // 设置超时,防止无文件时无限等待
    .end()
    // 聚合拉取到的文件生成Tar包
    .aggregate(new TarAggregationStrategy(false, true))
        // 用原始文件前缀作为聚合键,确保同一条消息的文件聚合成一个包
        .property("originalFilePrefix")
        .completionFromBatchConsumer()
        .eagerCheckCompletion()
        .parallelProcessing(false)
    .end()
    // 设置Tar包的最终路径和文件名
    .setHeader("CamelFileName", simple("${property.originalPad}/${property.originalFilePrefix}.tar"))
    // 将生成的Tar包输出到目标目录
    .to("file:${property.originalPad}")
    .log("Successfully created Tar package: ${header.CamelFileName}");

关键修改说明

  • 保留原始Header信息:通过setProperty把原消息的pad和file存到Exchange属性中,这样即使后续pollEnrich操作修改了Exchange的Body或部分Header,我们依然能拿到原始的路径和文件名参数,确保Tar包命名准确。
  • 优化PollEnrich拉取逻辑:添加maxMessagesPerPoll=-1让pollEnrich一次性拉取所有匹配${header.file}.*的文件(默认只拉取单个文件),同时设置timeout避免在目录下无匹配文件时无限等待。readLock=none可以根据你的实际场景调整,如果不需要文件锁定可以保留,需要的话换成合适的锁策略。
  • 精准聚合关联:把聚合的关联键从constant(true)改成property("originalFilePrefix"),这样每条源消息对应的文件集合会被独立聚合,不会和其他消息的文件混在一起生成错误的Tar包。
  • Tar包输出配置:用保存的原始属性设置CamelFileName,确保生成的Tar包直接输出到pad指定的目录,文件名符合file.tar的格式要求。

可选的异常处理优化

如果需要处理无匹配文件的场景,可以在pollEnrich之后添加判断逻辑,避免空聚合:

...
.end()
.choice()
    .when(body().isNull())
        .log("Warning: No files found matching prefix ${property.originalFilePrefix} in directory ${property.originalPad}")
    .otherwise()
        .aggregate(new TarAggregationStrategy(false, true))
            .property("originalFilePrefix")
            .completionFromBatchConsumer()
            .eagerCheckCompletion()
            .parallelProcessing(false)
        .end()
        .setHeader("CamelFileName", simple("${property.originalPad}/${property.originalFilePrefix}.tar"))
        .to("file:${property.originalPad}")
        .log("Successfully created Tar package: ${header.CamelFileName}")
.end()

内容的提问来源于stack exchange,提问作者Michiel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 03:12:37