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

Apicurio Registry反序列化报artifactId cannot be null错误及Kafka头格式疑问

问题:移除Apicurio Registry fallback配置后Kafka连接器报错"artifactId cannot be null"

Apicurio Registry中的Avro Schema

Artifact ID为hm.motor-value的Schema可通过以下命令获取:

curl --location 'http://apicurio-registry.svc/apis/registry/v2/groups/default/artifacts/hm.motor-value'

返回的Schema内容:

{
    "type": "record",
    "namespace": "com.hongbomiao",
    "name": "motor",
    "fields": [
        {
            "name": "timestamp",
            "type": "long"
        },
        {
            "name": "current",
            "type": "double"
        },
        {
            "name": "voltage",
            "type": "double"
        },
        {
            "name": "temperature",
            "type": "double"
        }
    ]
}

Spark生成Avro Kafka记录的代码

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{col, struct}
import org.apache.spark.sql.types.{DoubleType, LongType, StructType}
import org.apache.spark.sql.avro.functions.to_avro
import sttp.client3.{HttpClientSyncBackend, UriContext, basicRequest}

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 parquetSchema = new StructType()
      .add("timestamp", DoubleType)
      .add("current", DoubleType, nullable = true)
      .add("voltage", DoubleType, nullable = true)
      .add("temperature", DoubleType, nullable = true)

    val backend = HttpClientSyncBackend()
    val res = basicRequest
      .get(
        uri"http://apicurio-registry.svc:8080/apis/registry/v2/groups/default/artifacts/hm.motor-value"
      )
      .send(backend)
    val kafkaRecordValueSchema = res.body.fold(identity, identity)

    val df = spark.readStream
      .schema(parquetSchema)
      .option("maxFilesPerTrigger", 1)
      .parquet(folderPath)
      .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType))
      .select(to_avro(struct("*"), kafkaRecordValueSchema).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()
  }
}

JDBC Sink连接器配置

{
    "name": "hm-motor-jdbc-sink-kafka-connector",
    "connector.class": "io.aiven.connect.jdbc.JdbcSinkConnector",
    "tasks.max": 10,
    "topics": "hm.motor",
    "connection.url": "jdbc:postgresql://timescale.hm-timescale.svc:5432/hm_iot_db",
    "connection.user": "xxx",
    "connection.password": "xxx",
    "insert.mode": "insert",
    "table.name.format": "motor",
    "value.converter": "io.apicurio.registry.utils.converter.AvroConverter",
    "value.converter.apicurio.registry.url": "http://apicurio-registry.svc:8080/apis/registry/v2",
    "value.converter.apicurio.registry.fallback.artifact-id": "hm.motor-value"
}

问题场景

原配置包含value.converter.apicurio.registry.fallback.artifact-id时运行正常,但移除该配置后,连接器抛出如下错误:

org.apache.kafka.connect.errors.ConnectException: Tolerance exceeded in error handler
    at org.apache.kafka.connect.runtime.errors.RetryWithToleranceOperator.execAndHandleError(RetryWithToleranceOperator.java:230)
    ...
Caused by: java.lang.IllegalStateException: artifactId cannot be null
    at io.apicurio.registry.resolver.DefaultSchemaResolver.resolveSchemaByCoordinates(DefaultSchemaResolver.java:177)
    ...

根据官方逻辑,默认的TopicIdStrategy应通过主题名hm.motor自动匹配hm.motor-value的Artifact,但实际未生效。

已尝试的方案(2023年5月5日更新)

尝试手动添加Kafka记录头,但格式不正确,报错依旧:

val df = spark.readStream
  .schema(parquetSchema)
  .option("maxFilesPerTrigger", 1)
  .parquet(folderPath)
  .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType))
  .select(to_avro(struct("*"), kafkaRecordValueSchema).alias("value"))
  .withColumn(
    "headers",
    array(
      struct(
        lit("apicurio.registry.headers.value.groupId.name") as "key",
        lit("default").cast("binary") as "value"
      ),
      struct(
        lit("apicurio.registry.headers.value.artifactId.name") as "key",
        lit("hm.motor-value").cast("binary") as "value"
      )
    )
  )

解决方案

方式1:使用Apicurio官方SerDe序列化(推荐)

放弃Spark原生to_avro,改用Apicurio Registry的Avro SerDe,它会自动添加正确的Schema引用头信息,确保反序列化器能识别:

  1. 引入依赖(Maven示例):
<dependency>
    <groupId>io.apicurio</groupId>
    <artifactId>apicurio-registry-serdes-avro</artifactId>
    <version>2.4.0.Final</version>
</dependency>
  1. 修改Spark代码:
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.col
import org.apache.spark.sql.types.{DoubleType, LongType, StructType}
import io.apicurio.registry.serde.avro.AvroKafkaSerializer
import org.apache.kafka.common.serialization.StringSerializer

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 parquetSchema = 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(parquetSchema)
      .option("maxFilesPerTrigger", 1)
      .parquet(folderPath)
      .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType))

    // 配置Apicurio序列化参数
    val kafkaParams = Map(
      "bootstrap.servers" -> "hm-kafka-kafka-bootstrap.hm-kafka.svc:9092",
      "key.serializer" -> classOf[StringSerializer].getName,
      "value.serializer" -> classOf[AvroKafkaSerializer].getName,
      "apicurio.registry.url" -> "http://apicurio-registry.svc:8080/apis/registry/v2",
      "apicurio.registry.artifact-id" -> "hm.motor-value",
      "apicurio.registry.auto-register" -> "false" // Schema已提前注册,无需自动注册
    )

    // 用foreachBatch结合Apicurio SerDe发送数据
    df.writeStream
      .foreachBatch { (batchDF, _) =>
        import spark.implicits._
        batchDF
          .as[com.hongbomiao.motor] // 需提前用Avro工具生成对应POJO类
          .write
          .format("kafka")
          .options(kafkaParams)
          .option("topic", "hm.motor")
          .save()
      }
      .option("checkpointLocation", "/tmp/checkpoint")
      .start()
      .awaitTermination()
  }
}

注:需通过avro-maven-plugin等工具根据Avro Schema生成对应的POJO类,若不想生成POJO,可改用GenericRecord构建数据。

方式2:手动添加正确的Kafka记录头

若坚持使用Spark原生to_avro,需添加Apicurio反序列化器识别的标准头:

import org.apache.spark.sql.functions.{lit, struct, array}

// 修改df的headers生成逻辑
val df = spark.readStream
  .schema(parquetSchema)
  .option("maxFilesPerTrigger", 1)
  .parquet(folderPath)
  .withColumn("timestamp", (col("timestamp") * 1000).cast(LongType))
  .select(to_avro(struct("*"), kafkaRecordValueSchema).alias("value"))
  .withColumn(
    "headers",
    array(
      struct(
        lit("Content-Type") as "key",
        lit("application/vnd.apache.avro+json; artifactId=hm.motor-value; groupId=default").cast("binary") as "value"
      )
    )
  )

关键原因

你之前添加的头字段并非Apicurio反序列化器的标准识别字段,它只识别Content-Type、apicurio.registry.global-id这类特定头。同时,TopicIdStrategy生效的前提是序列化时已完成主题与Schema的关联注册,原生to_avro未做该关联,导致策略无法自动匹配Artifact。


内容的提问来源于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.22 20:17:02