NiFi 1.18.0 ConsumeAMQP处理器空指针异常及消费问题排查
NiFi ConsumeAMQP处理器NullPointerException问题修复方案
问题描述
使用NiFi 1.18.0默认的ConsumeAMQP处理器,RabbitMQ连接正常,但接收消息时抛出NullPointerException。异常发生后,RabbitMQ端消费者数量持续增长,推测是抛出异常的消费者无法正常终止。需要实现无异常消费,并解决convertMapToString方法引发的问题。
RabbitMQ消息详情
Exchange EXCHANGE.EX Routing Key ROUTING.TEST Redelivered ● Properties app_id: UI.XXX|TEST|0|1 user_id: XYZUSER timestamp: 123456 priority: 0 delivery_mode: 2 headers: ClassName: Test.UI.X.Y FileName: undefined LineNumber: 0 MethodName: MethodName singularityheader: notxdetect=True content_encoding: utf8 content_type: text/plain Payload 278 bytes Encoding: string TestLog| Testlogcontent| Testlog value....
异常日志
2022-11-30 13:22:56,818 ERROR [Timer-Driven Process Thread-9] o.a.nifi.amqp.processors.ConsumeAMQP ConsumeAMQP[id=b0520c61-32aa-386d-20ac-9a2a34c14a30] Processor failure java.lang.NullPointerException: null at org.apache.nifi.amqp.processors.ConsumeAMQP.convertMapToString(ConsumeAMQP.java:237) at org.apache.nifi.amqp.processors.ConsumeAMQP.buildHeaders(ConsumeAMQP.java:220) at org.apache.nifi.amqp.processors.ConsumeAMQP.buildAttributes(ConsumeAMQP.java:192) at org.apache.nifi.amqp.processors.ConsumeAMQP.processResource(ConsumeAMQP.java:172) at org.apache.nifi.amqp.processors.ConsumeAMQP.processResource(ConsumeAMQP.java:47) at org.apache.nifi.amqp.processors.AbstractAMQPProcessor.onTrigger(AbstractAMQPProcessor.java:223) at org.apache.nifi.processor.AbstractProcessor.onTrigger(AbstractProcessor.java:27) at org.apache.nifi.controller.StandardProcessorNode.onTrigger(StandardProcessorNode.java:1354) at org.apache.nifi.controller.tasks.ConnectableTask.invoke(ConnectableTask.java:246) at org.apache.nifi.controller.scheduling.TimerDrivenSchedulingAgent$1.run(TimerDrivenSchedulingAgent.java:102) at org.apache.nifi.engine.FlowEngine$2.run(FlowEngine.java:110) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) 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:750) 2022-11-30 13:22:57,001 ERROR [Timer-Driven Process Thread-9] o.a.n.c.r.StandardProcessSession Failed to asynchronously commit session StandardProcessSession[id=323331374] for ConsumeAMQP[id=b0520c61-32aa-386d-20ac-9a2a34c14a30] org.apache.nifi.processor.exception.FlowFileHandlingException: StandardFlowFileRecord[uuid=db0849b7-f695-4f7c-b5a7-6bffb6c33d8d,claim=StandardContentClaim [resourceClaim=StandardResourceClaim[id=1669802459297-401, container=default, section=401], offset=33821, length=147],offset=0,name=db0849b7-f695-4f7c-b5a7-6bffb6c33d8d,size=147] transfer relationship not specified. This FlowFile was created in this session and was not transferred to any Relationship via ProcessSession.transfer() at org.apache.nifi.controller.repository.StandardProcessSession.validateCommitState(StandardProcessSession.java:259) at org.apache.nifi.controller.repository.StandardProcessSession.checkpoint(StandardProcessSession.java:274) at org.apache.nifi.controller.repository.StandardProcessSession.commit(StandardProcessSession.java:556) at org.apache.nifi.controller.repository.StandardProcessSession.commitAsync(StandardProcessSession.java:510) at org.apache.nifi.processor.AbstractProcessor.onTrigger(AbstractProcessor.java:28) at org.apache.nifi.controller.StandardProcessorNode.onTrigger(StandardProcessorNode.java:1354) at org.apache.nifi.controller.tasks.ConnectableTask.invoke(ConnectableTask.java:246) at org.apache.nifi.controller.scheduling.TimerDrivenSchedulingAgent$1.run(TimerDrivenSchedulingAgent.java:102) at org.apache.nifi.engine.FlowEngine$2.run(FlowEngine.java:110) at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511) at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180) at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294) 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:750)
问题原因
从异常堆栈来看,NullPointerException发生在ConsumeAMQP.convertMapToString方法的237行。查看NiFi 1.18.0的源码,该方法负责将RabbitMQ消息的headers转换为字符串,但没有对header的value做空值检查。如果消息headers中存在null值或无法调用toString()的对象,就会触发NPE。异常导致处理器流程中断,FlowFile未被正确转移,同时消费者线程无法正常关闭,造成RabbitMQ端消费者数量泄漏。
修复方案
1. 临时规避:清理消息headers
- 检查消息生产者,确保发送的所有header值不为null,且均为可转换为字符串的类型。比如确认
FileName字段是否实际为null而非字符串"undefined"。 - 若无法修改生产者,可在RabbitMQ端配置Exchange过滤器,移除可能包含null值的headers。
2. 彻底修复:升级NiFi版本
NiFi官方在1.19.0及后续版本中修复了该问题,在convertMapToString方法中添加了空值判断,将null值转换为"null"字符串,同时优化了异常处理逻辑,避免消费者泄漏。直接升级到更高版本即可解决问题。
3. 自定义处理器修复
若无法升级NiFi,可修改ConsumeAMQP处理器代码,添加空值检查:
private String convertMapToString(Map<String, Object> map) { if (map == null) { return ""; } StringBuilder sb = new StringBuilder(); for (Map.Entry<String, Object> entry : map.entrySet()) { // 新增空值判断,避免NPE String value = entry.getValue() != null ? entry.getValue().toString() : "null"; sb.append(entry.getKey()).append("=").append(value).append(","); } if (sb.length() > 0) { sb.setLength(sb.length() - 1); } return sb.toString(); }
将修改后的代码编译打包成nar文件,替换NiFi lib目录下的默认AMQP处理器包,重启NiFi生效。
4. 解决消费者泄漏问题
- 升级或修复处理器后,异常处理逻辑会正常关闭消费者,无需额外操作。
- 临时解决:重启ConsumeAMQP处理器或NiFi节点,清理泄漏的消费者连接。
内容的提问来源于stack exchange,提问作者Efecan AHMETOĞLU
相关产品推荐
相关产品推荐

