Kafka Redshift连接器JSON_SR序列化错误排查求助
问题分析与解决方法
从错误日志里的Unknown magic byte!和Error deserializing JSON message for id -1可以明确核心问题:你发送到Kafka的是普通JSON字符串,但Redshift连接器配置的JsonSchemaConverter期望的是带Confluent Schema Registry元数据的序列化消息。Confluent的JSON_SR格式会在消息头部添加magic byte、schema ID等标识,普通JSON没有这些内容,导致转换器无法解析。
解决方法
方法1:按Schema Registry规范发送消息
如果需要保留Schema验证功能,必须用符合Confluent Schema Registry要求的方式发送消息:
- 使用Confluent REST Proxy发送时,需指定正确的请求头并携带Schema关联信息:
- 请求URL:
http://<rest-proxy-host>:8082/topics/stripe-connector-2 - 请求头:
Content-Type: application/vnd.kafka.jsonschema.v2+json - 请求体示例(替换你的Schema ID和业务数据):
{ "records": [ { "value_schema_id": 123, // 替换为你在Schema Registry中创建的主题Schema ID "value": { // 你的业务JSON数据,需严格匹配Schema定义 } } ] }
- 请求URL:
- 若用代码发送,需使用Confluent提供的
JsonSchemaSerializer,并配置Schema Registry地址,让它自动处理消息序列化与Schema关联。
方法2:修改连接器为普通JSON转换器(无需Schema验证)
如果不需要Schema Registry的验证能力,可直接调整Redshift连接器的转换器配置:
- 将连接器的
value.converter改为普通JSON转换器:value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false - 该配置会让连接器直接解析普通JSON字符串,无需依赖Schema Registry元数据。
额外检查项
- 确认Schema Registry地址配置正确:连接器的
value.converter.schema.registry.url需指向你的Schema Registry实例地址,且连接器能正常访问该地址。 - 确认发送的JSON结构与Schema定义完全匹配,避免因字段类型、数量不匹配引发后续错误。
内容的提问来源于stack exchange,提问作者C Ekanayake
相关产品推荐
相关产品推荐

