【问题标题】:Spring Cloud Stream kafka-streams application shuts down on deserialization exceptionSpring Cloud Stream kafka-streams 应用程序因反序列化异常而关闭
【发布时间】:2021-06-18 10:25:34
【问题描述】:

我正在使用

弹簧靴:2.3.5.RELEASE 春云:Hoxton.SR8

我正在尝试 spring cloud stream kafka-streams 应用程序。一切都运行良好,直到出现反序列化异常。应用程序每次都会关闭。

我想跳过不良记录并在 Kafka 主题中继续前进。但我无法做到这一点。 配置:

spring:
  application:
    name: statsprocessor.${ENV}.${INSTANCE_ID}
  cloud:
    stream:
      instance-index: ${INSTANCE_INDEX}
      instance-count: ${INSTANCE_COUNT}
      bindings:
        statsInput:
          destination: ${STORE_INPUT_TOPIC}
          group: statsprocessor.${ENV}
          consumer:
            concurrency: ${CONCURRENCY}
            partitioned: true
            useNativeDecoding: true
      kafka:
        streams:
          bindings:
            statsInput:
              consumer:
                keySerde: org.apache.kafka.common.serialization.Serdes$StringSerde
                valueSerde: per.shades.framework.kafka.serdes.CotsEventSerde
                startOffset: earliest
                applicationId: statsprocessor.${ENV}
                autoCommitOnError: false
                dlqName: ${STORE_INPUT_DLQ}
                useNativeDecoding: true
                configuration:
                  client.id: statsprocessor.${ENV}.${INSTANCE_ID}
          binder:
            auto-add-partitions: true
            auto-create-topics: true
            deserializationExceptionHandler: logAndContinue
            brokers:
              - ${KAFKA_URI}
            configuration:
              num.stream.threads: ${CONCURRENCY}
              buffered.records.per.partition: 500
              cache.max.bytes.buffering: 10485760
              commit.interval.ms: 500
              state.dir: ${KAFKA_STATE_DIR}
              replication.factor: ${DEFAULT_KAFKA_TOPIC_REPLICATION_FACTOR}
              reconnect.backoff.ms: 15000
              retry.backoff.ms: 10000
              producer.linger.ms: 100
              producer.acks: all
              producer.retries: 3
              producer.batch.size: 16384
              consumer.max.poll.records: 100
              consumer.session.timeout.ms: 60000

我得到的错误是

    Exception in thread "statsprocessor.local.1-StreamThread-1" org.apache.kafka.streams.errors.StreamsException: 
Exception caught in process. taskId=0_4, processor=KSTREAM-SOURCE-0000000000, topic=cots-event-store, partition=4, offset=0, stacktrace=org.apache.kafka.streams.errors.StreamsException:
 ClassCastException invoking Processor. Do the Processor's input types match the deserialized types? Check the Serde setup and change the default Serdes in StreamConfig or provide correct Serdes via method parameters. 
Make sure the Processor can accept the deserialized input of type key: unknown because key is null, and value: per.shades.model.events.CotsEvent.

Note that although incorrect Serdes are a common cause of error, the cast exception might have another cause (in user code, for example). For example, if a processor wires in a store, but casts the generics incorrectly, a class cast exception could be raised during processing, but the cause would not be wrong Serdes.
    ....
    ....
    Caused by: java.lang.ClassCastException: class per.shades.model.events.CotsEvent cannot be cast to class per.shades.model.stats.StatsMetadata (per.shades.model.events.CotsEvent and per.shades.model.stats.StatsMetadata are in unnamed module of loader 'app')

现在我正在使用此设置deserializationExceptionHandler: logAndContinue。 仍然没有效果。根据文档,它应该只记录错误并继续处理。即它应该跳过不良记录。但这不会发生。看到这个错误。

All stream threads have died. The instance will be in error state and should be closed.

我也用过

@Bean
public StreamsBuilderFactoryBeanCustomizer streamsBuilderFactoryBeanCustomizer()
{
    return streamsBuilderFactoryBean ->
    {
        streamsBuilderFactoryBean.getStreamsConfiguration()                    .put(org.apache.kafka.streams.StreamsConfig.DEFAULT_DESERIALIZATION_EXCEPTION_HANDLER_CLASS_CONFIG,
                        ContinueOnErrorHandler.class);
    };
}

处理程序类是

public class ContinueOnErrorHandler implements DeserializationExceptionHandler
{

    @Override
    public DeserializationHandlerResponse handle(ProcessorContext processorContext, ConsumerRecord<byte[], byte[]> consumerRecord, Exception e)
    {
        System.out.println(">>>>>>> We are here");
        return DeserializationHandlerResponse.CONTINUE;
    }

    @Override
    public void configure(Map<String, ?> map)
    {

    }
}

但这也行不通。它没有被调用。

我不想删除我的 Kafka 主题来摆脱不良记录。真的很难解决简单的反序列化错误。请帮忙!!

编辑:绑定代码:

@Configuration
public interface StatsStreamBindings
{
    String statsInput = "statsInput";

    @Input(statsInput)
    KStream<String, StatsMetadata> statsInput();

}

处理器签名

public void aggregateStats(KStream<String, StatsMetadata> inputStream)

【问题讨论】:

    标签: java spring apache-kafka apache-kafka-streams spring-cloud-stream


    【解决方案1】:

    原因:java.lang.ClassCastException:类 per.shades.model.events.CotsEvent 不能转换为类 per.shades.model.stats.StatsMetadata(per.shades.model.events.CotsEvent 和 per.shades .model.stats.StatsMetadata 位于加载程序“app”的未命名模块中)

    你没有得到DeserializationException,你得到的是ClassCastException invoking Processor;这意味着反序列化器成功反序列化为CotsEvent,但您的处理器需要StatsMetadata

    DeserializationExceptionHandler 只处理DeserializationExceptions。

    【讨论】:

    • 好的。我怀疑是不是这样。但是,如果它不是反序列化错误,我该如何捕获它?因为我还添加了一个未捕获的异常错误处理程序。但这并没有被调用。请查看我的编辑,还添加了绑定代码和处理器签名。我收到的错误是我无法在代码中捕获的。所以请指教
    • 看起来很清楚。您已经配置了一个“CotsEventSerde”,但您期待一个“StatsMetadata”。您需要反序列化为预期的类型。
    • 是的,我知道。但是我不想改变我的反序列化器。我想确保错误消息不会破坏我的应用程序。因此,如果我发送错误的消息类型,应用程序应该能够跳过错误消息。目前,应用程序每次启动时都会卡住处理相同的错误消息。它无法跳过那些坏消息
    • 那么问题就完全错误了。它与反序列化错误无关。我建议您删除此内容并提出一个新问题,清楚地解释您想要实现的目标。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-11-19
    • 2019-02-15
    • 1970-01-01
    • 2019-10-30
    • 1970-01-01
    相关资源
    最近更新 更多