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

