【问题标题】:Periodic NPE In Kafka Streams Processor ContextKafka 流处理器上下文中的周期性 NPE
【发布时间】:2016-08-21 19:03:33
【问题描述】:

使用 kafka-streams 0.10.0.0,我在转发消息时会定期在 StreamTask 中看到空指针异常。它在调用的 10% 到 50% 之间变化。 NPE出现在这种方法中:

public <K, V> void forward(K key, V value) {
    ProcessorNode thisNode = currNode;
    try {
        for (ProcessorNode childNode : (List<ProcessorNode<K, V>>) thisNode.children()) {
            currNode = childNode;
            childNode.process(key, value);
        }
    } finally {
        currNode = thisNode;
    }
}

似乎在某些情况下,thisNode 字段为空。知道可能是什么原因造成的吗?堆栈跟踪如下。

[ERROR] 2016-08-21 14:50:39.288 [StreamThread-1] StreamedMetricMeter - Forwarding failed
java.lang.NullPointerException
    at org.apache.kafka.streams.processor.internals.StreamTask.forward(StreamTask.java:336) ~[kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:187) ~[kafka-streams-0.10.0.0.jar:?]
    at com.heliosapm.streams.metrics.processors.AbstractStreamedMetricProcessor.forward(AbstractStreamedMetricProcessor.java:552) [classes/:?]
    at com.heliosapm.streams.metrics.processors.impl.StreamedMetricMeter.doProcess(StreamedMetricMeter.java:89) [classes/:?]
    at com.heliosapm.streams.metrics.processors.impl.StreamedMetricMeter.doProcess(StreamedMetricMeter.java:1) [classes/:?]
    at com.heliosapm.streams.metrics.processors.AbstractStreamedMetricProcessor.process(AbstractStreamedMetricProcessor.java:166) [classes/:?]
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:68) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.StreamTask.forward(StreamTask.java:338) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:187) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:64) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:174) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:320) [kafka-streams-0.10.0.0.jar:?]
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:218) [kafka-streams-0.10.0.0.jar:?]

【问题讨论】:

  • 你能分享你的拓扑代码吗?你试过0.10.0.1 吗?
  • 想通了。见答案。感谢您的参与。那个程序员错误是如此令人震惊,一旦我意识到,我没有使用 0.10.0.1 进行测试。
  • 顺便说一下,这个问题也出现在0.10.2.1。您的回答是救命稻草,谢谢!

标签: apache-kafka apache-kafka-streams


【解决方案1】:

问题是我的ProcessorSuppliers 在每次调用get 时都返回相同的处理器实例。反过来,Kafka Streams 引擎试图创建多个处理器实例,我毫无疑问地创建了多线程垃圾箱火灾。注意同样粗心.... ProcessorSupplier.get() 应该在每次调用时返回一个新的处理器实例。

【讨论】:

猜你喜欢
  • 2014-08-04
  • 2018-09-10
  • 2019-06-02
  • 1970-01-01
  • 1970-01-01
  • 2015-11-27
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多