You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

AWS Glue Schema Registry:同一Topic多事件类型反序列化优化咨询

问题描述

我们在同一个Kafka Topic中存储多种事件类型,已通过自定义命名策略为该Topic生成不同的Avro Schema。但消费者端若缺少对应事件的Avro文件,会触发反序列化失败。目前需添加所有事件类型的Schema才能正常消费,但实际仅需处理其中一种,希望找到优化方案。

若仅添加所需Schema,会触发如下异常:

java.lang.IllegalStateException: This error handler cannot process 'SerializationException's directly; please consider configuring an 'ErrorHandlingDeserializer' in the value and/or key deserializer
at org.springframework.kafka.listener.SeekUtils.seekOrRecover(SeekUtils.java:194)
at org.springframework.kafka.listener.SeekToCurrentErrorHandler.handle(SeekToCurrentErrorHandler.java:112)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1598)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1210)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: org.apache.kafka.common.errors.SerializationException: Error deserializing key/value for partition bazaar_identity_user-0
at offset 2505. If needed, please seek past the record to continue consumption.
Caused by: com.amazonaws.services.schemaregistry.exception.AWSSchemaRegistryException: Exception occurred while de-serializing Avro message
at com.amazonaws.services.schemaregistry.deserializers.avro.AvroDeserializer.deserialize(AvroDeserializer.java:103)
at com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryDeserializationFacade.deserialize(GlueSchemaRegistryDeserializationFacade.java:172)
at com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer.deserializeByHeaderVersionByte(GlueSchemaRegistryKafkaDeserializer.java:160)
at com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer.deserialize(GlueSchemaRegistryKafkaDeserializer.java:116)
at org.apache.kafka.common.serialization.Deserializer.deserialize(Deserializer.java:60)
at org.apache.kafka.clients.consumer.internals.Fetcher.parseRecord(Fetcher.java:1387)
at org.apache.kafka.clients.consumer.internals.Fetcher.access$3400(Fetcher.java:133)
at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.fetchRecords(Fetcher.java:1618)
at org.apache.kafka.clients.consumer.internals.Fetcher$CompletedFetch.access$1700(Fetcher.java:1454)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchRecords(Fetcher.java:687)
at org.apache.kafka.clients.consumer.internals.Fetcher.fetchedRecords(Fetcher.java:638)
at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1272)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1233)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1206)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doPoll(KafkaMessageListenerContainer.java:1410)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1249)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1161)
at java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:539)
at java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)
at java.base/java.lang.Thread.run(Thread.java:833)
Caused by: com.google.common.util.concurrent.UncheckedExecutionException: java.lang.NullPointerException
at com.google.common.cache.LocalCache$Segment.get(LocalCache.java:2051)
at com.google.common.cache.LocalCache.get(LocalCache.java:3951)
at com.google.common.cache.LocalCache.getOrLoad(LocalCache.java:3974)
at com.google.common.cache.LocalCache$LocalLoadingCache.get(LocalCache.java:4935)
at com.amazonaws.services.schemaregistry.deserializers.avro.AvroDeserializer.deserialize(AvroDeserializer.java:93)
... 19 common frames omitted
Caused by: java.lang.NullPointerException: null
优化方案

1. 配置ErrorHandlingDeserializer跳过无效消息

按照异常提示,用Spring Kafka的ErrorHandlingDeserializer包裹Glue Schema Registry的反序列化器,遇到无法解析的消息时直接跳过,不中断消费流程。

Spring Boot配置示例

# 用ErrorHandlingDeserializer包裹值反序列化器
spring.kafka.consumer.value-deserializer=org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
# 指定实际的反序列化器实现
spring.kafka.consumer.properties.spring.deserializer.value.delegate.class=com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer
# 配置错误处理器
spring.kafka.listener.error-handler=customSeekToCurrentErrorHandler

自定义错误处理器

import org.springframework.kafka.listener.SeekToCurrentErrorHandler;
import org.springframework.util.backoff.FixedBackOff;
import org.apache.kafka.common.errors.SerializationException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

@Bean
public SeekToCurrentErrorHandler customSeekToCurrentErrorHandler() {
    Logger log = LoggerFactory.getLogger(SeekToCurrentErrorHandler.class);
    return new SeekToCurrentErrorHandler((record, exception) -> {
        // 仅跳过序列化异常对应的消息
        if (exception.getCause() instanceof SerializationException) {
            log.warn("跳过无法反序列化的消息,Topic: {}, Partition: {}, Offset: {}", 
                     record.topic(), record.partition(), record.offset());
        } else {
            // 其他异常抛出,按原有逻辑处理
            throw new RuntimeException(exception);
        }
    }, new FixedBackOff(0L, 0L)); // 不重试,直接跳过
}

2. 提前过滤目标事件类型

在消费拉取阶段就过滤掉不需要的事件,避免触发反序列化异常。核心是通过消息头中的Schema ID识别事件类型:

  1. 从Glue Schema Registry获取目标事件对应的Schema ID
  2. 自定义反序列化器或拦截器,先读取消息头中的Schema ID,非目标ID直接返回null

自定义反序列化器示例

import com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer;
import org.apache.kafka.common.serialization.Deserializer;

public class TargetEventDeserializer implements Deserializer<YourTargetEvent> {
    private final GlueSchemaRegistryKafkaDeserializer delegate = new GlueSchemaRegistryKafkaDeserializer();
    // 替换为你的目标事件Schema ID
    private static final String TARGET_SCHEMA_ID = "your-target-schema-id";

    @Override
    public YourTargetEvent deserialize(String topic, byte[] data) {
        // 从消息头中提取Schema ID(格式参考Glue Schema Registry序列化规范)
        String schemaId = extractSchemaId(data);
        if (!TARGET_SCHEMA_ID.equals(schemaId)) {
            return null; // 跳过非目标事件
        }
        return (YourTargetEvent) delegate.deserialize(topic, data);
    }

    // 实现Schema ID提取逻辑,具体格式根据Glue的序列化规则调整
    private String extractSchemaId(byte[] data) {
        // 示例:前N字节为Schema ID标识,实际需参考Glue文档
        return new String(data, 0, 16);
    }
}

3. 使用兼容式通用Schema

如果所有事件类型有共同的顶层结构,可以定义一个包含union类型的通用Avro Schema,兼容所有事件类型。消费时先反序列化为通用结构,再判断是否为目标事件并转换。这种方式需要上游事件结构具备一定兼容性。


内容的提问来源于stack exchange,提问作者Shakeel Hussain

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.19 09:45:37