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识别事件类型:
- 从Glue Schema Registry获取目标事件对应的Schema ID
- 自定义反序列化器或拦截器,先读取消息头中的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
相关产品推荐
相关产品推荐

