Delta Table流数据至Kafka Topic实现及40403报错排查
互联网上多为从Kafka Topic流数据到Delta Table的示例,而我需将Delta Table数据流式传输至Kafka Topic,是否可行?以下是我尝试的代码:
val schemaRegistryAddr = "https://..." val avroSchema = buildSchema(topic) //defined this method val Df = spark.readStream.format("delta").load("path..") .withColumn("key", col("lskey").cast(StringType)) .withColumn("topLevelRecord",struct(col("col1"),col("col2")...) .select( to_avro($"key", lit("topic-key"), schemaRegistryAddr).as("key"), to_avro($"topLevelRecord", lit("topic-value"), schemaRegistryAddr, avroSchema).as("value")) Df.writeStream .format("kafka") .option("checkpointLocation",checkpointPath) .option("kafka.bootstrap.servers", bootstrapServers) .option("kafka.security.protocol", "SSL") .option("kafka.ssl.keystore.location", kafkaKeystoreLocation) .option("kafka.ssl.keystore.password", keystorePassword) .option("kafka.ssl.truststore.location", kafkaTruststoreLocation) .option("topic",topic) .option("batch.size",262144) .option("linger.ms",5000) .trigger(ProcessingTime("25 seconds")) .start()
运行时出现错误:org.spark_project.confluent.kafka.schemaregistry.client.rest.exceptions.RestClientException: Schema not found; error code: 40403,但使用批量生产者写入同一Topic可成功,请问流式写入遗漏了什么?另外我发现一篇旧博客称当前Structured Streaming API不支持'kafka'格式。
解答
首先明确:Delta Table流式写入Kafka完全可行,Structured Streaming从Spark 2.3版本起就已经支持Kafka作为流式输出源,你看到的旧博客信息已经过时。
针对你遇到的40403报错,排查方向如下:
Schema Registry的subject名称不匹配:你调用
to_avro时传入的lit("topic-key")和lit("topic-value")是Schema Registry中的subject名称,而批量写入成功说明批量代码用的subject是正确的(通常格式为{topic-name}-key和{topic-name}-value)。检查Schema Registry中是否存在topic-key和topic-value这两个subject,或者改成和批量代码一致的subject名称。流式任务的权限/网络问题:流式任务的执行上下文可能和批量任务不同,比如运行流式任务的账号没有Schema Registry的读权限,或者流式环境无法连通Schema Registry地址。验证流式任务所在节点能否访问Schema Registry,以及所用账号是否有权限读取目标subject的schema。
Schema自动注册配置缺失:批量写入时可能自动完成了schema注册,但流式任务默认可能没有开启自动注册。如果使用的是Confluent的Avro工具类,需要给
to_avro添加auto.register.schemas=true参数,确保流式任务能自动注册schema(如果目标subject尚未存在)。Checkpoint遗留状态干扰:如果之前的流式任务失败过,checkpoint目录中可能留存了旧的schema状态信息,导致重启时用错误的schema去查找。尝试删除checkpoint目录后重新启动任务。
内容的提问来源于stack exchange,提问作者Don Sam

