Quarkus对接Confluent Avro Schema失败问题求助
Quarkus对接Kafka使用Confluent Avro Schema时的序列化错误解决思路
问题描述
在Quarkus对接Kafka时,尝试使用Confluent Avro Schema Registry中已创建的名为default.avro的Avro Schema,但出现序列化异常,无法正常发送消息到Kafka。
配置详情
application.properties中的Kafka输出通道配置片段:
mp.messaging.outgoing.logistics-events-example-out.topic=logistics.events.eta.example.test mp.messaging.outgoing.logistics-events-example-out.transforms=SetSchemaMetadata mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.type=org.apache.kafka.connect.transforms.SetSchemaMetadata$Value mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.schema.name=default.avro mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.schema.version=1
依赖信息
pom.xml中相关依赖配置:
<!-- AVRO --> <!-- the extension --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-confluent-registry-avro</artifactId> </dependency> <!-- Confluent registry libraries use Jakarta REST client --> <dependency> <groupId>io.quarkus</groupId> <artifactId>quarkus-rest-client-reactive</artifactId> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-avro-serializer</artifactId> <version>7.3.0</version> <exclusions> <exclusion> <groupId>jakarta.ws.rs</groupId> <artifactId>jakarta.ws.rs-api</artifactId> </exclusion> </exclusions> </dependency> <dependency> <groupId>io.confluent</groupId> <artifactId>kafka-schema-registry-client</artifactId> <version>7.3.0</version> </dependency>
错误日志
sampled=false] io.smallrye.reactive.messaging.kafka.impl.KafkaSink lambda$writeMessageToKafka$4 SRMSG18212: Message org.eclipse.microprofile.reactive.messaging.Message$5@32c7bed1 from channel 'test-events-example-out' was not sent to Kafka topic 'logistics.events.eta.example.test' - nacking message: org.apache.kafka.common.errors.SerializationException: Error retrieving Avro schema{"type":"record","name":"ExampleChangedEvent","namespace":"Example","fields":[{"name":"countryCode","type":{"type":"string","avro.java.string":"String"}},{"name":"occurredOn","type":{"type":"string","avro.java.string":"String"}}]} at io.confluent.kafka.serializers.AbstractKafkaSchemaSerDe.toKafkaException(AbstractKafkaSerDe.java:253) at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:168) at io.confluent.kafka.serializers.KafkaAvroSerializer.serialize(KafkaAvroSerializer.java:61) at org.apache.kafka.common.serialization.Serializer.serialize(Serializer.java:62) at io.smallrye.reactive.messaging.kafka.fault.SerializerWrapper.lambda$serialize$1(SerializerWrapper.java:56) at io.smallrye.reactive.messaging.kafka.fault.SerializerWrapper.wrapSerialize(SerializerWrapper.java:81) at io.smallrye.reactive.messaging.kafka.fault.SerializerWrapper.serialize(SerializerWrapper.java:56) at org.apache.kafka.clients.producer.KafkaProducer.doSend(KafkaProducer.java:1000) at org.apache.kafka.clients.producer.KafkaProducer.send(KafkaProducer.java:947) at io.smallrye.reactive.messaging.kafka.impl.ReactiveKafkaProducer.lambda$send$4(ReactiveKafkaProducer.java:164) at io.smallrye.context.impl.wrappers.SlowContextualConsumer.accept(SlowContextualConsumer.java:21) at io.smallrye.mutiny.operators.uni.builders.UniCreateWithEmitter.subscribe(UniCreateWithEmitter.java:22) at io.smallrye.mutiny.operators.AbstractUni.subscribe(AbstractUni.java:36) at io.smallrye.mutiny.operators.uni.UniOnItemTransformToUni$UniOnItemTransformToUniProcessor.performInnerSubscription(UniOnItemTransformToUni.java:81) at io.smallrye.mutiny.operators.uni.UniOnItemTransformToUni$UniOnItemTransformToUniProcessor.onItem(UniOnItemTransformToUni.java:57) at io.smallrye.mutiny.operators.uni.UniOperatorProcessor.onItem(UniOperatorProcessor.java:47) at io.smallrye.mutiny.operators.uni.UniMemoizeOp.forwardTo(UniMemoizeOp.java:123) at io.smallrye.mutiny.operators.uni.UniMemoizeOp.subscribe(UniMemoizeOp.java:67) at io.smallrye.mutiny.operators.AbstractUni.subscribe(AbstractUni.java:36) at io.smallrye.mutiny.operators.uni.UniRunSubscribeOn.lambda$subscribe$0(UniRunSubscribeOn.java:27) at java.base@21.0.2/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1144) at java.base@21.0.2/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:642) at java.base@21.0.2/java.lang.Thread.runWith(Thread.java:1596) at java.base@21.0.2/java.lang.Thread.run(Thread.java:1583) at org.graalvm.nativeimage.builder/com.oracle.svm.core.thread.PlatformThreads.threadStartRoutine(PlatformThreads.java:833) at org.graalvm.nativeimage.builder/com.oracle.svm.core.posix.thread.PosixPlatformThreads.pthreadStartRoutine(PosixPlatformThreads.java:211) Caused by: io.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Subject 'example.name-value' not found.; error code: 40401 at io.confluent.kafka.schemaregistry.client.rest.RestService.sendHttpRequest(RestService.java:302) at io.confluent.kafka.schemaregistry.client.rest.RestService.httpRequest(RestService.java:372) at io.confluent.kafka.schemaregistry.client.rest.RestService.lookUpSubjectVersion(RestService.java:463) at io.confluent.kafka.schemaregistry.client.rest.RestService.lookUpSubjectVersion(RestService.java:448) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getIdFromRegistry(CachedSchemaRegistryClient.java:350) at io.confluent.kafka.schemaregistry.client.CachedSchemaRegistryClient.getId(CachedSchemaRegistryClient.java:565) at io.confluent.kafka.serializers.AbstractKafkaAvroSerializer.serializeImpl(AbstractKafkaAvroSerializer.java:137) ... 24 more
关键错误分析:日志明确提示Subject 'example.name-value' not found,说明序列化器尝试查找的schema subject与Registry中创建的default.avro不匹配,而是自动使用了Avro对象的namespace.name-value作为subject名称。
解决思路
1. 替换为Quarkus原生配置指定schema subject
Quarkus的SmallRye Reactive Messaging不支持Kafka Connect的SetSchemaMetadata转换,该配置无法生效。需改用Quarkus Confluent Registry Avro扩展的原生配置:
# 指定value序列化器为Avro mp.messaging.outgoing.logistics-events-example-out.value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer # 强制指定要使用的schema subject名称 mp.messaging.outgoing.logistics-events-example-out.avro.schema.subject=default.avro # 配置Schema Registry地址(若未配置过) mp.messaging.outgoing.logistics-events-example-out.schema.registry.url=http://你的registry地址:8081
2. 确认Registry中的schema存在性
通过Registry API验证default.avro subject及版本1是否存在:
curl http://你的registry地址:8081/subjects/default.avro/versions/1
若返回404,需核对Registry中的实际subject名称和版本号。
3. 移除无效的Kafka Connect转换配置
删除原配置中无效的SetSchemaMetadata相关项:
mp.messaging.outgoing.logistics-events-example-out.transforms=SetSchemaMetadata mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.type=org.apache.kafka.connect.transforms.SetSchemaMetadata$Value mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.schema.name=default.avro mp.messaging.outgoing.logistics-events-example-out.transforms.SetSchemaMetadata.schema.version=1
4. 确保消息对象与schema定义一致
检查发送的ExampleChangedEvent对象结构,字段名、类型、命名空间需与Registry中default.avro的定义完全匹配,避免序列化器自动注册新schema导致subject不匹配。
内容的提问来源于stack exchange,提问作者tom vercammen
相关产品推荐
相关产品推荐

