Skip to content

Commit 348cdfd

Browse files
wooppyalex
andauthored
Backport delay in Redis acknowledgement of spec (#915)
Co-authored-by: pyalex <moskalenko.alexey@gmail.com>
1 parent 565ba46 commit 348cdfd

2 files changed

Lines changed: 15 additions & 2 deletions

File tree

core/src/main/java/feast/core/config/FeatureStreamConfig.java

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ public class FeatureStreamConfig {
4545

4646
String DEFAULT_KAFKA_REQUEST_TIMEOUT_MS_CONFIG = "15000";
4747
int DEFAULT_SPECS_TOPIC_PARTITIONING = 1;
48-
short DEFAULT_SPECS_TOPIC_REPLICATION = 3;
48+
short DEFAULT_SPECS_TOPIC_REPLICATION = 1;
4949

5050
@Bean
5151
public KafkaAdmin admin(FeastProperties feastProperties) {

storage/connectors/redis/src/main/java/feast/storage/connectors/redis/writer/RedisFeatureSink.java

Lines changed: 14 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -110,7 +110,20 @@ public PCollection<FeatureSetReference> prepareWrite(
110110
"At least one RedisConfig or RedisClusterConfig must be provided to Redis Sink");
111111
}
112112
specsView = featureSetSpecs.apply(ParDo.of(new ReferenceToString())).apply(View.asMultimap());
113-
return featureSetSpecs.apply(Keys.create());
113+
return featureSetSpecs
114+
.apply(
115+
"DummyDelay",
116+
ParDo.of(
117+
new DoFn<
118+
KV<FeatureSetReference, FeatureSetSpec>,
119+
KV<FeatureSetReference, FeatureSetSpec>>() {
120+
@ProcessElement
121+
public void process(ProcessContext c) throws InterruptedException {
122+
Thread.sleep(1000);
123+
c.output(c.element());
124+
}
125+
}))
126+
.apply(Keys.create());
114127
}
115128

116129
@Override

0 commit comments

Comments
 (0)