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

Spark Structured Streaming多查询报错:Race while writing batch 0 求助

问题描述

我需要从多个Kafka Topic读取数据,完成聚合后写入HDFS。通过遍历Kafka Topic列表创建多个流查询时,单个查询运行正常,但同时运行多个查询时出现报错。已经为每个Topic设置了独立的checkpoint目录(参考多篇文章避免同类问题),但问题仍未解决。

报错信息

java.lang.IllegalStateException: Race while writing batch 0

代码片段

object CombinedDcAggStreaming {
 def main(args: Array[String]): Unit = {
 val jobConfigFile = "configPath"
 /* Read input configuration */
 val jobProps = Util.loadProperties(jobConfigFile).asScala
 val sparkConfigFile = jobProps.getOrElse("spark_config_file", throw new RuntimeException("Can't find spark property file"))
 val kafkaConfigFile = jobProps.getOrElse("kafka_config_file", throw new RuntimeException("Can't find kafka property file"))
 val sparkProps = Util.loadProperties(sparkConfigFile).asScala
 val kafkaProps = Util.loadProperties(kafkaConfigFile).asScala
 val topicList = Seq("topic_1", "topic_2")
 val avroSchemaFile = jobProps.getOrElse("schema_file", throw new RuntimeException("Can't find schema file..."))
 val checkpointLocation = jobProps.getOrElse("checkpoint_location", throw new RuntimeException("Can't find check point directory..."))
 val triggerInterval = jobProps.getOrElse("triggerInterval", throw new RuntimeException("Can't find trigger interval..."))
 val outputPath = jobProps.getOrElse("output_path", throw new RuntimeException("Can't find output directory..."))
 val outputFormat = jobProps.getOrElse("output_format", throw new RuntimeException("Can't find output format...")) //"parquet"
 val outputMode = jobProps.getOrElse("output_mode", throw new RuntimeException("Can't find output mode...")) //"append"
 val partitionByCols = jobProps.getOrElse("partition_by_columns", throw new RuntimeException("Can't find partition by columns...")).split(",").toSeq
 val spark = SparkSession.builder.appName("streaming").master("local[4]").getOrCreate()
 sparkProps.foreach(prop => spark.conf.set(prop._1, prop._2))
 topicList.foreach( topicId => {
 kafkaProps.update("subscribe", topicId)
 val schemaPath = avroSchemaFile + "/" + topicId + ".avsc"
 val dimensionMap = ConfigUtils.getDimensionMap(jobConfig)
 val measureMap = ConfigUtils.getMeasureMap(jobConfig)
 val source= Source.fromInputStream(Util.getInputStream(schemaPath)).getLines.mkString
 val schemaParser = new Schema.Parser
 val schema = schemaParser.parse(source)
 val sqlTypeSchema = SchemaConverters.toSqlType(schema).dataType.asInstanceOf[StructType]
 val kafkaStreamData = spark
 .readStream
 .format("kafka")
 .options(kafkaProps)
 .load()
 val udfDeserialize = udf(deserialize(source), DataTypes.createStructType(sqlTypeSchema.fields))
 val transformedDeserializedData = kafkaStreamData.select("value").as(Encoders.BINARY)
 .withColumn("rows", udfDeserialize(col("value")))
 .select("rows.*")
 .withColumn("end_time", (col("end_time") / 1000).cast(LongType))
 .withColumn("timestamp", from_unixtime(col("end_time"),"yyyy-MM-dd HH").cast(TimestampType))
 .withColumn("year", from_unixtime(col("end_time"),"yyyy").cast(IntegerType))
 .withColumn("month", from_unixtime(col("end_time"),"MM").cast(IntegerType))
 .withColumn("day", from_unixtime(col("end_time"),"dd").cast(IntegerType))
 .withColumn("hour",from_unixtime(col("end_time"),"HH").cast(IntegerType))
 .withColumn("topic_id", lit(topicId))
 val groupBycols: Array[String] = dimensionMap.keys.toArray[String] ++ partitionByCols.toArray[String]
 val aggregatedData = AggregationUtils.aggregateDFWithWatermarking(transformedDeserializedData, groupBycols, "timestamp", "10 minutes", measureMap) //Watermarking time -> 10. minutes, window => window("timestamp", "5 minutes")
 val query = aggregatedData
 .writeStream
 .trigger(Trigger.ProcessingTime(triggerInterval))
 .outputMode("update")
 .format("console")
 .partitionBy(partitionByCols: _*)
 .option("path", outputPath)
 .option("checkpointLocation", checkpointLocation + "//" + topicId)
 .start()
 })
 spark.streams.awaitAnyTermination()
 def deserialize(source: String): Array[Byte] => Option[Row] = (data: Array[Byte]) => {
 try {
 val parser = new Schema.Parser
 val schema = parser.parse(source)
 val recordInjection: Injection[GenericRecord, Array[Byte]] = GenericAvroCodecs.toBinary(schema)
 val record = recordInjection.invert(data).get
 val objectArray = new Array[Any](record.asInstanceOf[GenericRecord].getSchema.getFields.size)
 record.getSchema.getFields.asScala.foreach(field => {
 val fieldVal = record.get(field.pos()) match {
 case x: org.apache.avro.util.Utf8 => x.toString
 case y: Any => y
 case _ => None
 }
 objectArray(field.pos()) = fieldVal
 })
 Some(Row(objectArray: _*))
 } catch {
 case ex: Exception => {
 log.info(s"Failed to parse schema with error: ${ex.printStackTrace()}")
 None
 }
 }
 }
 }
 }

