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

如何正确配置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信息:
    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"))
    
    或者通过SparkSession全局配置Registry:
    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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 05:23:09