DataService to save last 24h of data
This commit is contained in:
@@ -18,17 +18,11 @@ object DBService {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class DBService(log: Logger[IO]) {
|
class DBService(log: Logger[IO]) extends DataServiceTrait {
|
||||||
private val dateFormatter = DateTimeFormatter.ofPattern("yyyyMMdd_HHmm")
|
private val dateFormatter = DateTimeFormatter.ofPattern("yyyyMMdd_HHmm")
|
||||||
private val dataPath = "./data"
|
private val dataPath = "./data"
|
||||||
private val nonDuplicatedLines = 34 // takes only first 34 lines of data as rest after 'Zosēni' is duplicated
|
private val nonDuplicatedLines = 34 // takes only first 34 lines of data as rest after 'Zosēni' is duplicated
|
||||||
|
|
||||||
def readFile(fileName: String): IO[List[String]] = {
|
|
||||||
val file = new File(dataPath, fileName)
|
|
||||||
val sourceResource = Resource.fromAutoCloseable(IO(Source.fromFile(file)))
|
|
||||||
sourceResource.use(source => IO(source.getLines().take(nonDuplicatedLines).toList)).handleError(_ => List.empty)
|
|
||||||
}
|
|
||||||
|
|
||||||
private def readFileNames(path: String): IO[List[String]] =
|
private def readFileNames(path: String): IO[List[String]] =
|
||||||
IO(new File(path).listFiles.toList.map(_.getName))
|
IO(new File(path).listFiles.toList.map(_.getName))
|
||||||
.handleError(_ => List.empty)
|
.handleError(_ => List.empty)
|
||||||
@@ -46,14 +40,6 @@ class DBService(log: Logger[IO]) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
def getInRange(from: LocalDateTime, to: LocalDateTime): IO[List[String]] = {
|
|
||||||
for {
|
|
||||||
fileNames <- readFileNames(dataPath)
|
|
||||||
.map (_.filter (inRange (_, from, to)))
|
|
||||||
fileLines <- fileNames.traverse(readFile)
|
|
||||||
} yield fileLines.flatten
|
|
||||||
}
|
|
||||||
|
|
||||||
def save(fileName: String, content: String): IO[Either[Throwable, String]] = {
|
def save(fileName: String, content: String): IO[Either[Throwable, String]] = {
|
||||||
val path = Paths.get(s"$dataPath/$fileName")
|
val path = Paths.get(s"$dataPath/$fileName")
|
||||||
IO(Files.writeString(path, content))
|
IO(Files.writeString(path, content))
|
||||||
@@ -64,8 +50,22 @@ class DBService(log: Logger[IO]) {
|
|||||||
)
|
)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def readFile(fileName: String): IO[List[String]] = {
|
||||||
|
val file = new File(dataPath, fileName)
|
||||||
|
val sourceResource = Resource.fromAutoCloseable(IO(Source.fromFile(file)))
|
||||||
|
sourceResource.use(source => IO(source.getLines().take(nonDuplicatedLines).toList)).handleError(_ => List.empty)
|
||||||
|
}
|
||||||
|
|
||||||
|
def getInRange(from: LocalDateTime, to: LocalDateTime): IO[List[String]] = {
|
||||||
|
for {
|
||||||
|
fileNames <- readFileNames(dataPath)
|
||||||
|
.map (_.filter (inRange (_, from, to)))
|
||||||
|
fileLines <- fileNames.traverse(readFile)
|
||||||
|
} yield fileLines.flatten
|
||||||
|
}
|
||||||
|
|
||||||
// dates in which we have saved data
|
// dates in which we have saved data
|
||||||
def getDates(): IO[List[LocalDate]] = {
|
def getDates: IO[List[LocalDate]] = {
|
||||||
val formatter = DateTimeFormatter.ofPattern("yyyyMMdd")
|
val formatter = DateTimeFormatter.ofPattern("yyyyMMdd")
|
||||||
for {
|
for {
|
||||||
fileNames <- readFileNames(dataPath)
|
fileNames <- readFileNames(dataPath)
|
||||||
|
|||||||
@@ -0,0 +1,71 @@
|
|||||||
|
package db
|
||||||
|
|
||||||
|
import cats.effect.{Clock, IO, Ref}
|
||||||
|
import cats.implicits.toTraverseOps
|
||||||
|
import fetch.FileNameService
|
||||||
|
import org.typelevel.log4cats.Logger
|
||||||
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
|
|
||||||
|
import java.time.{Instant, LocalDate, LocalDateTime, ZoneId}
|
||||||
|
|
||||||
|
trait DataServiceTrait {
|
||||||
|
def save(fileName: String, content: String): IO[Either[Throwable, String]]
|
||||||
|
def readFile(fileName: String): IO[List[String]]
|
||||||
|
def getInRange(from: LocalDateTime, to: LocalDateTime): IO[List[String]]
|
||||||
|
def getDates: IO[List[LocalDate]]
|
||||||
|
def getDateFileNames(date: LocalDate): IO[List[String]]
|
||||||
|
}
|
||||||
|
|
||||||
|
object DataService {
|
||||||
|
def of(dbService: DBService): IO[DataService] = {
|
||||||
|
for {
|
||||||
|
log <- Slf4jLogger.create[IO]
|
||||||
|
fileNameService = new FileNameService()
|
||||||
|
fileNames <- fileNameService.generateLast24Hours
|
||||||
|
contents <- fileNames.traverse(fileName => dbService.readFile(fileName)
|
||||||
|
.map(content => (fileName, content)))
|
||||||
|
state = contents.toMap
|
||||||
|
stateRef <- Ref.of[IO, Map[String, List[String]]](state)
|
||||||
|
} yield new DataService(dbService, new FileNameService(), log, stateRef)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
class DataService private(
|
||||||
|
dbService: DBService,
|
||||||
|
fileNameService: FileNameService,
|
||||||
|
log: Logger[IO],
|
||||||
|
private val state: Ref[IO, Map[String, List[String]]]
|
||||||
|
) extends DataServiceTrait {
|
||||||
|
|
||||||
|
private def logState: IO[Unit] = {
|
||||||
|
state.get.flatMap(currentState => log.info(s"State keys: ${currentState.keys.size}"))
|
||||||
|
}
|
||||||
|
|
||||||
|
private def filterState: IO[Unit] = {
|
||||||
|
fileNameService.generateLast24Hours.flatMap { last24Hours =>
|
||||||
|
state.update(st => st.filterKeys(last24Hours.contains).toMap)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
def save(fileName: String, content: String): IO[Either[Throwable, String]] = {
|
||||||
|
dbService.save(fileName, content).flatMap {
|
||||||
|
case Right(savedFileName) =>
|
||||||
|
state.update(st => st.updated(savedFileName, content.split("\n").toList)) *>
|
||||||
|
filterState *>
|
||||||
|
logState.as(Right(savedFileName))
|
||||||
|
case e@Left(_) => IO.pure(e)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
def readFile(fileName: String): IO[List[String]] = dbService.readFile(fileName)
|
||||||
|
|
||||||
|
def getInRange(from: LocalDateTime, to: LocalDateTime): IO[List[String]] = dbService.getInRange(from, to)
|
||||||
|
|
||||||
|
def getDates: IO[List[LocalDate]] = dbService.getDates
|
||||||
|
|
||||||
|
def getDateFileNames(date: LocalDate): IO[List[String]] = dbService.getDateFileNames(date)
|
||||||
|
|
||||||
|
// TODO implement getting full data from state
|
||||||
|
def getLast24Hours: IO[List[String]] = {
|
||||||
|
state.get.map(_.keys.toList.sorted)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -27,7 +27,7 @@ object Main {
|
|||||||
for {
|
for {
|
||||||
log <- Slf4jLogger.create[IO]
|
log <- Slf4jLogger.create[IO]
|
||||||
dbService <- DBService.of
|
dbService <- DBService.of
|
||||||
dates <- dbService.getDates()
|
dates <- dbService.getDates
|
||||||
_ <- log.info(s"$dates")
|
_ <- log.info(s"$dates")
|
||||||
} yield ()
|
} yield ()
|
||||||
}
|
}
|
||||||
@@ -41,9 +41,17 @@ object Main {
|
|||||||
} yield ()
|
} yield ()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private def testDataService: IO[Unit] = {
|
||||||
|
for {
|
||||||
|
dbService <- DBService.of
|
||||||
|
dataService <- DataService.of(dbService)
|
||||||
|
} yield ()
|
||||||
|
}
|
||||||
|
|
||||||
def main(args: Array[String]): Unit = {
|
def main(args: Array[String]): Unit = {
|
||||||
// testGetInRange.unsafeRunSync()
|
// testGetInRange.unsafeRunSync()
|
||||||
// testGetDates.unsafeRunSync()
|
// testGetDates.unsafeRunSync()
|
||||||
testGetDate.unsafeRunSync()
|
// testGetDate.unsafeRunSync()
|
||||||
|
testDataService.unsafeRunSync()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,23 +1,23 @@
|
|||||||
package fetch
|
package fetch
|
||||||
|
|
||||||
import cats.effect.IO
|
import cats.effect.IO
|
||||||
import db.DBService
|
import db.{DBService, DataServiceTrait}
|
||||||
import fs2.Stream
|
import fs2.Stream
|
||||||
import org.typelevel.log4cats.Logger
|
import org.typelevel.log4cats.Logger
|
||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
|
|
||||||
|
|
||||||
object FileFetchScheduler {
|
object FileFetchScheduler {
|
||||||
def of(dbService: DBService, fetch: FetchServiceTrait): IO[FileFetchScheduler] = {
|
def of(dataService: DataServiceTrait, fetch: FetchServiceTrait): IO[FileFetchScheduler] = {
|
||||||
Scheduler.of.flatMap { scheduler =>
|
Scheduler.of.flatMap { scheduler =>
|
||||||
Slf4jLogger.create[IO].map {
|
Slf4jLogger.create[IO].map {
|
||||||
new FileFetchScheduler(dbService, fetch, new FileNameService(), scheduler, _)
|
new FileFetchScheduler(dataService, fetch, new FileNameService(), scheduler, _)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class FileFetchScheduler(dbService: DBService, fetch: FetchServiceTrait, fileNameService: FileNameService, scheduler: Scheduler, log: Logger[IO]) {
|
class FileFetchScheduler(dataService: DataServiceTrait, fetch: FetchServiceTrait, fileNameService: FileNameService, scheduler: Scheduler, log: Logger[IO]) {
|
||||||
def run: Stream[IO, Unit] = {
|
def run: Stream[IO, Unit] = {
|
||||||
val fetchTask = fileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
val fetchTask = fileNameService.generateCurrentHour.flatMap(fetch.fetchSingleFile)
|
||||||
scheduler.scheduleTask(fetchTask)
|
scheduler.scheduleTask(fetchTask)
|
||||||
@@ -25,7 +25,7 @@ class FileFetchScheduler(dbService: DBService, fetch: FetchServiceTrait, fileNam
|
|||||||
case Left(fetchErr) =>
|
case Left(fetchErr) =>
|
||||||
log.error(s"Fetch error: $fetchErr")
|
log.error(s"Fetch error: $fetchErr")
|
||||||
case Right((name, content)) =>
|
case Right((name, content)) =>
|
||||||
dbService.save(name, content).attempt.flatMap {
|
dataService.save(name, content).attempt.flatMap {
|
||||||
case Left(err) => log.error(s"error: $err")
|
case Left(err) => log.error(s"error: $err")
|
||||||
case Right(saveResult) => saveResult match {
|
case Right(saveResult) => saveResult match {
|
||||||
case Left(err) => log.error(s"error: $err")
|
case Left(err) => log.error(s"error: $err")
|
||||||
|
|||||||
@@ -1,18 +1,19 @@
|
|||||||
package server
|
package server
|
||||||
import cats.effect._
|
import cats.effect._
|
||||||
import cats.implicits.catsSyntaxTuple2Parallel
|
import cats.implicits.catsSyntaxTuple2Parallel
|
||||||
import db.DBService
|
import db.{DBService, DataService}
|
||||||
import fetch.{FetchService, FileFetchScheduler, StatefulFetchService}
|
import fetch.{FetchService, FileFetchScheduler, StatefulFetchService}
|
||||||
|
|
||||||
object Main extends IOApp {
|
object Main extends IOApp {
|
||||||
def run(args: List[String]): IO[ExitCode] = {
|
def run(args: List[String]): IO[ExitCode] = {
|
||||||
for {
|
for {
|
||||||
dbService <- DBService.of
|
dbService <- DBService.of
|
||||||
|
dataService <- DataService.of(dbService)
|
||||||
fetch <- FetchService.of
|
fetch <- FetchService.of
|
||||||
statefulFetch <- StatefulFetchService.of(fetch)
|
statefulFetch <- StatefulFetchService.of(fetch)
|
||||||
fileFetchScheduler <- FileFetchScheduler.of(dbService, statefulFetch)
|
fileFetchScheduler <- FileFetchScheduler.of(dataService, statefulFetch)
|
||||||
schedulerTask = fileFetchScheduler.run.compile.drain
|
schedulerTask = fileFetchScheduler.run.compile.drain
|
||||||
server <- Server.of(dbService, statefulFetch)
|
server <- Server.of(dataService, statefulFetch)
|
||||||
serverTask = server.run
|
serverTask = server.run
|
||||||
exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success)
|
exitCode <- (serverTask, schedulerTask).parMapN((_, _) => ExitCode.Success)
|
||||||
} yield exitCode
|
} yield exitCode
|
||||||
|
|||||||
@@ -3,7 +3,7 @@ package server
|
|||||||
import cats.effect._
|
import cats.effect._
|
||||||
import cats.implicits.toTraverseOps
|
import cats.implicits.toTraverseOps
|
||||||
import com.comcast.ip4s.IpLiteralSyntax
|
import com.comcast.ip4s.IpLiteralSyntax
|
||||||
import db.DBService
|
import db.DataService
|
||||||
import fetch.FetchServiceTrait
|
import fetch.FetchServiceTrait
|
||||||
import parse.{Parser, WeatherData}
|
import parse.{Parser, WeatherData}
|
||||||
import server.ValidateRoutes.{AggKey, CityList, DateTimeRange, ValidDate}
|
import server.ValidateRoutes.{AggKey, CityList, DateTimeRange, ValidDate}
|
||||||
@@ -22,18 +22,19 @@ import parse.Aggregate.{AggregateKey, UserQuery}
|
|||||||
import org.http4s.circe.jsonEncoder
|
import org.http4s.circe.jsonEncoder
|
||||||
import org.typelevel.log4cats.Logger
|
import org.typelevel.log4cats.Logger
|
||||||
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
import org.typelevel.log4cats.slf4j.Slf4jLogger
|
||||||
|
|
||||||
import scala.concurrent.duration.DurationInt
|
import scala.concurrent.duration.DurationInt
|
||||||
|
|
||||||
|
|
||||||
object Server {
|
object Server {
|
||||||
def of(dbService: DBService, fetch: FetchServiceTrait): IO[Server] = {
|
def of(dataService: DataService, fetch: FetchServiceTrait): IO[Server] = {
|
||||||
Slf4jLogger.create[IO].map {
|
Slf4jLogger.create[IO].map {
|
||||||
new Server(dbService, fetch, _)
|
new Server(dataService, fetch, _)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
class Server(dbService: DBService, fetch: FetchServiceTrait, log: Logger[IO]) {
|
class Server(dataService: DataService, fetch: FetchServiceTrait, log: Logger[IO]) {
|
||||||
|
|
||||||
// Define the extension method `pretty` for Json
|
// Define the extension method `pretty` for Json
|
||||||
implicit class JsonPrettyPrinter(json: Json) {
|
implicit class JsonPrettyPrinter(json: Json) {
|
||||||
@@ -47,7 +48,7 @@ class Server(dbService: DBService, fetch: FetchServiceTrait, log: Logger[IO]) {
|
|||||||
|
|
||||||
// http://0.0.0.0:8080/api/query/20230414_2200-20230501_1230/Liepāja,Rēzekne/tempMax/max
|
// http://0.0.0.0:8080/api/query/20230414_2200-20230501_1230/Liepāja,Rēzekne/tempMax/max
|
||||||
case GET -> Root / "query" / DateTimeRange(from, to) / CityList(cities) / field / AggKey(key) =>
|
case GET -> Root / "query" / DateTimeRange(from, to) / CityList(cities) / field / AggKey(key) =>
|
||||||
dbService.getInRange(from, to)
|
dataService.getInRange(from, to)
|
||||||
.map(Parser.queryData(UserQuery(cities, field, key), _))
|
.map(Parser.queryData(UserQuery(cities, field, key), _))
|
||||||
.flatMap(result => Ok(result.asJson.pretty))
|
.flatMap(result => Ok(result.asJson.pretty))
|
||||||
|
|
||||||
@@ -58,7 +59,7 @@ class Server(dbService: DBService, fetch: FetchServiceTrait, log: Logger[IO]) {
|
|||||||
fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList
|
fetchServiceError = fetchResultEither.left.toOption.map(e => s"FetchServiceError: ${e.getMessage}").toList
|
||||||
fetchResult = fetchResultEither.getOrElse(List.empty)
|
fetchResult = fetchResultEither.getOrElse(List.empty)
|
||||||
(fetchErrors, successDownloads) = fetchResult.partitionMap(identity)
|
(fetchErrors, successDownloads) = fetchResult.partitionMap(identity)
|
||||||
saveResults <- successDownloads.traverse { case (name, content) => dbService.save(name, content) }
|
saveResults <- successDownloads.traverse { case (name, content) => dataService.save(name, content) }
|
||||||
(saveErrors, successSaves) = saveResults.partitionMap(identity)
|
(saveErrors, successSaves) = saveResults.partitionMap(identity)
|
||||||
// successes = successDownloads.map(s => s"fetched: ${s._1}") ++ successSaves.map(s => s"saved: $s")
|
// successes = successDownloads.map(s => s"fetched: ${s._1}") ++ successSaves.map(s => s"saved: $s")
|
||||||
successes = successSaves
|
successes = successSaves
|
||||||
@@ -76,19 +77,24 @@ class Server(dbService: DBService, fetch: FetchServiceTrait, log: Logger[IO]) {
|
|||||||
|
|
||||||
// http://0.0.0.0:8080/api/show/all_dates
|
// http://0.0.0.0:8080/api/show/all_dates
|
||||||
case GET -> Root / "show" / "all_dates" =>
|
case GET -> Root / "show" / "all_dates" =>
|
||||||
dbService.getDates().flatMap(dates =>
|
dataService.getDates.flatMap(dates =>
|
||||||
Ok(dates.asJson.pretty)
|
Ok(dates.asJson.pretty)
|
||||||
)
|
)
|
||||||
|
|
||||||
// http://0.0.0.0:8080/api/show/date/20230423
|
// http://0.0.0.0:8080/api/show/date/20230423
|
||||||
case GET -> Root / "show" / "date" / ValidDate(date) =>
|
case GET -> Root / "show" / "date" / ValidDate(date) =>
|
||||||
dbService.getDateFileNames(date).flatMap(fileNames =>
|
dataService.getDateFileNames(date).flatMap(fileNames =>
|
||||||
Ok(fileNames.asJson.pretty)
|
Ok(fileNames.asJson.pretty)
|
||||||
)
|
)
|
||||||
|
|
||||||
// http://0.0.0.0:8080/api/show/file/20230423_12:30.csv
|
// http://0.0.0.0:8080/api/show/file/20230423_12:30.csv
|
||||||
case GET -> Root / "show" / "file" / (fileName: String) =>
|
case GET -> Root / "show" / "file" / (fileName: String) =>
|
||||||
dbService.readFile(fileName).flatMap(content => Ok(content.asJson))
|
dataService.readFile(fileName).flatMap(content => Ok(content.asJson))
|
||||||
|
|
||||||
|
// http://0.0.0.0:8080/api/getLast24hours
|
||||||
|
case GET -> Root / "getLast24hours" => {
|
||||||
|
dataService.getLast24Hours.flatMap(content => Ok(content.asJson.pretty))
|
||||||
|
}
|
||||||
|
|
||||||
// http://0.0.0.0:8080/api/help
|
// http://0.0.0.0:8080/api/help
|
||||||
case GET -> Root / "help" => {
|
case GET -> Root / "help" => {
|
||||||
|
|||||||
Reference in New Issue
Block a user