【问题标题】:KafkaStreams: Handling Deserialize exception in KStream-KTable JoinKafkaStreams:处理 KStream-KTable Join 中的反序列化异常
【发布时间】:2020-06-06 07:02:29
【问题描述】:

假设我们正在 KStream 和 KTable 之间进行内部连接,如下所示:

        StreamsBuilder sb = new StreamsBuilder();
        JsonSerde<SensorMetaData> sensorMetaDataJsonSerde = new JsonSerde<>(SensorMetaData.class);
        KTable<String, String> kTable = sb.stream("sensorMetadata",
                Consumed.with(Serdes.String(), sensorMetaDataJsonSerde)).toTable();

        KStream<String, String> kStream = sb.stream("sensorValues",
                Consumed.with(Serdes.String(), Serdes.String()));

        KStream<String, String> joined = kStream.join(kTable, (left, right)->{return getJoinedOutput(left, right);});

关于应用程序的几点:

  1. SensorMetaData 是一个 POJO

        public class SensorMetaData{        
        String sensorId;
        String sensorMetadata;      
        }
    
  2. DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG 设置为 org.apache.kafka.streams.errors.LogAndContinueExceptionHandler

  3. 如果反序列化失败,JsonSerde 类将抛出 SerializationException。

当我运行应用程序并向两个主题发送消息时,加入按预期工作。

现在我更改了 SensorMetaData 的架构如下,并在新节点上重新部署了应用程序

public class SensorMetaData{            
String sensorId;    
MetadataTag[] metadataTags;         
}

应用程序启动后,当我向 sensorValues 主题(流主题)发送消息时,应用程序正在关闭并出现 org.apache.kafka.common.errors.SerializationException。查看堆栈跟踪,我意识到由于 SensorMetaData 中的架构更改,它在执行连接时无法反序列化 SensorMetaData。 Deserialize 方法中的断点显示,它试图反序列化来自主题“app-KSTREAM-TOTABLE-STATE-STORE-0000000002-changelog”的数据。

所以问题是即使 DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG 设置为 org.apache.kafka.streams.errors.LogAndContinueExceptionHandler ,为什么应用程序会关闭而不是跳过错误记录(即具有旧模式的记录)?

但是,当应用程序在读取主题“sensorMetadata”(即 sb.stream("sensorMetadata"))时遇到错误记录时,它会成功跳过记录并发出警告“由于反序列化错误而跳过记录”。

为什么 join 没有跳过这里的不良记录?如何处理这种情况。我希望应用程序跳过记录并继续运行而不是关闭。这是堆栈跟踪

at kafkastream.JsonSerde$2.deserialize(JsonSerde.java:51)
    at org.apache.kafka.streams.state.internals.ValueAndTimestampDeserializer.deserialize(ValueAndTimestampDeserializer.java:54)
    at org.apache.kafka.streams.state.internals.ValueAndTimestampDeserializer.deserialize(ValueAndTimestampDeserializer.java:27)
    at org.apache.kafka.streams.state.StateSerdes.valueFrom(StateSerdes.java:160)
    at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.outerValue(MeteredKeyValueStore.java:207)
    at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.lambda$get$2(MeteredKeyValueStore.java:133)
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:821)
    at org.apache.kafka.streams.state.internals.MeteredKeyValueStore.get(MeteredKeyValueStore.java:133)
    at org.apache.kafka.streams.processor.internals.ProcessorContextImpl$KeyValueStoreReadWriteDecorator.get(ProcessorContextImpl.java:465)
    at org.apache.kafka.streams.kstream.internals.KTableSourceValueGetterSupplier$KTableSourceValueGetter.get(KTableSourceValueGetterSupplier.java:49)
    at org.apache.kafka.streams.kstream.internals.KStreamKTableJoinProcessor.process(KStreamKTableJoinProcessor.java:77)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.lambda$process$2(ProcessorNode.java:142)
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:806)
    at org.apache.kafka.streams.processor.internals.ProcessorNode.process(ProcessorNode.java:142)
    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:101)
    at org.apache.kafka.streams.processor.internals.StreamTask.lambda$process$3(StreamTask.java:383)
    at org.apache.kafka.streams.processor.internals.metrics.StreamsMetricsImpl.maybeMeasureLatency(StreamsMetricsImpl.java:806)
    at org.apache.kafka.streams.processor.internals.StreamTask.process(StreamTask.java:383)
    at org.apache.kafka.streams.processor.internals.AssignedStreamsTasks.process(AssignedStreamsTasks.java:475)
    at org.apache.kafka.streams.processor.internals.TaskManager.process(TaskManager.java:550)
    at org.apache.kafka.streams.processor.internals.StreamThread.runOnce(StreamThread.java:802)
    at org.apache.kafka.streams.processor.internals.StreamThread.runLoop(StreamThread.java:697)
    at org.apache.kafka.streams.processor.internals.StreamThread.run(StreamThread.java:670)
INFO stream-client [app-814c1c5b-a899-4cbf-8d85-2ed6eba81ccb] State transition from ERROR to PENDING_SHUTDOWN 

【问题讨论】:

    标签: apache-kafka-streams


    【解决方案1】:

    Kafka 在读取 RocksDB 文件时不使用 DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG 中的处理程序(请参阅堆栈跟踪提到类 StateSerdes)。这就是为什么它适用于来自源主题的记录,但在反序列化表中的数据时失败。

    我对 Kafka 不是很有经验,但我一遍又一遍地听到:如果发生变化,请将新格式的数据复制到另一个主题或删除数据、重置偏移量并重新处理。

    在这种情况下,也许最好删除 KTable 文件,用于 ktable 的内部主题,让应用程序重新生成具有新结构的 KTable。

    几个月前的这篇博客对删除数据的过程进行了更多解释:https://www.confluent.io/blog/data-reprocessing-with-kafka-streams-resetting-a-streams-application/

    分享一点见解:kafka 是一个非常复杂的野兽。要在生产环境中成功管理它,您需要构建大量工具和代码来维护它,并且(通常)更改您的部署过程以适应 Kafka。

    【讨论】:

      猜你喜欢
      • 2018-02-23
      • 1970-01-01
      • 2020-10-17
      • 1970-01-01
      • 1970-01-01
      • 2020-08-15
      • 1970-01-01
      • 2012-10-05
      • 2020-09-18
      相关资源
      最近更新 更多