Fetch weather stations in cron from ftp server
This commit is contained in:
+16
-12
@@ -5,6 +5,7 @@ import data.DataService
|
|||||||
import db.{DBConnection, PostgresService}
|
import db.{DBConnection, PostgresService}
|
||||||
import fetch.csv.{FetchService, FileNameService}
|
import fetch.csv.{FetchService, FileNameService}
|
||||||
import fetch.dmi
|
import fetch.dmi
|
||||||
|
import fetch.lvgmc
|
||||||
import scheduler.Scheduler
|
import scheduler.Scheduler
|
||||||
import server.Server
|
import server.Server
|
||||||
|
|
||||||
@@ -22,25 +23,28 @@ object Main extends IOApp {
|
|||||||
|
|
||||||
scheduler <- Scheduler.of
|
scheduler <- Scheduler.of
|
||||||
dataService <- DataService.of
|
dataService <- DataService.of
|
||||||
fetchService <- FetchService.of
|
fetchLegacyService <- FetchService.of
|
||||||
|
fetchLvgmcService <- lvgmc.FetchService.of
|
||||||
|
|
||||||
fetchCsvList = new FileNameService().generateCurrentHour
|
|
||||||
.flatMap(fetchService.fetchSingleFile)
|
|
||||||
.flatMap {
|
|
||||||
case Right((name, content)) => postgresService.save(name, content)
|
|
||||||
case Left(_) => IO.unit // ignore error
|
|
||||||
}
|
|
||||||
fetchCsvTask = scheduler.scheduleTask("Fetch CSV", List(31), fetchCsvList).compile.drain
|
|
||||||
|
|
||||||
cleanupTask = scheduler.scheduleTask("Cleanup", List(1), dataService.deleteOldForecasts()).compile.drain
|
fetchWeatherStations = for {
|
||||||
|
fileName <- new FileNameService().generateCurrentHour
|
||||||
|
stationDataStr <- fetchLvgmcService.fetchWeatherStations()
|
||||||
|
_ <- IO.println(fileName)
|
||||||
|
// TODO refactor after legacy deleted. No need for fileName, instead pass only year
|
||||||
|
_ <- postgresService.save(fileName, stationDataStr)
|
||||||
|
} yield ()
|
||||||
|
fetchStationsTask = scheduler.scheduleTask("Fetch Weather Stations", List(11,13,23,30), fetchWeatherStations).compile.drain
|
||||||
|
|
||||||
|
cleanupTask = scheduler.scheduleTask("Cleanup old Grib", List(41), dataService.deleteOldForecasts()).compile.drain
|
||||||
|
|
||||||
fetchGrib <- dmi.FetchService.of(dataService)
|
fetchGrib <- dmi.FetchService.of(dataService)
|
||||||
fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(3), fetchGrib.fetchRecentForecasts()).compile.drain
|
fetchGribTask = scheduler.scheduleTask("Fetch Grib", List(43), fetchGrib.fetchRecentForecasts()).compile.drain
|
||||||
|
|
||||||
server <- Server.of(postgresService, dataService, fetchService)
|
server <- Server.of(postgresService, dataService, fetchLegacyService)
|
||||||
serverTask = server.run
|
serverTask = server.run
|
||||||
|
|
||||||
exitCode <- (serverTask, fetchCsvTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ExitCode.Success)
|
exitCode <- (serverTask, fetchStationsTask, cleanupTask, fetchGribTask).parMapN((_, _, _, _) => ExitCode.Success)
|
||||||
} yield exitCode
|
} yield exitCode
|
||||||
|
|
||||||
program.handleErrorWith { error =>
|
program.handleErrorWith { error =>
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ import org.typelevel.log4cats.Logger
|
|||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
|
|
||||||
import java.io.ByteArrayOutputStream
|
import java.io.ByteArrayOutputStream
|
||||||
|
import java.nio.charset.StandardCharsets
|
||||||
|
|
||||||
final case class LVGMCServerConfig(
|
final case class LVGMCServerConfig(
|
||||||
username: String,
|
username: String,
|
||||||
@@ -85,4 +86,8 @@ class FetchService(log: Logger[IO]) {
|
|||||||
} yield content
|
} yield content
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def fetchWeatherStations(): IO[String] = {
|
||||||
|
fetchFile("Latvija_faktiskais_laiks.csv").map(result => new String(result, StandardCharsets.UTF_8))
|
||||||
|
}
|
||||||
}
|
}
|
||||||
@@ -2,12 +2,18 @@ package fetch.lvgmc
|
|||||||
|
|
||||||
import cats.effect.IO
|
import cats.effect.IO
|
||||||
import cats.effect.unsafe.implicits.global
|
import cats.effect.unsafe.implicits.global
|
||||||
|
import db.{DBConnection, PostgresService}
|
||||||
|
import fetch.csv.FileNameService
|
||||||
|
|
||||||
|
import java.nio.charset.StandardCharsets
|
||||||
|
|
||||||
object FetchServiceTest {
|
object FetchServiceTest {
|
||||||
def main(args: Array[String]): Unit = {
|
def main(args: Array[String]): Unit = {
|
||||||
fetchFile().unsafeRunSync()
|
// fetchFile().unsafeRunSync()
|
||||||
|
fetchAndSave().unsafeRunSync()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
// Eiropa_LTV_pilsetas_nakama_dn.csv
|
// Eiropa_LTV_pilsetas_nakama_dn.csv
|
||||||
// Eiropa_LTV_pilsetas_tekosa_dn.csv
|
// Eiropa_LTV_pilsetas_tekosa_dn.csv
|
||||||
// Latvija_LTV_pilsetas_nakama_dnn.csv
|
// Latvija_LTV_pilsetas_nakama_dnn.csv
|
||||||
@@ -22,4 +28,25 @@ object FetchServiceTest {
|
|||||||
} yield ()
|
} yield ()
|
||||||
program
|
program
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private def fetchAndSave (): IO[Unit] = {
|
||||||
|
var program = for {
|
||||||
|
fileName <- new FileNameService().generateCurrentHour
|
||||||
|
_ <- IO.println(fileName)
|
||||||
|
|
||||||
|
fetch <- FetchService.of
|
||||||
|
result <- fetch.fetchFile("Latvija_faktiskais_laiks.csv")
|
||||||
|
resultStr = new String(result, StandardCharsets.UTF_8)
|
||||||
|
_ <- IO.println(s"fetched file, size: ${result.length}")
|
||||||
|
// _ <- IO.println(resultStr)
|
||||||
|
|
||||||
|
transactor <- DBConnection.transactor[IO]
|
||||||
|
postgresService <- PostgresService.of(transactor)
|
||||||
|
|
||||||
|
_ <- postgresService.save(fileName, resultStr)
|
||||||
|
} yield ()
|
||||||
|
program
|
||||||
|
}
|
||||||
|
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user