【问题标题】:Kafka Stream forward method throwing NullPointerException because ProcessorNode currentNode object is nullKafka Stream forward方法抛出NullPointerException,因为ProcessorNode currentNode对象为null
【发布时间】:2020-08-30 14:40:26
【问题描述】:

我们编写了 Kafka 流应用程序,它从源主题读取数据,在 2 个处理器中执行一些业务逻辑,然后将输出写入接收主题。

以下是创建拓扑、添加源、处理器(请注意我们添加 2 个处理器)和接收器的代码

Topology topology = new Topology();
topology.addSource("sourceProcessor", "source-topic")
        .addProcessor("Process", ()->fileExtractProcessorObject , "sourceProcessor")
        .addProcessor("XMLJSON", () -> xmlJson , "Process")
        .addSink("sinkProcessor", "sink-topic", "XMLJSON");
KafkaStreams streams = new KafkaStreams(topology, getProperties(appConfig,sourceConfig));
        streams.start()

下面是第一个处理器代码

public class FileExtractProcessor implements Processor<String, byte[]> {
    private ProcessorContext context;
    public FileExtractProcessor() {
    }

    @Override
    public void init(ProcessorContext context) {
        this.context = context;
    }

    @Override
    public void process(String key, byte[] bytearray) {
            /* business logic */

            //This is pojo class which will be forwarded to next processor as value
            ProcessData pData = new ProcessData();

                // Setting pojo object and forwarding to next processor
                pData.setValue(value);
                pData.setXmlData(files);
                pData.setAppconfig(ac);

                context.forward(value.getRunKeyId(), pData, To.child("XMLJSON"));
                context.commit();    
    }
}

在上面的代码中,当我们调用 forward 方法时,我们在下面的代码中得到空指针异常。

org.apache.kafka.streams.processor.internals.ProcessorContextImpl#forward(K, V, org.apache.kafka.streams.processor.To)

ProcessorNode child = this.currentNode().getChild(sendTo); 

java.lang.NullPointerException
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:183)
    at com.myapp.FileExtractProcessor.process(FileExtractProcessor.java:78)
    at com.myapp.FileExtractProcessor.process(FileExtractProcessor.java:22)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:118)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:201)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:180)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:133)
    at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:87)
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:429)
    at org.apache.kafka.streams.processor.internals.AssignedStreamsTasks.process(AssignedStreamsTasks.java:474)
    at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:536)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:792)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:698)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:671)

我们正在使用 kafka-stream-2.4.0.jar 你能帮我们在这里失踪的地方吗?无法弄清楚为什么 currentNode 对象为 Null。 非常感谢早期的帮助。

【问题讨论】:

  • 接收器主题的字符串中有错字:.addSink("sinkProcessor", ""sink-topic, "XMLJSON");。也许这已经解决了?
  • 这不是问题。那只是写问题时的拼写错误。运行第一个处理器时的问题
  • @chanduram,你如何创建处理器fileExtractProcessorObject?你会在任何地方分享吗?
  • @user207421 将问题标记为重复似乎不合适。可以撤消吗?
  • 可能是 KafkaStream 中的错误?你为什么要使用To.child()?您只有一个孩子,因此可以省略此参数。 -- 你能用 TopologyTestDriver 重现这个问题吗?如果是,你能提交一份错误报告吗?

标签: java apache-kafka apache-kafka-streams topology


【解决方案1】:

根据代码:https://github.com/apache/kafka/blob/2.4/streams/src/main/java/org/apache/kafka/streams/processor/internals/ProcessorContextImpl.java#L183 错误意味着currentNode() 返回null

这表明你违反了ProcessorSupplier的规则,在get()上返回一个新的Processor实例;当您将() -&gt; fileExtractProcessorObject 传递给addProcessor 时,这似乎成立,每次都返回相同的对象引用——相反,您每次都需要通过() -&gt; new FileExtractProcessor() 创建一个新实例。

我创建了一个工单来改进错误消息:https://issues.apache.org/jira/browse/KAFKA-10036

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2011-03-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2012-06-05
    相关资源
    最近更新 更多