Azure Event Hub Avro Schema Registry Java反序列化异常求助
问题:Azure Schema Registry Avro反序列化抛出IllegalStateException异常
我用Gradle脚本从.avsc文件生成了Avro的Order类,相关文件和代码如下:
Avro Schema文件(src/main/resources/avro/Order.avsc)
{ "namespace": "com.azure.schemaregistry.samples", "type": "record", "name": "Order", "fields": [ { "name": "id", "type": "string" }, { "name": "amount", "type": "double" } ]}
Gradle构建脚本
import org.apache.avro.tool.SpecificCompilerTool buildscript { dependencies { // Add the Avro code generation to the build dependencies so that it can be used in a Gradle task. classpath group: 'org.apache.avro', name: 'avro-tools', version: '1.11.1' } } plugins { id 'java' id 'org.springframework.boot' version '3.1.8' id 'io.spring.dependency-management' version '1.1.4' } def avroSchemasDir = "src/main/resources/avro" def avroCodeGenerationDir = "build/avro" group = 'com.eventhub' version = '0.0.1-SNAPSHOT' java { sourceCompatibility = '17' } repositories { mavenCentral() } dependencies { implementation 'org.springframework.boot:spring-boot-starter-web' implementation 'org.apache.kafka:kafka-streams' implementation 'org.springframework.kafka:spring-kafka' compileOnly 'org.projectlombok:lombok' annotationProcessor 'org.projectlombok:lombok' testImplementation 'org.springframework.boot:spring-boot-starter-test' testImplementation 'org.springframework.kafka:spring-kafka-test' implementation group: 'com.microsoft.azure', name: 'msal4j', version: '1.14.2' implementation group: 'org.apache.commons', name: 'commons-lang3', version: '3.12.0' implementation group: 'com.azure', name: 'azure-data-schemaregistry', version: '1.4.2' implementation group: 'com.azure', name: 'azure-identity', version: '1.10.4' implementation group: 'com.azure', name: 'azure-data-schemaregistry-avro', version: '1.0.0-beta.5' implementation "org.apache.avro:avro:1.11.0" } tasks.named('bootBuildImage') { builder = 'paketobuildpacks/builder-jammy-base:latest' } tasks.register('avroCodeGeneration') { // Define the task inputs and outputs for the Gradle up-to-date checks. inputs.dir(avroSchemasDir) outputs.dir(avroCodeGenerationDir) // The Avro code generation logs to the standard streams. Redirect the standard streams to the Gradle log. logging.captureStandardOutput(LogLevel.INFO); logging.captureStandardError(LogLevel.ERROR) doLast { // Run the Avro code generation. new SpecificCompilerTool().run(System.in, System.out, System.err, List.of( "-encoding", "UTF-8", "-string", "-fieldVisibility", "private", "-noSetters", "schema", "$projectDir/$avroSchemasDir".toString(), "$projectDir/$avroCodeGenerationDir".toString() )) } } tasks.withType(JavaCompile).configureEach { // Make Java compilation tasks depend on the Avro code generation task. dependsOn('avroCodeGeneration') }
Java反序列化代码
public class Main { public static void main(String[] args) throws IOException { Order order = deserialize(); System.out.println(order.getAmount()); } public static Order deserialize() throws IOException { TokenCredential tokenCredential = new ClientSecretCredentialBuilder() .tenantId("tenant") .clientId("id") .clientSecret("secret") .build(); SchemaRegistryAsyncClient schemaRegistryAsyncClient = new SchemaRegistryClientBuilder() .fullyQualifiedNamespace("https://test.servicebus.windows.net") .credential(tokenCredential) .buildAsyncClient(); SchemaRegistryAvroSerializer schemaRegistryAvroSerializer = new SchemaRegistryAvroSerializerBuilder() .schemaRegistryAsyncClient(schemaRegistryAsyncClient) .schemaGroup("test-group") .avroSpecificReader(true) .autoRegisterSchema(true) .buildSerializer(); InputStream inputStream = getDataToDeserialize(); return schemaRegistryAvroSerializer.deserialize(inputStream, TypeReference.createInstance(Order.class)); } static InputStream getDataToDeserialize() throws IOException { Order order = new Order("10", 100.0); System.out.println("OUTPUT : " + order.getSchema().toString()); return new ByteArrayInputStream(order.toByteBuffer().array()); } }
执行时抛出异常:
Exception in thread "main" java.lang.IllegalStateException: Illegal format: unsupported record format indicator in payload at com.azure.data.schemaregistry.avro.SchemaRegistryAvroSerializer.lambda$deserializeAsync$2(SchemaRegistryAvroSerializer.java:79) at reactor.core.publisher.FluxFlatMap.trySubscribeScalarMap(FluxFlatMap.java:152) at reactor.core.publisher.MonoFlatMap.subscribeOrReturn(MonoFlatMap.java:53) at reactor.core.publisher.Mono.subscribe(Mono.java:4480) at reactor.core.publisher.Mono.block(Mono.java:1711) at com.azure.data.schemaregistry.avro.SchemaRegistryAvroSerializer.deserialize(SchemaRegistryAvroSerializer.java:50)
参考官方示例编写但仍报错,求排查问题。
问题原因与解决方案
核心原因
SchemaRegistryAvroSerializer要求的序列化 payload 必须包含Azure Schema Registry的格式前缀(用来标识schema的ID等元数据),但当前getDataToDeserialize方法直接用Avro原生的toByteBuffer()生成了纯Avro二进制数据,没有添加这个前缀,导致序列化器无法识别格式,抛出异常。
解决方案
方案1:用同一个SchemaRegistryAvroSerializer序列化数据
修改getDataToDeserialize方法,通过SchemaRegistryAvroSerializer来序列化对象,而不是用Avro原生方法:
static InputStream getDataToDeserialize(SchemaRegistryAvroSerializer serializer) throws IOException { Order order = new Order("10", 100.0); System.out.println("OUTPUT : " + order.getSchema().toString()); // 使用SchemaRegistryAvroSerializer序列化,生成带前缀的payload byte[] serializedBytes = serializer.serialize(order, TypeReference.createInstance(Order.class)); return new ByteArrayInputStream(serializedBytes); }
然后在deserialize方法中调用时传入serializer:
InputStream inputStream = getDataToDeserialize(schemaRegistryAvroSerializer);
方案2:手动添加格式前缀(不推荐,仅用于理解原理)
Azure Schema Registry的Avro payload格式为:
- 1字节的格式指示器(固定为
0x00,表示使用schema ID) - 4字节的Big-Endian整数(schema ID的长度)
- N字节的UTF-8编码的schema ID字符串
- 剩余部分为纯Avro二进制数据
但手动拼接容易出错,推荐使用方案1。
额外注意点
- 检查
azure-data-schemaregistry-avro版本:当前使用的是1.0.0-beta.5,建议升级到稳定版(比如1.4.2,和azure-data-schemaregistry版本保持一致),beta版本可能存在已知问题。 - 确认Schema Registry的fullyQualifiedNamespace格式:正确格式应为
<namespace>.servicebus.windows.net,不需要https://前缀,当前的"https://test.servicebus.windows.net"需要去掉https://。
内容的提问来源于stack exchange,提问作者Prashant Aghara
相关产品推荐
相关产品推荐

