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

如何为Kafka Connect及Spark注册使用AVRO Schema?JDBC连接器实操咨询

注册AVRO Schema并在JDBC连接器、Spark中使用的完整指南

我来帮你一步步解决注册AVRO Schema、配置JDBC源连接器以及在Spark中读取数据的问题,结合你给出的Schema和连接器信息详细说明:

一、先将你的AVRO Schema注册到Schema Registry

Confluent的JDBC连接器依赖Schema Registry来管理AVRO Schema,所以第一步要把你的自定义Schema上传到Registry中。

假设你的Schema Registry地址是http://localhost:8081(根据实际部署调整),可以用curl命令完成注册:

curl -X POST -H "Content-Type: application/vnd.schemaregistry.v1+json" \
  --data '{
    "schema": "{ \"type\": \"record\", \"name\": \"myrecord\", \"fields\": [ { \"name\": \"int1\", \"type\": \"int\" }, { \"name\": \"str1\", \"type\": \"string\" }, { \"name\": \"str2\", \"type\": \"string\" } ] }"
  }' \
  http://localhost:8081/subjects/mssql-source-value/versions

这里的mssql-source-value是subject名称,通常遵循<连接器输出的topic名>-value的命名规则,你可以根据实际的Kafka topic名称调整。执行成功后,会返回包含id的响应(比如{"id":1}),这个ID后续会用到。

二、配置JDBC源连接器使用已注册的AVRO Schema

现在修改你的JDBC源连接器配置,添加AVRO转换和Schema Registry相关参数,确保连接器输出的数据严格匹配你注册的Schema:

{
  "name": "mssql-source",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
    "connection.url": "jdbc:sqlserver://<你的MSSQL地址>:1433;databaseName=<数据库名>",
    "connection.user": "<用户名>",
    "connection.password": "<密码>",
    "table.whitelist": "<要同步的表名>",
    "mode": "incrementing",
    "incrementing.column.name": "<自增列名>",
    "topic.prefix": "mssql-",
    // 以下是AVRO核心配置
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "http://localhost:8081",
    // 指定使用我们注册的Schema,而非自动推断
    "value.converter.use.latest.version": false,
    "value.converter.schema.id": 1, // 替换成你刚才注册得到的Schema ID
    // 也可以用subject关联(二选一)
    // "value.converter.schema.subject.name.strategy": "io.confluent.kafka.serializers.subject.TopicNameStrategy"
  }
}

关键说明:

  • 替换<>里的MSSQL连接信息、表名等为你的实际参数
  • value.converter必须设置为Confluent的Avro转换器,不能用默认的JSON转换器
  • 若用Schema ID关联,要保证ID和注册时一致;若用subject,需确保subject名称和注册时完全匹配

三、在Spark中读取AVRO格式的数据

根据数据存储位置,分两种场景说明:

场景1:从Kafka读取AVRO数据

首先要引入对应依赖(以Spark Shell为例):

spark-shell --packages org.apache.spark:spark-sql-kafka-0-10_2.12:3.3.0,io.confluent:kafka-avro-serializer:7.3.0

然后编写代码读取并解析:

import org.apache.spark.sql.avro.functions.from_avro
import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder()
  .appName("ReadAVROFromKafka")
  .getOrCreate()

// 配置Schema Registry地址
val schemaRegistryUrl = "http://localhost:8081"
// 从Schema Registry拉取已注册的Schema(也可以直接硬编码你的Schema字符串)
val avroSchema = spark.read.format("avro")
  .option("schemaRegistryUrl", schemaRegistryUrl)
  .option("subject", "mssql-source-value")
  .load()
  .schema.json

// 读取Kafka数据并解析AVRO结构
val df = spark.readStream
  .format("kafka")
  .option("kafka.bootstrap.servers", "<Kafka地址>:9092")
  .option("subscribe", "mssql-<你的表名>")
  .load()
  .select(from_avro($"value", avroSchema, schemaRegistryUrl).as("data"))
  .select("data.*")

df.show()

场景2:直接读取AVRO文件

如果数据已保存为AVRO文件,直接用Spark的AVRO数据源即可:

val df = spark.read
  .format("avro")
  .option("avroSchema", """{ "type": "record", "name": "myrecord", "fields": [ { "name": "int1", "type": "int" }, { "name": "str1", "type": "string" }, { "name": "str2", "type": "string" } ] }""")
  .load("<AVRO文件路径>")

df.show()

常见问题提醒

  • 确保Schema Registry、Kafka、JDBC连接器的版本兼容(Confluent组件尽量使用同一大版本)
  • 如果连接器自动推断的Schema和你注册的不一致,会触发Schema兼容性错误,此时要确保自定义Schema和数据库表结构完全匹配
  • Spark中使用的AVRO依赖版本要和Spark版本对应

内容的提问来源于stack exchange,提问作者Mahmoud Elbably

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.22 09:19:44