在我的驼峰路线中,我使用队列中的消息;每条消息都包含标题“pad”(路径)和文件前缀。例如。:
消息1:pad=“/some/dir”,file=“AAA”
消息2:pad=“/另一个/dir”,file=“BRD”
根据邮件,我希望创建一个文件:
消息1:/some/dir/AAA。tar(包含所有文件/部分/目录/AAA*)
消息2:/other/dir/BRD。tar(包含/other/dir/BRD.tar中的所有文件)
目录和文件名收集在另一个路由中。
到目前为止,我有一条骆驼路线:
from("broker1:files.queue")
.log("starting with message ${header.file}")
.pollEnrich()
.simple("file:${header.pad}?antInclude=${header.file}.*")
.aggregate(new TarAggregationStrategy(false,true))
.constant(true)
.completionFromBatchConsumer()
.eagerCheckCompletion()
.parallelProcessing(false)
.setHeader("file", simple("${header.file}"))
.setHeader("pad", simple("${header.pad}"))
.log("tarring to: ${header.pad}${header.file}.tar")
.setHeader(Exchange.FILE_NAME, simple("${header.file}.tar"))
.setHeader(Exchange.FILE_PATH, simple("${header.pad}"))
.to("file://ignored")
.log("Going to do other stuff here on ${header.file}");
我这里有几个问题:
-运行此路由时,在看到日志行“tarring to”之前,我看到多个“start with message”行
-日志行“tarring to”实际上表示“.tar”,标题为空。。。
-创建的“.tar”文件存储在“..ignored”中,并包含来自每个jms消息文件头的一个文件。
这让我相信,聚合发生在我没有预料到的水平上;我想聚合pollEnrich的结果,而不是队列中其他消息的结果。为什么,我怎样才能让它表现得像我想要的那样?
另一种是丢失的报头;这可能是由于聚合了错误的项目。。。无论如何,我认为聚合中的setHeader()应该设置它们,但它们还是丢失了;我如何保存它们?
我对camel编程比较陌生;所以请原谅我的缺点;代码中的缩进是
认为
范围应为;这可能是完全错误的。我正在使用camel-2.20.1,但可以切换到任何其他版本。
编辑
阅读让我改变了一点路线;如评论中所述;现在看起来是这样的:(TarAggregationStrategy()是在我的CamelContext中创建的,并添加到注册表中)
from("broker1:files.queue")
.log("starting with message ${header.file}")
.pollEnrich()
.simple("file:${header.pad}?antInclude=${header.file}.*")
.aggregationStrategyRef("tarAggregationStrategy")
.log("tarring to: ${header.pad}${header.file}.tar")
.setHeader(Exchange.FILE_NAME, simple("${header.file}.tar"))
.setHeader(Exchange.FILE_PATH, simple("${header.pad}"))
.to("file://ignored")
.log("Going to do other stuff here on ${header.file}");
现在看来情况确实好转了;除了由于无法根据堆栈跟踪创建临时文件而没有发生实际的tar之外:
org.apache.camel.component.file.GenericFileOperationFailedException: Could not make temp file (c9db039a-1585-4e63-85dc-e21ca268b290)
at org.apache.camel.processor.aggregate.tarfile.TarAggregationStrategy.aggregate(TarAggregationStrategy.java:174)
at org.apache.camel.processor.PollEnricher.process(PollEnricher.java:280)
at org.apache.camel.processor.RedeliveryErrorHandler.process(RedeliveryErrorHandler.java:548)
at org.apache.camel.processor.CamelInternalProcessor.process(CamelInternalProcessor.java:201)
at org.apache.camel.processor.Pipeline.process(Pipeline.java:138)
at org.apache.camel.processor.Pipeline.process(Pipeline.java:101)
at org.apache.camel.processor.CamelInternalProcessor.process(CamelInternalProcessor.java:201)
at org.apache.camel.processor.DelegateAsyncProcessor.process(DelegateAsyncProcessor.java:97)
at org.apache.camel.component.jms.EndpointMessageListener.onMessage(EndpointMessageListener.java:112)
at org.springframework.jms.listener.AbstractMessageListenerContainer.doInvokeListener(AbstractMessageListenerContainer.java:719)
at org.springframework.jms.listener.AbstractMessageListenerContainer.invokeListener(AbstractMessageListenerContainer.java:679)
at org.springframework.jms.listener.AbstractMessageListenerContainer.doExecuteListener(AbstractMessageListenerContainer.java:649)
at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.doReceiveAndExecute(AbstractPollingMessageListenerContainer.java:317)
at org.springframework.jms.listener.AbstractPollingMessageListenerContainer.receiveAndExecute(AbstractPollingMessageListenerContainer.java:255)
at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.invokeListener(DefaultMessageListenerContainer.java:1166)
at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.executeOngoingLoop(DefaultMessageListenerContainer.java:1158)
at org.springframework.jms.listener.DefaultMessageListenerContainer$AsyncMessageListenerInvoker.run(DefaultMessageListenerContainer.java:1055)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
Caused by: java.io.IOException: Could not make temp file (c9db039a-1585-4e63-85dc-e21ca268b290)
at org.apache.camel.processor.aggregate.tarfile.TarAggregationStrategy.addFileToTar(TarAggregationStrategy.java:199)
at org.apache.camel.processor.aggregate.tarfile.TarAggregationStrategy.aggregate(TarAggregationStrategy.java:167)
... 19 more
我注意到,after(和)之间的部分无法创建临时文件,实际上是正文的内容(我可以将其留空,但没有明显的原因,我填写了文件id)