如何通过Camel Stream组件从UNIX命名管道读取完整XML文档消息?
解决Camel消费UNIX命名管道XML流时拆分完整文档的问题
嘿Dave,我刚好碰到过类似的场景,给你几个可行的思路参考下:
思路1:使用聚合器(Aggregate)组件按XML文档边界聚合消息
因为Stream组件默认按行拆分,那我们可以把这些行消息聚合起来,直到识别出一个完整的XML文档结束标记,再输出完整消息。
具体步骤:
- 先用
stream:file组件读取命名管道(UNIX命名管道可以当作文件路径处理,比如/tmp/my-named-pipe) - 接着用
aggregate组件,自定义聚合策略:- 聚合策略里维护一个字符串缓冲区,每次收到行消息就追加到缓冲区
- 检查缓冲区是否包含完整的XML结束标签(比如你的XML根元素是
<order>,就检查是否有</order>) - 一旦检测到完整文档,就输出这个缓冲区的内容,并重置缓冲区准备下一个文档
示例代码片段:
from("stream:file?fileName=/tmp/my-named-pipe&charset=UTF-8") .aggregate(constant(true), new XmlDocumentAggregationStrategy()) .completionPredicate(exchange -> { String body = exchange.getIn().getBody(String.class); // 这里替换成你的XML根元素结束标签 return body.contains("</root>"); }) .completionTimeout(5000) // 超时兜底,防止死等 .to("direct:processCompleteXml"); // 自定义聚合策略类 class XmlDocumentAggregationStrategy implements AggregationStrategy { @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { if (oldExchange == null) { return newExchange; } String oldBody = oldExchange.getIn().getBody(String.class); String newBody = newExchange.getIn().getBody(String.class); oldExchange.getIn().setBody(oldBody + newBody); return oldExchange; } }
思路2:给Stream组件指定XML文档级别的分隔符
如果你的所有XML文档都有固定的结束标记,而且这个标记不会出现在XML内容里(比如根元素固定),可以直接给Stream组件设置delimiter参数,把完整的XML结束标签作为分隔符,这样Stream组件会自动把每个分隔符之间的内容当作一个完整消息。
示例配置:
// 注意:delimiter需要转义特殊字符,比如XML的尖括号 from("stream:file?fileName=/tmp/my-named-pipe&charset=UTF-8&delimiter=</root>") .process(exchange -> { // 这里的body就是完整的XML文档,可能需要补回被截断的结束标签 String body = exchange.getIn().getBody(String.class); exchange.getIn().setBody(body + "</root>"); }) .to("direct:processCompleteXml");
注意:这个方法的局限性是如果XML内容里包含和结束标签一样的字符串(比如CDATA块里),会导致错误拆分,所以适合XML结构固定且内容不会出现根结束标签的场景。
思路3:自定义DataFormat处理XML流
如果上面的方法都不满足需求,可以自定义一个Camel DataFormat,底层用SAX或者StAX解析器来读取流,每当解析完一个完整的XML文档,就输出这个文档的内容。这种方法最可靠,能正确处理各种XML格式(包括带CDATA、命名空间的情况)。
大致实现步骤:
- 实现
DataFormat接口,在unmarshal方法中用StAX解析器读取输入流 - 每次解析完一个完整的XML文档,就把对应的XML字符串写入输出Exchange
- 在路由中使用这个自定义DataFormat
示例代码框架:
from("stream:file?fileName=/tmp/my-named-pipe&charset=UTF-8") .unmarshal(new XmlStreamDataFormat()) .to("direct:processCompleteXml"); // 自定义XML流DataFormat class XmlStreamDataFormat implements DataFormat { @Override public void marshal(Exchange exchange, Object graph, OutputStream stream) throws Exception { // 不需要实现marshal,因为我们只消费 } @Override public Object unmarshal(Exchange exchange, InputStream stream) throws Exception { XMLInputFactory factory = XMLInputFactory.newInstance(); XMLEventReader reader = factory.createXMLEventReader(stream); StringBuilder xmlBuilder = new StringBuilder(); int depth = 0; while (reader.hasNext()) { XMLEvent event = reader.nextEvent(); if (event.isStartElement()) { depth++; } xmlBuilder.append(event.toString()); if (event.isEndElement()) { depth--; if (depth == 0) { // 一个完整XML文档解析完成,返回 return xmlBuilder.toString(); } } } return null; } }
另外还要注意UNIX命名管道的特性:如果管道没有写入端,Stream组件可能会阻塞或者退出,你可以设置stream:file的readTimeout参数来处理这种情况。
内容的提问来源于stack exchange,提问作者Dave Hinkle
相关产品推荐
相关产品推荐

