You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Confluent Cloud SQS连接器:JSON字符串转JSON对象适配Schema Registry验证

解决方案

要解决SQS消息体JSON字符串无法通过Schema Registry验证的问题,核心是在连接器拉取数据后,将字符串格式的Body解析为结构化JSON对象,再配合JSON Schema转换器完成Schema验证与JSON_SR格式输出。以下是具体实现步骤与配置:

关键配置思路

  1. 提取SQS消息Body:使用官方SMT将SQS消息中的Body字段提取为Kafka记录的值。
  2. 解析JSON字符串为结构化对象:利用Kafka Connect自带的JacksonTransform SMT,将提取到的JSON字符串解析为可被Schema Registry识别的JSON结构。
  3. 配置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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 03:57:52