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

Apache Camel动态路由源实现:按PDU_ID聚合文件并执行REST删除

Got it, let's break down what's going wrong and fix this route. The problem with chaining two from() calls is that each from() creates a separate, standalone route in Camel. Your second from() would start watching that file path as soon as the app boots, not when a message comes through your queue. That's why the dynamic routing isn't working.

Here's the revised route that correctly triggers file processing based on each queue message:

public class ToDeleteTarAndDeleteRoute extends RouteBuilder {
    @Override
    public void configure() throws Exception {
        from("broker1:todelete.message_ids.queue")
            // Clean up the path to ensure it ends with a slash (avoids messy path concatenation)
            .setProperty("cleanPadPath", simple("${header.pad.replaceAll('/*$', '/')}"))
            .setProperty("pduId", simple("${header.pdu_id}"))
            
            // Dynamically fetch all files matching the pdu_id prefix in the pad path
            // noop=true keeps source files intact (we'll delete them via REST later)
            // consumer.batch=true fetches all matching files in one go
            .pollEnrich()
                .simple("file://${exchangeProperty.cleanPadPath}?fileName=${exchangeProperty.pduId}.*&noop=true&consumer.batch=true")
                .timeout(1000) // Adjust timeout based on how long you want to wait for the batch
            
            // Aggregate the fetched files into a single Tar archive
            .aggregate(new TarAggregationStrategy())
                .constant(true) // Group all files from the batch into one Tar
                .completionFromBatchConsumer() // Finish aggregation when all files are processed
                .eagerCheckCompletion()
            
            // Save the Tar to the same pad path (using dynamic endpoint with toD)
            .toD("file://${exchangeProperty.cleanPadPath}?fileName=${exchangeProperty.pduId}.tar")
            
            .log("${exchangeProperty.pduId} successfully archived to Tar")
            
            // Prepare and send DELETE request to your REST API
            .setHeader(Exchange.HTTP_METHOD, constant("DELETE"))
            .setHeader("Connection", constant("Close"))
            .toD("http:127.0.0.1/restfuldb${exchangeProperty.pduId}?httpClient.authenticationPreemptive=true")
            
            .log("${exchangeProperty.pduId} archived and source files deleted via REST API");
    }
}

Key fixes and explanations:

  • Single entry point: The queue is now the only from() — each message triggers the full flow for that specific pdu_id and path.
  • pollEnrich for dynamic file fetching: This component lets us pull files from a dynamically generated URI (using values from the queue message headers). The consumer.batch=true flag ensures we get all matching files at once, not one by one.
  • Exchange properties: We store the cleaned path and pdu_id in properties to avoid reprocessing headers multiple times, and to handle path formatting cleanly (like ensuring the path ends with a slash).
  • toD instead of to: For endpoints where the URI depends on message data, use toD() (dynamic to) instead of static to() — this tells Camel to resolve the URI at runtime.
  • noop=true: Prevents the file component from deleting source files automatically — we rely on your REST API call to handle deletion later.

A couple of quick notes to adjust based on your setup:

  • Tweak the timeout value in pollEnrich to match how long you want Camel to wait for all matching files to be collected.
  • Double-check your REST API URI in the final toD to make sure it correctly targets the files associated with the pdu_id.

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 06:32:09