Spark使用Confluent Schema Registry客户端遇"schemaType"未识别字段错误
Confluent Schema Registry 客户端读取JSON Schema时
schemaType字段无法识别问题解决方案 错误根因
你遇到的反序列化报错是因为使用的confluent-kafka-schemaregistry-client版本过低:
- 低版本的
io.confluent.kafka.schemaregistry.client.rest.entities.Schema实体类没有定义schemaType字段,也没有配置Jackson忽略未知属性,所以当注册中心返回包含schemaType的响应时,Jackson在反序列化阶段就会直接抛出异常,还没到你手动提取schema字段的步骤。
可选解决方案
方案1:升级客户端版本(最推荐)
直接升级confluent-kafka-schemaregistry-client到5.5.0及以上版本,该版本开始的Schema实体类已经原生支持schemaType字段,不需要修改任何代码即可正常运行,原有逻辑可以直接拿到schema内容。
方案2:手动发送HTTP请求提取schema字段(不升级依赖可选)
如果不方便升级依赖,可以绕过官方RestService,自己调用注册中心HTTP接口,手动解析返回结果仅提取schema字段,Scala示例代码如下:
import com.fasterxml.jackson.databind.ObjectMapper import com.fasterxml.jackson.module.scala.DefaultScalaModule import java.net.HttpURLConnection import java.net.URL val mapper = new ObjectMapper().registerModule(DefaultScalaModule) val url = new URL("http://localhost:8081/subjects/jsontesting1/versions/latest") val conn = url.openConnection().asInstanceOf[HttpURLConnection] conn.setRequestMethod("GET") val responseJson = mapper.readTree(conn.getInputStream) val jsonSchema = responseJson.get("schema").asText() conn.disconnect()
方案3:配置Jackson全局忽略未知属性(临时修复)
如果必须使用原有RestService,可以在调用接口前通过反射修改RestService内部ObjectMapper的配置,开启忽略未知属性,代码如下:
import io.confluent.kafka.schemaregistry.client.rest.RestService import com.fasterxml.jackson.databind.DeserializationFeature import java.lang.reflect.Field val restService = new RestService("http://localhost:8081") // 反射获取RestService内部的ObjectMapper实例 val mapperField: Field = classOf[RestService].getDeclaredField("objectMapper") mapperField.setAccessible(true) val objectMapper = mapperField.get(restService).asInstanceOf[com.fasterxml.jackson.databind.ObjectMapper] objectMapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false) // 后续调用就不会报错了 val valueRestResponseSchema = restService.getLatestVersion("jsontesting1") val jsonSchema = valueRestResponseSchema.getSchema
内容的提问来源于stack exchange,提问作者mythic
相关产品推荐
相关产品推荐

