【发布时间】: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);});
关于应用程序的几点:
-
SensorMetaData 是一个 POJO
public class SensorMetaData{ String sensorId; String sensorMetadata; } DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG 设置为 org.apache.kafka.streams.errors.LogAndContinueExceptionHandler
如果反序列化失败,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
【问题讨论】: