Skip to content

Commit 679f242

Browse files
committed
keep same amount of partitions
Signed-off-by: Oleksii Moskalenko <moskalenko.alexey@gmail.com>
1 parent dd36544 commit 679f242

1 file changed

Lines changed: 2 additions & 2 deletions

File tree

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -70,8 +70,8 @@ class RedisSinkRelation(override val sqlContext: SQLContext, config: SparkRedisC
7070
override def insert(data: DataFrame, overwrite: Boolean): Unit = {
7171
// repartition for deduplication
7272
val dataToStore =
73-
if (config.repartitionByEntity)
74-
data.repartition(config.entityColumns.map(col): _*).localCheckpoint()
73+
if (config.repartitionByEntity && data.rdd.getNumPartitions > 1)
74+
data.repartition(data.rdd.getNumPartitions, config.entityColumns.map(col): _*).localCheckpoint()
7575
else data
7676

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

0 commit comments

Comments
 (0)