如何正确配置Spark→Kafka→JDBC Sink Connector与Avro?
问题:Kafka JDBC Sink无法匹配Apicurio Registry中的Avro Schema
环境与代码
Spark生产Kafka消息代码
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{col, struct} import org.apache.spark.sql.avro.functions.to_avro import org.apache.spark.sql.types.{DoubleType, LongType, StructType} object IngestFromS3ToKafka { def main(args: Array[String]): Unit = { val spark: SparkSession = SparkSession .builder() .master("local[*]") .appName("ingest-from-s3-to-kafka") .config("spark.ui.port", "4040") .getOrCreate() val folderPath = "s3a://hongbomiao-bucket/iot/" val parquet_schema = new StructType() .add("timestamp", DoubleType) .add("current", DoubleType, nullable = true) .add("voltage", DoubleType, nullable = true) .add("temperature", DoubleType, nullable = true) val df = spark.readStream .schema(parquet_schema) .option("maxFilesPerTrigger", 1) .parquet(folderPath) .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType)) .select(to_avro(struct("*")).alias("value")) val query = df.writeStream .format("kafka") .option( "kafka.bootstrap.servers", "hm-kafka-kafka-bootstrap.hm-kafka.svc:9092" ) .option("topic", "hm.motor") .option("checkpointLocation", "/tmp/checkpoint") .start() query.awaitTermination() } }
Apicurio Registry创建Schema命令
curl --location 'http://apicurio-registry-apicurio-registry.hm-apicurio-registry.svc:8080/apis/registry/v2/groups/hm-group/artifacts' \ --header 'Content-type: application/json; artifactType=AVRO' \ --header 'X-Registry-ArtifactId: hm-iot' \ --data '{ "type": "record", "namespace": "com.hongbomiao", "name": "hm.motor", "fields": [ { "name": "timestamp", "type": "long" }, { "name": "current", "type": "double" }, { "name": "voltage", "type": "double" }, { "name": "temperature", "type": "double" } ] }'
JDBC Sink Connector配置
{ "name": "hm-motor-jdbc-sink-kafka-connector", "config": { "connector.class": "io.aiven.connect.jdbc.JdbcSinkConnector", "tasks.max": 1, "topics": "hm.motor", "connection.url": "jdbc:postgresql://timescale.hm-timescale.svc:5432/hm_iot_db", "connection.user": "${file:/opt/kafka/external-configuration/hm-iot-db-credentials-volume/iot-db-credentials.properties:timescaledb_user}", "connection.password": "${file:/opt/kafka/external-configuration/hm-iot-db-credentials-volume/iot-db-credentials.properties:timescaledb_password}", "insert.mode": "upsert", "table.name.format": "motor", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://apicurio-registry-apicurio-registry.hm-apicurio-registry.svc:8080/apis/ccompat/v6", "transforms": "convertTimestamp", "transforms.convertTimestamp.type": "org.apache.kafka.connect.transforms.TimestampConverter$Value", "transforms.convertTimestamp.field": "timestamp", "transforms.convertTimestamp.target.type": "Timestamp" } }
问题现象
Kafka Connect日志提示找不到Content ID -1330532454的内容,但实际在Apicurio Registry中创建的Schema Content ID是26。尝试将Schema名称改为hm.motor-value(遵循默认的TopicNameStrategy)后,问题依旧。
解决方案
1. 修正Spark的Avro序列化逻辑
Spark默认的to_avro函数不会自动向Schema Registry注册Schema,而是使用本地生成的Schema CRC32哈希值(即日志中的负数ID)作为标识,而非Registry分配的Content ID。需要修改Spark代码,让其使用Apicurio Registry进行序列化:
- 添加Apicurio Registry依赖:确保Spark应用包含
io.apicurio:apicurio-registry-serdes-avro:2.4.0.Final(版本根据实际环境调整) - 修改
to_avro调用,指定Registry信息:
或者通过SparkSession全局配置Registry:import org.apache.spark.sql.functions.lit // ... 其他代码不变 val df = spark.readStream .schema(parquet_schema) .option("maxFilesPerTrigger", 1) .parquet(folderPath) .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType)) // 指定Registry URL、分组、Artifact ID .select(to_avro( struct("*"), lit("http://apicurio-registry-apicurio-registry.hm-apicurio-registry.svc:8080/apis/registry/v2"), lit("hm-group"), lit("hm.motor-value") ).alias("value"))val spark: SparkSession = SparkSession .builder() .master("local[*]") .appName("ingest-from-s3-to-kafka") .config("spark.ui.port", "4040") .config("spark.sql.avro.registry.url", "http://apicurio-registry-apicurio-registry.hm-apicurio-registry.svc:8080/apis/registry/v2") .getOrCreate()
2. 对齐Schema的字段可空性
Spark定义的parquet_schema中current、voltage、temperature字段为nullable=true,但Apicurio中的Schema将这些字段定义为非可空的double类型,二者不匹配会导致解析失败。更新Apicurio中的Schema:
{ "type": "record", "namespace": "com.hongbomiao", "name": "hm.motor-value", "fields": [ { "name": "timestamp", "type": "long" }, { "name": "current", "type": ["null", "double"], "default": null }, { "name": "voltage", "type": ["null", "double"], "default": null }, { "name": "temperature", "type": ["null", "double"], "default": null } ] }
3. 验证消息中的Schema ID
使用Kafka控制台消费者验证消息携带的Schema ID是否为Registry中的26:
kafka-console-consumer.sh \ --bootstrap-server hm-kafka-kafka-bootstrap.hm-kafka.svc:9092 \ --topic hm.motor \ --from-beginning \ --formatter io.confluent.kafka.formatter.AvroMessageFormatter \ --property schema.registry.url=http://apicurio-registry-apicurio-registry.hm-apicurio-registry.svc:8080/apis/ccompat/v6
如果输出正常且Schema ID匹配,说明序列化逻辑已修正。
4. 确认Connector的Schema策略
确保JDBC Sink配置中,value.converter的Schema策略与Spark端一致:
- 若使用TopicNameStrategy,确保Registry中Artifact ID为
hm.motor-value的Schema存在且正确 - 无需额外配置策略参数,默认即为TopicNameStrategy
内容的提问来源于stack exchange,提问作者Hongbo Miao
相关产品推荐
相关产品推荐

