From c3f85fb477e2f7e488e0a3b01da4b50abcf391af Mon Sep 17 00:00:00 2001 From: Guntis Smaukstelis Date: Thu, 13 Feb 2025 13:15:46 +0200 Subject: [PATCH] Create just one scheduler instance --- src/main/scala/app/Main.scala | 7 +++---- src/main/scala/data/DataService.scala | 2 +- src/main/scala/fetchDMI/FetchService.scala | 1 + src/main/scala/fetchDMI/Scheduler.scala | 8 ++++---- 4 files changed, 9 insertions(+), 9 deletions(-) diff --git a/src/main/scala/app/Main.scala b/src/main/scala/app/Main.scala index f32bbeb..3ef82a6 100644 --- a/src/main/scala/app/Main.scala +++ b/src/main/scala/app/Main.scala @@ -18,12 +18,11 @@ object Main extends IOApp { fileFetchScheduler <- FileFetchScheduler.of(postgresService, fetch) fetchCsvTask = fileFetchScheduler.run.compile.drain - scheduler <- fetchDMI.Scheduler.of("Cleanup", List(1)) - cleanupTask = scheduler.scheduleTask(DataService.deleteOldForecasts()).compile.drain + scheduler <- fetchDMI.Scheduler.of + cleanupTask = scheduler.scheduleTask("Cleanup", List(1), DataService.deleteOldForecasts()).compile.drain - scheduler <- fetchDMI.Scheduler.of("Fetch Grib", List(2)) fetchGrib <- fetchDMI.FetchService.of - fetchGribTask = scheduler.scheduleTask(fetchGrib.fetchRecentForecasts()).compile.drain + fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(2), fetchGrib.fetchRecentForecasts()).compile.drain server <- Server.of(postgresService, fetch) serverTask = server.run diff --git a/src/main/scala/data/DataService.scala b/src/main/scala/data/DataService.scala index 20af4ee..6ef12af 100644 --- a/src/main/scala/data/DataService.scala +++ b/src/main/scala/data/DataService.scala @@ -54,13 +54,13 @@ object DataService { val oldThreshold = nowUTC.minusHours(maxHours) for { + _ <- IO.println("start cleanup") fileList <- getFileList() fileDateList = fileList.flatMap(fileName => getTimeFromName(fileName).map(extracted => (fileName, extracted._1)) ) deleteList = fileDateList.filter(_._2.isBefore(oldThreshold)).map(_._1) _ <- deleteList.traverse(name => Files[IO].delete(Path(s"$FOLDER/${name}"))) - // TODO change to log _ <- deleteList.traverse(name => IO.println(s"delete: $name")) } yield deleteList } diff --git a/src/main/scala/fetchDMI/FetchService.scala b/src/main/scala/fetchDMI/FetchService.scala index af17840..d8a0609 100644 --- a/src/main/scala/fetchDMI/FetchService.scala +++ b/src/main/scala/fetchDMI/FetchService.scala @@ -103,6 +103,7 @@ class FetchService(log: Logger[IO]) { for { dateTimeList <- generateFetchList() resultList <- fetchFromList(dateTimeList) + _ <- IO.println("finish grib downloads") } yield resultList } diff --git a/src/main/scala/fetchDMI/Scheduler.scala b/src/main/scala/fetchDMI/Scheduler.scala index e7b9822..777dda8 100644 --- a/src/main/scala/fetchDMI/Scheduler.scala +++ b/src/main/scala/fetchDMI/Scheduler.scala @@ -11,12 +11,12 @@ import scala.concurrent.duration._ object Scheduler { - def of(name: String, minutes: List[Int]): IO[Scheduler] = { - Slf4jLogger.create[IO].map(logger => new Scheduler(name, minutes, logger)) + def of: IO[Scheduler] = { + Slf4jLogger.create[IO].map(logger => new Scheduler(logger)) } } -class Scheduler(name: String, minutes: List[Int], log: Logger[IO]) { +class Scheduler(log: Logger[IO]) { private def durationToNext(targetMinute: Int)(implicit clock: Clock[IO]): IO[FiniteDuration] = { clock.realTime.map { duration => val now = LocalTime.ofSecondOfDay((duration.toMillis / 1000) % (24 * 60 * 60)) @@ -28,7 +28,7 @@ class Scheduler(name: String, minutes: List[Int], log: Logger[IO]) { } } - def scheduleTask(task: IO[List[String]]): Stream[IO, List[String]] = { + def scheduleTask(name: String, minutes: List[Int], task: IO[List[String]]): Stream[IO, List[String]] = { def streamForMinute(minute: Int): Stream[IO, List[String]] = { Stream.eval(durationToNext(minute)).flatMap { delay => Stream.eval(log.info(s"$name scheduler in: ${delay.toMinutes} min")) *>