Azure Data Lake Camel消费者directoryName选项失效求助
解决Camel Azure Data Lake路由中directoryName不生效及获取文件名问题
一、directoryName配置不生效的解决方案
问题核心
使用pollEnrich()时,默认会创建独立的Exchange调用目标端点,原Exchange中的Header不会自动传递;同时若URI拼接逻辑不当,组件无法正确识别directoryName与fileName的组合关系。
可行方案
方案1:动态拼接URI参数
直接在pollEnrich的simple表达式中明确指定directoryName,确保组件能正确解析路径组合:
from("direct:updateValues") .pollEnrich() // 若getFrom()已包含directoryName,先移除再重新添加(避免重复参数) .simple(messageRouteConfig.getFrom().replace("&directoryName=/foo/bar", "") + "&directoryName=/foo/bar&fileName=${body}") .log(LoggingLevel.INFO, log, "message content (file content): ${body}") .log(LoggingLevel.INFO, log, "Starting processing file : ${header.CamelAzureStorageDataLakeFileName}")
如果directoryName是动态值,也可以从Header读取:
from("direct:updateValues") .setHeader("targetDir", constant("/foo/bar")) .pollEnrich() .simple(messageRouteConfig.getFrom() + "&directoryName=${header.targetDir}&fileName=${body}") // 后续处理逻辑
方案2:传递原Exchange的Header到pollEnrich请求
通过自定义逻辑将原Exchange的Header传递给pollEnrich的目标请求,确保组件能读取到CamelAzureStorageDataLakeDirectoryName:
from("direct:updateThings") .setHeader("CamelAzureStorageDataLakeDirectoryName", constant("/foo/bar")) .pollEnrich(exchange -> { String uri = messageRouteConfig.getFrom() + "&fileName=${body}"; Exchange targetExchange = exchange.getContext().createExchange(); // 复制原Header到目标Exchange targetExchange.getIn().setHeaders(exchange.getIn().getHeaders()); targetExchange.getIn().setBody(exchange.getIn().getBody()); return exchange.getContext().createProducerTemplate().request(uri, targetExchange); }) // 后续处理逻辑
二、获取已读取文件名的方法
${file:name}是Camel本地文件组件的表达式,不适用于Azure Data Lake组件。Azure Data Lake组件在读取文件后会设置专属Header,你需要使用${header.CamelAzureStorageDataLakeFileName}来获取文件名。
注意:pollEnrich默认不会将目标Exchange的Header合并到主Exchange中,需自定义AggregationStrategy来合并Header:
from("direct:updateValues") .pollEnrich( simple(messageRouteConfig.getFrom() + "&directoryName=/foo/bar&fileName=${body}"), new AggregationStrategy() { @Override public Exchange aggregate(Exchange oldExchange, Exchange newExchange) { if (newExchange != null) { // 合并目标Exchange的Header到主Exchange oldExchange.getIn().setHeaders(newExchange.getIn().getHeaders()); // 设置文件内容到主Exchange的Body oldExchange.getIn().setBody(newExchange.getIn().getBody()); } return oldExchange; } } ) .log(LoggingLevel.INFO, log, "message content (file content): ${body}") .log(LoggingLevel.INFO, log, "Starting processing file : ${header.CamelAzureStorageDataLakeFileName}")
内容的提问来源于stack exchange,提问作者Paehrin
相关产品推荐
相关产品推荐

