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

使用AWS Glue Schema Registry反序列化器读取AVRO消息失败

问题描述

搭建了基于MSK、MSK Connect、Debezium Postgres源连接器和AWS Glue Schema Registry的Kafka管道,生产者端可正常将带Schema的AVRO记录发布到Glue Schema Registry,但使用Confluent的kafka-avro-console-consumer读取消息时出现空指针错误,无法完成反序列化。

连接器配置

# Glue Schema Registry Specific Converters    
"key.converter" = "com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter"    
"key.converter.schemas.enable"   = false    
"value.converter"= "com.amazonaws.services.schemaregistry.kafkaconnect.AWSKafkaAvroConverter"    
"value.converter.schemas.enable" = true

"key.converter.region"                        = "us-east-1"    
"key.converter.registry.name"                 = "<REGISTRY NAME>"    
"key.converter.compatibility"                 = "FULL"    
"key.converter.schemaAutoRegistrationEnabled" = true    
"key.converter.dataFormat"="AVRO"    
"key.converter.avroRecordType"="GENERIC_RECORD"    
"key.converter.schemaNameGenerationClass" = "<SCHEMA NAME GENERATION CLASS>"

"value.converter.region"                        = "us-east-1"    
"value.converter.registry.name"                 = "<REGISTRY NAME>"    
"value.converter.compatibility"                 = "FULL"    
"value.converter.schemaAutoRegistrationEnabled" = true    
"value.converter.dataFormat"="AVRO"    
"value.converter.avroRecordType"="GENERIC_RECORD"    
"key.converter.schemaNameGenerationClass" = "<SCHEMA NAME GENERATION CLASS>"

消费者命令

kafka-avro-console-consumer --bootstrap-server <bootstrap_server_url> \
--consumer.config client.properties \
--property schema.registry.url=https://glue.us-east-1.amazonaws.com \
--property print.key=true \
--property print.value=true \
--key-deserializer com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer \
--value-deserializer com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer
--topic platform_avro_users --from-beginning

消费者配置文件client.properties

security.protocol=SASL_SSL
sasl.mechanism=AWS_MSK_IAM
sasl.jaas.config=software.amazon.msk.auth.iam.IAMLoginModule required;
sasl.client.callback.handler.class=software.amazon.msk.auth.iam.IAMClientCallbackHandler

key.deserializer=com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer
key.deserializer.region=us-east-1
key.deserializer.registry.name=<REGISTRY NAME>
key.deserializer.avroRecordType=GENERIC_RECORD
key.deserializer.schemaNameGenerationClass=<SCHEMANAME GENERATION CLASS NAME>

value.deserializer=com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryKafkaDeserializer
value.deserializer.region=us-east-1
value.deserializer.registry.name=<REGISTRY NAME>
value.deserializer.avroRecordType=GENERIC_RECORD
value.deserializer.schemaNameGenerationClass=<SCHEMANAME GENERATION CLASS NAME>

错误信息

Processed a total of 1 messages
ERROR Unknown error when running consumer:  (kafka.tools.ConsoleConsumer$:44)
java.lang.NullPointerException: Cannot invoke "com.amazonaws.services.schemaregistry.deserializers.GlueSchemaRegistryDeserializationFacade.deserialize(com.amazonaws.services.schemaregistry.common.AWSDeserializerInput)" because "this.glueSchemaRegistryDeserializationFacade" is null
    at com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer.deserializeByHeaderVersionByte(AWSKafkaAvroDeserializer.java:149)
    at com.amazonaws.services.schemaregistry.deserializers.avro.AWSKafkaAvroDeserializer.deserialize(AWSKafkaAvroDeserializer.java:114)
    at io.confluent.kafka.formatter.AvroMessageFormatter$AvroMessageDeserializer.deserializeKey(AvroMessageFormatter.java:125)
    at io.confluent.kafka.formatter.SchemaMessageFormatter.writeTo(SchemaMessageFormatter.java:157)
    at kafka.tools.ConsoleConsumer$.process(ConsoleConsumer.scala:116)
    at kafka.tools.ConsoleConsumer$.run(ConsoleConsumer.scala:76)
    at kafka.tools.ConsoleConsumer$.main(ConsoleConsumer.scala:53)
    at kafka.tools.ConsoleConsumer.main(ConsoleConsumer.scala)
解决方案

空指针的核心原因是命令行指定的反序列化器与配置文件中的反序列化器冲突,且AWSKafkaAvroDeserializer未被正确初始化配置参数。以下是修复步骤:

  • 统一反序列化器配置:删除消费者命令行中的--key-deserializer和--value-deserializer参数,直接使用client.properties中配置的GlueSchemaRegistryKafkaDeserializer,该类会自动处理Glue Schema Registry的初始化逻辑,避免手动指定导致的配置缺失。修改后的命令:

    kafka-avro-console-consumer --bootstrap-server <bootstrap_server_url> \
    --consumer.config client.properties \
    --property print.key=true \
    --property print.value=true \
    --topic platform_avro_users --from-beginning
    
  • 移除冲突参数:删除命令行中的schema.registry.url=https://glue.us-east-1.amazonaws.com,该参数属于Confluent Schema Registry,Glue Schema Registry的反序列化器不需要此配置。

  • 验证JAR包位置:确保AWS Glue Schema Registry的所有相关JAR包(含依赖包)都放入Confluent的share/java/kafka-serde-tools目录,或通过CLASSPATH环境变量显式指定JAR路径,避免类加载失败。

  • 检查配置一致性:确认client.properties中的schemaNameGenerationClass值与生产者端使用的类完全一致,避免Schema匹配失败。

内容的提问来源于stack exchange,提问作者Pragya Rai

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 02:07:07