Scala-Akka程序执行sbt run传入路径无法定位CSV文件问题
问题排查与修复方案
核心问题分析
- 目录有效性校验缺失:
new File(directoryPath).listFiles()在目录不存在、无访问权限时会返回null,直接导致无文件被读取,且无错误提示。 - 未过滤CSV文件:会读取目录下所有文件(包括非CSV),可能引发后续解析异常。
- 文件数统计逻辑错误:用
sensors.size(传感器数量)替代实际处理的文件数,完全不符合需求。 - 统计计算错误:
SensorStats中的min和max被错误用平均值填充,未实际计算湿度的极值。 - 空行/无效值处理缺失:CSV空行会导致数组越界,"NaN"字符串会被
toDoubleOption解析为Double.NaN,被误判为有效数值。
修复后的完整代码
import java.io.File import akka.actor.ActorSystem import akka.stream.ActorMaterializer import akka.stream.scaladsl.{FileIO, Framing, Sink, Source} import akka.util.ByteString import scala.collection.mutable import scala.concurrent.ExecutionContext.Implicits.global object HumiditySensorStatistics { case class HumidityData(sum: Double, count: Int, min: Double, max: Double) { def avg: Option[Double] = if (count > 0) Some(sum / count) else None } case class SensorStats(min: Option[Double], avg: Option[Double], max: Option[Double]) def main(args: Array[String]): Unit = { if (args.isEmpty) { println("请提供目录路径作为参数") System.exit(1) } val directoryPath = args(0) val dir = new File(directoryPath) // 校验目录合法性 if (!dir.exists() || !dir.isDirectory) { println(s"错误:目录 $directoryPath 不存在或不是有效目录") System.exit(1) } implicit val system: ActorSystem = ActorSystem("HumiditySensorStatistics") implicit val materializer: ActorMaterializer = ActorMaterializer() val sensors = mutable.Map[String, HumidityData]() var failedMeasurements = 0 var processedFiles = 0 println(s"正在目录中查找CSV文件:$directoryPath") // 仅读取目录下的CSV文件 val csvFiles = dir.listFiles().filter(f => f.isFile && f.getName.endsWith(".csv")) val fileSource = Source.fromIterator(() => csvFiles.iterator) val measurementSource = fileSource.flatMapConcat { f => processedFiles += 1 // 统计实际处理的文件数 FileIO.fromPath(f.toPath) } .via(Framing.delimiter(ByteString("\n"), maximumFrameLength = 1024, allowTruncation = true)) .map(_.utf8String.trim) .filter(_.nonEmpty) // 过滤空行 .drop(1) // 跳过表头行 .map(line => { val fields = line.split(",").map(_.trim) // 处理字段前后空格 (fields(0), fields(1)) }) val sink = Sink.foreach[(String, String)](data => { val sensorId = data._1 val humidityStr = data._2 // 明确区分"NaN"无效值与合法数值 val humidity = if (humidityStr.equalsIgnoreCase("NaN")) None else humidityStr.toDoubleOption humidity match { case Some(h) if !h.isNaN => sensors.update(sensorId, sensors.getOrElse(sensorId, HumidityData(0.0, 0, Double.MaxValue, Double.MinValue)) match { case HumidityData(sum, count, currMin, currMax) => HumidityData(sum + h, count + 1, math.min(currMin, h), math.max(currMax, h)) }) case _ => failedMeasurements += 1 } }) measurementSource.runWith(sink).onComplete(_ => { val numMeasurementsProcessed = sensors.values.map(_.count).sum val numFailedMeasurements = failedMeasurements println(s"已处理文件数:$processedFiles") println(s"已处理测量数:$numMeasurementsProcessed") println(s"失败测量数:$numFailedMeasurements") val statsBySensor = sensors.map { case (sensorId, humidityData) => val stats = SensorStats( min = if (humidityData.count > 0) Some(humidityData.min) else None, avg = humidityData.avg, max = if (humidityData.count > 0) Some(humidityData.max) else None ) (sensorId, stats) } println("平均湿度最高的传感器:") println("sensor-id,min,avg,max") statsBySensor.toList.sortBy(_._2.avg)(Ordering.Option[Double].reverse).foreach { case (sensorId, stats) => println(s"$sensorId,${stats.min.map(_.toInt).getOrElse("NaN")},${stats.avg.map(_.toInt).getOrElse("NaN")},${stats.max.map(_.toInt).getOrElse("NaN")}") } system.terminate() }) } }
额外注意事项
- 目录路径确认:执行
sbt run时工作目录为项目根目录,若CSV在src/main/scala/data,需用命令sbt "run src/main/scala/data";若将data目录移至项目根目录,可简化为sbt "run data"。 - CSV格式规范:确保所有CSV文件第一行为
sensor-id,humidity格式,字段间尽量避免多余空格(修复后代码已兼容带空格的字段)。 - 无效值处理:修复后的代码明确识别"NaN"字符串为无效测量值,避免统计误差。
内容的提问来源于stack exchange,提问作者Always_a_learner
相关产品推荐
相关产品推荐

