Create just one scheduler instance
This commit is contained in:
@@ -18,12 +18,11 @@ object Main extends IOApp {
|
|||||||
fileFetchScheduler <- FileFetchScheduler.of(postgresService, fetch)
|
fileFetchScheduler <- FileFetchScheduler.of(postgresService, fetch)
|
||||||
fetchCsvTask = fileFetchScheduler.run.compile.drain
|
fetchCsvTask = fileFetchScheduler.run.compile.drain
|
||||||
|
|
||||||
scheduler <- fetchDMI.Scheduler.of("Cleanup", List(1))
|
scheduler <- fetchDMI.Scheduler.of
|
||||||
cleanupTask = scheduler.scheduleTask(DataService.deleteOldForecasts()).compile.drain
|
cleanupTask = scheduler.scheduleTask("Cleanup", List(1), DataService.deleteOldForecasts()).compile.drain
|
||||||
|
|
||||||
scheduler <- fetchDMI.Scheduler.of("Fetch Grib", List(2))
|
|
||||||
fetchGrib <- fetchDMI.FetchService.of
|
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)
|
server <- Server.of(postgresService, fetch)
|
||||||
serverTask = server.run
|
serverTask = server.run
|
||||||
|
|||||||
@@ -54,13 +54,13 @@ object DataService {
|
|||||||
val oldThreshold = nowUTC.minusHours(maxHours)
|
val oldThreshold = nowUTC.minusHours(maxHours)
|
||||||
|
|
||||||
for {
|
for {
|
||||||
|
_ <- IO.println("start cleanup")
|
||||||
fileList <- getFileList()
|
fileList <- getFileList()
|
||||||
fileDateList = fileList.flatMap(fileName =>
|
fileDateList = fileList.flatMap(fileName =>
|
||||||
getTimeFromName(fileName).map(extracted => (fileName, extracted._1))
|
getTimeFromName(fileName).map(extracted => (fileName, extracted._1))
|
||||||
)
|
)
|
||||||
deleteList = fileDateList.filter(_._2.isBefore(oldThreshold)).map(_._1)
|
deleteList = fileDateList.filter(_._2.isBefore(oldThreshold)).map(_._1)
|
||||||
_ <- deleteList.traverse(name => Files[IO].delete(Path(s"$FOLDER/${name}")))
|
_ <- deleteList.traverse(name => Files[IO].delete(Path(s"$FOLDER/${name}")))
|
||||||
// TODO change to log
|
|
||||||
_ <- deleteList.traverse(name => IO.println(s"delete: $name"))
|
_ <- deleteList.traverse(name => IO.println(s"delete: $name"))
|
||||||
} yield deleteList
|
} yield deleteList
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -103,6 +103,7 @@ class FetchService(log: Logger[IO]) {
|
|||||||
for {
|
for {
|
||||||
dateTimeList <- generateFetchList()
|
dateTimeList <- generateFetchList()
|
||||||
resultList <- fetchFromList(dateTimeList)
|
resultList <- fetchFromList(dateTimeList)
|
||||||
|
_ <- IO.println("finish grib downloads")
|
||||||
} yield resultList
|
} yield resultList
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -11,12 +11,12 @@ import scala.concurrent.duration._
|
|||||||
|
|
||||||
|
|
||||||
object Scheduler {
|
object Scheduler {
|
||||||
def of(name: String, minutes: List[Int]): IO[Scheduler] = {
|
def of: IO[Scheduler] = {
|
||||||
Slf4jLogger.create[IO].map(logger => new Scheduler(name, minutes, logger))
|
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] = {
|
private def durationToNext(targetMinute: Int)(implicit clock: Clock[IO]): IO[FiniteDuration] = {
|
||||||
clock.realTime.map { duration =>
|
clock.realTime.map { duration =>
|
||||||
val now = LocalTime.ofSecondOfDay((duration.toMillis / 1000) % (24 * 60 * 60))
|
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]] = {
|
def streamForMinute(minute: Int): Stream[IO, List[String]] = {
|
||||||
Stream.eval(durationToNext(minute)).flatMap { delay =>
|
Stream.eval(durationToNext(minute)).flatMap { delay =>
|
||||||
Stream.eval(log.info(s"$name scheduler in: ${delay.toMinutes} min")) *>
|
Stream.eval(log.info(s"$name scheduler in: ${delay.toMinutes} min")) *>
|
||||||
|
|||||||
Reference in New Issue
Block a user