【问题标题】:Kafka-Streams throwing NullPointerException when consumingKafka-Streams 在消费时抛出 NullPointerException
【发布时间】:2017-02-01 15:44:59
【问题描述】:

我有这个问题:

当我使用处理器 API 从主题中消费时,当处理器内部使用方法 context().forward(K, V) 时,Kafka Streams 会引发空指针异常。

这是它的堆栈跟踪:

Exception in thread "StreamThread-1" java.lang.NullPointerException
at org.apache.kafka.streams.processor.internals.StreamTask.forward(StreamTask.java:336)
at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:187)
at org.apache.kafka.streams.processor.ProcessorContext$forward.call(Unknown Source)
at org.codehaus.groovy.runtime.callsite.CallSiteArray.defaultCall(CallSiteArray.java:48)
at org.codehaus.groovy.runtime.callsite.AbstractCallSite.call(AbstractCallSite.java:113)
at org.codehaus.groovy.runtime.callsite.AbstractCallSite.call(AbstractCallSite.java:133)
at com.bnsf.ltf.processor.ConversionProcessor.process(ConversionProcessor.groovy:23)
at com.bnsf.ltf.processor.ConversionProcessor.process(ConversionProcessor.groovy)
at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:68)
at org.apache.kafka.streams.processor.internals.StreamTask.forward(StreamTask.java:338)
at org.apache.kafka.streams.processor.internals.ProcessorContextImpl.forward(ProcessorContextImpl.java:187)
at org.apache.kafka.streams.processor.internals.SourceNode.process(SourceNode.java:64)
at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:174)
at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:320)
at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:218)

我的 Gradle 依赖项如下所示:

compile('org.codehaus.groovy:groovy-all')
compile('org.apache.kafka:kafka-streams:0.10.0.0')

更新:我尝试使用版本 0.10.0.1,但仍然抛出相同的错误。

这是我正在构建的拓扑的代码...

 topologyBuilder.addSource('inboundTopic', stringDeserializer, stringDeserializer, conversionConfiguration.inTopic)
    .addProcessor('conversionProcess', new ProcessorSupplier() {
        @Override
        Processor get() {
            return conversionProcessor
        }
    }, 'inboundTopic')
    .addSink('outputTopic', conversionConfiguration.outTopic, stringSerializer, stringSerializer, 'conversionProcess')

    stream = new KafkaStreams(topologyBuilder, streamConfig)
    stream.start()

我的处理器如下所示:

@Override
void process(String key, String message) {
    // Call to a service and the return of the service is set on the
    // converted local variable named converted
    context().forward(key, converted)
    context().commit()
}

【问题讨论】:

标签: groovy kafka-consumer-api apache-kafka-streams


【解决方案1】:

直接提供您的处理器。

.addProcessor('conversionProcess', () -> new MyProcessor(), 'inboundTopic')

MyProcessor 应该反过来继承自 AbstractProcessor

【讨论】:

    猜你喜欢
    • 2020-04-02
    • 1970-01-01
    • 2020-12-30
    • 1970-01-01
    • 1970-01-01
    • 2020-10-13
    • 2020-08-14
    • 1970-01-01
    • 2015-11-03
    相关资源
    最近更新 更多