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引用头信息,确保反序列化器能识别:
- 引入依赖(Maven示例):
<dependency> <groupId>io.apicurio</groupId> <artifactId>apicurio-registry-serdes-avro</artifactId> <version>2.4.0.Final</version> </dependency>
- 修改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
相关产品推荐
相关产品推荐

