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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 09:31:34