解决建议

这个Race while writing batch 0错误大概率是多个流查询写入同一输出路径时的资源竞争导致的——你现在所有查询共用同一个outputPath,哪怕checkpoint独立,写入HDFS时也会因为同时操作同一目录下的分区文件引发冲突。下面是具体的修复方案:

1. 为每个Topic分配独立的输出路径

既然给每个Topic设了独立的checkpoint,输出路径也要做同样的区分,避免文件写入冲突:

.option("path", outputPath + "/" + topicId)

这样每个查询的结果会写入outputPath/topic_1、outputPath/topic_2这类独立目录,从根源上解决竞争问题。

2. 修复Kafka配置的复用问题

你在遍历Topic时直接修改了同一个kafkaProps对象:kafkaProps.update("subscribe", topicId)。由于Scala的可变Map(或Java Properties)是共享状态,这会导致多个查询的Kafka配置互相干扰。

建议每次循环时创建配置副本,保证每个查询的配置独立:

val topicKafkaProps = new mutable.HashMap[String, String]() ++ kafkaProps
topicKafkaProps.update("subscribe", topicId)
// 使用独立的配置创建流
val kafkaStreamData = spark
 .readStream
 .format("kafka")
 .options(topicKafkaProps)
 .load()

3. 优化UDF中的Schema解析逻辑

你的deserialize方法每次调用都会重新解析Avro Schema,这不仅浪费性能,还可能引发潜在的并发问题。可以把Schema解析逻辑提到UDF外部,每个Topic只解析一次:

// 循环内解析好Schema后
val schema = schemaParser.parse(source)
// 直接传入解析后的Schema给UDF
val udfDeserialize = udf(deserialize(schema), sqlTypeSchema)

// 修改deserialize方法参数
def deserialize(schema: Schema): Array[Byte] => Option[Row] = (data: Array[Byte]) => {
 try {
 val recordInjection: Injection[GenericRecord, Array[Byte]] = GenericAvroCodecs.toBinary(schema)
 // 剩余逻辑不变
 } catch {
 // 异常处理不变
 }
}

4. 检查输出模式与分区的兼容性

你当前使用outputMode("update")配合partitionBy,虽然Spark支持这种组合,但多查询写入同一分区目录时更容易触发冲突。确保每个查询的分区路径完全独立,或者如果业务允许,可以考虑将多个Topic的流合并为一个流处理,减少独立查询的数量。

优先尝试第一个方案(独立输出路径),这是解决当前报错最直接的办法,之后再逐步排查其他潜在问题。

内容的提问来源于stack exchange,提问作者Siddharth Goel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 10:07:52