Skip to content

Commit ad67392

Browse files
authored
Features are not being ingested due to max age overflow (#1209)
Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com>
1 parent 3826bd9 commit ad67392

5 files changed

Lines changed: 7 additions & 7 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,7 +72,7 @@ object BatchPipeline extends BasePipeline {
7272
.option("namespace", featureTable.name)
7373
.option("project_name", featureTable.project)
7474
.option("timestamp_column", config.source.eventTimestampColumn)
75-
.option("max_age", config.featureTable.maxAge.getOrElse(0))
75+
.option("max_age", config.featureTable.maxAge.getOrElse(0L))
7676
.save()
7777

7878
config.deadLetterPath match {

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -95,7 +95,7 @@ case class FeatureTable(
9595
project: String,
9696
entities: Seq[Field],
9797
features: Seq[Field],
98-
maxAge: Option[Int] = None
98+
maxAge: Option[Long] = None
9999
)
100100

101101
case class IngestionJobConfig(

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -87,7 +87,7 @@ object StreamingPipeline extends BasePipeline with Serializable {
8787
.option("namespace", featureTable.name)
8888
.option("project_name", featureTable.project)
8989
.option("timestamp_column", config.source.eventTimestampColumn)
90-
.option("max_age", config.featureTable.maxAge.getOrElse(0))
90+
.option("max_age", config.featureTable.maxAge.getOrElse(0L))
9191
.save()
9292

9393
config.deadLetterPath match {

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -24,7 +24,7 @@ case class SparkRedisConfig(
2424
iteratorGroupingSize: Int = 1000,
2525
timestampPrefix: String = "_ts",
2626
repartitionByEntity: Boolean = true,
27-
maxAge: Int = 0,
27+
maxAge: Long = 0,
2828
expiryPrefix: String = "_ex"
2929
)
3030

@@ -43,6 +43,6 @@ object SparkRedisConfig {
4343
entityColumns = parameters.getOrElse(ENTITY_COLUMNS, "").split(","),
4444
timestampColumn = parameters.getOrElse(TS_COLUMN, "event_timestamp"),
4545
repartitionByEntity = parameters.getOrElse(ENTITY_REPARTITION, "true") == "true",
46-
maxAge = parameters.get(MAX_AGE).map(_.toInt).getOrElse(0)
46+
maxAge = parameters.get(MAX_AGE).map(_.toLong).getOrElse(0)
4747
)
4848
}

spark/ingestion/src/test/scala/feast/ingestion/BatchPipelineIT.scala

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -131,7 +131,7 @@ class BatchPipelineIT extends SparkSpec with ForAllTestContainer {
131131
val gen = rowGenerator(startDate, endDate)
132132
val rows = generateDistinctRows(gen, 1000, groupByEntity)
133133
val tempPath = storeAsParquet(sparkSession, rows)
134-
val maxAge = 86400 * 2
134+
val maxAge = 86400L * 30
135135
val configWithMaxAge = config.copy(
136136
source = FileSource(tempPath, Map.empty, "eventTimestamp"),
137137
featureTable = config.featureTable.copy(maxAge = Some(maxAge)),
@@ -162,7 +162,7 @@ class BatchPipelineIT extends SparkSpec with ForAllTestContainer {
162162

163163
})
164164

165-
val increasedMaxAge = 86400 * 3
165+
val increasedMaxAge = 86400L * 60
166166
val configWithSecondFeatureTable = config.copy(
167167
source = FileSource(tempPath, Map.empty, "eventTimestamp"),
168168
featureTable = config.featureTable.copy(

0 commit comments

Comments
 (0)