Skip to content

Commit eaaf233

Browse files
authored
Optimize memory footprint for Spark Ingestion Job (#1265)
Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com>
1 parent 9e5c41e commit eaaf233

4 files changed

Lines changed: 21 additions & 13 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -69,7 +69,7 @@ object BatchPipeline extends BasePipeline {
6969
}
7070

7171
val validRows = projected
72-
.mapPartitions(metrics.incrementRead)
72+
.map(metrics.incrementRead)
7373
.filter(rowValidator.allChecks)
7474

7575
validRows.write
@@ -85,7 +85,7 @@ object BatchPipeline extends BasePipeline {
8585
case Some(path) =>
8686
projected
8787
.filter(!rowValidator.allChecks)
88-
.mapPartitions(metrics.incrementDeadLetters)
88+
.map(metrics.incrementDeadLetters)
8989
.write
9090
.format("parquet")
9191
.mode(SaveMode.Append)

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -107,7 +107,7 @@ object StreamingPipeline extends BasePipeline with Serializable {
107107
implicit def rowEncoder: Encoder[Row] = RowEncoder(rowsAfterValidation.schema)
108108

109109
rowsAfterValidation
110-
.mapPartitions(metrics.incrementRead)
110+
.map(metrics.incrementRead)
111111
.filter(if (config.doNotIngestInvalidRows) expr("_isValid") else rowValidator.allChecks)
112112
.write
113113
.format("feast.ingestion.stores.redis")
@@ -122,7 +122,7 @@ object StreamingPipeline extends BasePipeline with Serializable {
122122
case Some(path) =>
123123
rowsAfterValidation
124124
.filter("!_isValid")
125-
.mapPartitions(metrics.incrementDeadLetters)
125+
.map(metrics.incrementDeadLetters)
126126
.write
127127
.format("parquet")
128128
.mode(SaveMode.Append)

spark/ingestion/src/main/scala/feast/ingestion/metrics/IngestionPipelineMetrics.scala

Lines changed: 16 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -22,20 +22,28 @@ import org.apache.spark.sql.Row
2222

2323
class IngestionPipelineMetrics extends Serializable {
2424

25-
def incrementDeadLetters(rowIterator: Iterator[Row]): Iterator[Row] = {
26-
val materialized = rowIterator.toArray
25+
def incrementDeadLetters(row: Row): Row = {
2726
if (metricSource.nonEmpty)
28-
metricSource.get.METRIC_DEADLETTER_ROWS_INSERTED.inc(materialized.length)
27+
metricSource.get.METRIC_DEADLETTER_ROWS_INSERTED.inc()
2928

30-
materialized.toIterator
29+
row
3130
}
3231

33-
def incrementRead(rowIterator: Iterator[Row]): Iterator[Row] = {
34-
val materialized = rowIterator.toArray
32+
def incrementRead(row: Row): Row = {
3533
if (metricSource.nonEmpty)
36-
metricSource.get.METRIC_ROWS_READ_FROM_SOURCE.inc(materialized.length)
34+
metricSource.get.METRIC_ROWS_READ_FROM_SOURCE.inc()
3735

38-
materialized.toIterator
36+
row
37+
}
38+
39+
def incrementRead(inc: Long): Unit = {
40+
if (metricSource.nonEmpty)
41+
metricSource.get.METRIC_ROWS_READ_FROM_SOURCE.inc(inc)
42+
}
43+
44+
def incrementDeadLetters(inc: Long): Unit = {
45+
if (metricSource.nonEmpty)
46+
metricSource.get.METRIC_DEADLETTER_ROWS_INSERTED.inc(inc)
3947
}
4048

4149
private lazy val metricSource: Option[IngestionPipelineMetricSource] = {

spark/ingestion/src/main/scala/feast/ingestion/stores/redis/RedisSinkRelation.scala

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,7 @@ class RedisSinkRelation(override val sqlContext: SQLContext, config: SparkRedisC
7171
// repartition for deduplication
7272
val dataToStore =
7373
if (config.repartitionByEntity)
74-
data.repartition(config.entityColumns.map(col): _*)
74+
data.repartition(config.entityColumns.map(col): _*).localCheckpoint()
7575
else data
7676

7777
dataToStore.foreachPartition { partition: Iterator[Row] =>

0 commit comments

Comments
 (0)