|
@@ -24,13 +24,13 @@ class AkkaStreamFileProcessingImpl extends AkkaStreamFileProcessingFuture {
|
24
|
24
|
|
25
|
25
|
implicit val logger: Logger = Logger(LoggerFactory.getLogger(this.getClass))
|
26
|
26
|
|
27
|
|
- override def getPersonByIdFuture(personID: String): Future[Person] = {
|
|
27
|
+ override def getPersonByIdFuture(personID: String): Future[Option[Person]] = {
|
28
|
28
|
val source: Source[Map[String, String], _] = buildAndValidateSource(inputFile = nameBasics)
|
29
|
29
|
|
30
|
30
|
val start: Long = System.currentTimeMillis()
|
31
|
|
- val personFuture: Future[Person] = source
|
|
31
|
+ val personFuture: Future[Option[Person]] = source
|
32
|
32
|
.via(flow = filterPersonByIdFlow(personID = personID))
|
33
|
|
- .runWith(Sink.head[Person])
|
|
33
|
+ .runWith(Sink.headOption[Person])
|
34
|
34
|
|
35
|
35
|
personFuture.andThen({
|
36
|
36
|
case Failure(exception) => logger.error(s"!${exception.printStackTrace()}")
|
|
@@ -39,28 +39,28 @@ class AkkaStreamFileProcessingImpl extends AkkaStreamFileProcessingFuture {
|
39
|
39
|
personFuture
|
40
|
40
|
}
|
41
|
41
|
|
42
|
|
- override def getPersonByNameFuture(primaryName: String): Future[Person] = {
|
|
42
|
+ override def getPersonByNameFuture(primaryName: String): Future[Option[Person]] = {
|
43
|
43
|
val source: Source[Map[String, String], _] = buildAndValidateSource(inputFile = nameBasics)
|
44
|
44
|
|
45
|
45
|
val start: Long = System.currentTimeMillis()
|
46
|
|
- val personFuture: Future[Person] = source
|
|
46
|
+ val personFuture: Future[Option[Person]] = source
|
47
|
47
|
.via(flow = filterPersonByNameFlow(primaryName = primaryName))
|
48
|
|
- .runWith(Sink.head[Person])
|
|
48
|
+ .runWith(Sink.headOption[Person])
|
49
|
49
|
|
50
|
50
|
personFuture
|
51
|
51
|
}
|
52
|
52
|
|
53
|
|
- override def getTvSerieByIdFuture(tvSerieID: String): Future[TvSerie] = {
|
|
53
|
+ override def getTvSerieByIdFuture(tvSerieID: String): Future[Option[TvSerie]] = {
|
54
|
54
|
val source: Source[Map[String, String], _] = buildAndValidateSource(inputFile = titleBasics)
|
55
|
55
|
|
56
|
56
|
val start: Long = System.currentTimeMillis()
|
57
|
|
- val tvSerieFuture: Future[TvSerie] = source
|
|
57
|
+ val tvSerieFuture: Future[Option[TvSerie]] = source
|
58
|
58
|
.via(flow = filterTvSerieByIdFlow(tvSerieID = tvSerieID))
|
59
|
|
- .runWith(Sink.head[TvSerie])
|
|
59
|
+ .runWith(Sink.headOption[TvSerie])
|
60
|
60
|
|
61
|
61
|
tvSerieFuture.onComplete({
|
62
|
62
|
case Failure(exception) => logger.info(s"$exception")
|
63
|
|
- case Success(value: TvSerie) =>
|
|
63
|
+ case Success(value: Option[TvSerie]) =>
|
64
|
64
|
logger.info(s"$value")
|
65
|
65
|
logger.info(s"SUCCESS, elapsed time:${(System.currentTimeMillis() - start) / 1000} sec")
|
66
|
66
|
})
|