使用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

