Confluent Cloud SQS连接器:JSON字符串转JSON对象适配Schema Registry验证
解决方案
要解决SQS消息体JSON字符串无法通过Schema Registry验证的问题,核心是在连接器拉取数据后,将字符串格式的Body解析为结构化JSON对象,再配合JSON Schema转换器完成Schema验证与JSON_SR格式输出。以下是具体实现步骤与配置:
关键配置思路
- 提取SQS消息Body:使用官方SMT将SQS消息中的
Body字段提取为Kafka记录的值。 - 解析JSON字符串为结构化对象:利用Kafka Connect自带的
JacksonTransformSMT,将提取到的JSON字符串解析为可被Schema Registry识别的JSON结构。 - 配置JSON Schema转换器:指定使用
JsonSchemaConverter,关联Schema Registry地址与认证信息,开启Schema验证,最终输出JSON_SR格式的消息。
完整连接器配置示例
name=sqs-source-connector connector.class=io.confluent.connect.sqs.SqsSourceConnector tasks.max=1 kafka.topic=your-target-topic sqs.queue.url=https://sqs.your-region.amazonaws.com/123456789012/your-queue-name aws.access.key.id=your-aws-access-key aws.secret.access.key=your-aws-secret-key # 转换步骤:提取Body字段 → 解析JSON字符串 transforms=extractBody,parseJson transforms.extractBody.type=org.apache.kafka.connect.transforms.ExtractField$Value transforms.extractBody.field=Body transforms.parseJson.type=org.apache.kafka.connect.transforms.JacksonTransform$Value transforms.parseJson.spec=parseJson(value) # JSON Schema转换器配置(对应JSON_SR格式) value.converter=io.confluent.connect.json.JsonSchemaConverter value.converter.schema.registry.url=https://your-schema-registry-url value.converter.basic.auth.credentials.source=USER_INFO value.converter.basic.auth.user.info=your-schema-registry-api-key:your-schema-registry-api-secret value.converter.validate.schema=true
注意事项
- 确保Ruby on Rails应用发送到SQS的JSON字符串结构,与Schema Registry中已注册的Schema完全兼容(字段名、类型、嵌套结构一致),否则仍会触发Schema不兼容报错。
JacksonTransform是Apache Kafka Connect原生SMT,无需额外安装插件,可直接在Confluent Cloud中使用。- 若Schema Registry开启了严格兼容性模式,需保证新消息的Schema与主题下现有Schema兼容(如仅添加可选字段、不修改已有字段类型)。
内容的提问来源于stack exchange,提问作者Ignotus
相关产品推荐
相关产品推荐

