Skip to content

Commit ba55d85

Browse files
authored
Deadletter path is being incorrectly joined (#1144)
* fix deadletter path join Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com> * cleanup Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com>
1 parent 32ddf19 commit ba55d85

2 files changed

Lines changed: 5 additions & 8 deletions

File tree

spark/ingestion/src/main/scala/feast/ingestion/BatchPipeline.scala

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -16,14 +16,12 @@
1616
*/
1717
package feast.ingestion
1818

19-
import java.nio.file.Paths
20-
2119
import feast.ingestion.sources.bq.BigQueryReader
2220
import feast.ingestion.sources.file.FileReader
2321
import feast.ingestion.validation.{RowValidator, TypeCheck}
22+
import org.apache.commons.lang.StringUtils
2423
import org.apache.spark.SparkEnv
25-
import org.apache.spark.sql.{Column, SaveMode, SparkSession}
26-
import org.apache.spark.sql.functions.col
24+
import org.apache.spark.sql.{SaveMode, SparkSession}
2725

2826
/**
2927
* Batch Ingestion Flow:
@@ -83,7 +81,7 @@ object BatchPipeline extends BasePipeline {
8381
.write
8482
.format("parquet")
8583
.mode(SaveMode.Append)
86-
.save(Paths.get(path, SparkEnv.get.conf.getAppId).toString)
84+
.save(StringUtils.stripEnd(path, "/") + "/" + SparkEnv.get.conf.getAppId)
8785
case _ => None
8886
}
8987

spark/ingestion/src/main/scala/feast/ingestion/StreamingPipeline.scala

Lines changed: 2 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -16,13 +16,12 @@
1616
*/
1717
package feast.ingestion
1818

19-
import java.nio.file.Paths
20-
2119
import feast.ingestion.registry.proto.ProtoRegistryFactory
2220
import org.apache.spark.sql.{DataFrame, Row, SaveMode, SparkSession}
2321
import org.apache.spark.sql.functions.udf
2422
import feast.ingestion.utils.ProtoReflection
2523
import feast.ingestion.validation.{RowValidator, TypeCheck}
24+
import org.apache.commons.lang.StringUtils
2625
import org.apache.spark.SparkEnv
2726
import org.apache.spark.sql.streaming.StreamingQuery
2827
import org.apache.spark.sql.avro._
@@ -97,7 +96,7 @@ object StreamingPipeline extends BasePipeline with Serializable {
9796
.write
9897
.format("parquet")
9998
.mode(SaveMode.Append)
100-
.save(Paths.get(path, SparkEnv.get.conf.getAppId).toString)
99+
.save(StringUtils.stripEnd(path, "/") + "/" + SparkEnv.get.conf.getAppId)
101100
case _ =>
102101
batchDF
103102
.filter(!validator.checkAll)

0 commit comments

Comments
 (0)