moved filter on next phase

This commit is contained in:
Sandro La Bruzzo 2021-10-28 16:13:24 +02:00
parent 1be9aa0a5f
commit d9cbca83f7
1 changed files with 2 additions and 1 deletions

View File

@ -75,7 +75,7 @@ object SparkResolveRelation {
if (targetResolved != null && targetResolved._1.nonEmpty)
currentRelation.setTarget(targetResolved._1)
currentRelation
}.filter(r => !r.getSource.startsWith("unresolved") && !r.getTarget.startsWith("unresolved"))
}
.write
.mode(SaveMode.Overwrite)
.save(s"$workingPath/relation_resolved")
@ -88,6 +88,7 @@ object SparkResolveRelation {
fs.rename(new Path(s"$graphBasePath/relation"), new Path(s"$workingPath/relation"))
spark.read.load(s"$workingPath/relation_resolved").as[Relation]
.filter(r => !r.getSource.startsWith("unresolved") && !r.getTarget.startsWith("unresolved"))
.map(r => mapper.writeValueAsString(r))
.write
.option("compression", "gzip")