diff --git a/core/src/main/java/feast/core/config/FeatureStreamConfig.java b/core/src/main/java/feast/core/config/FeatureStreamConfig.java index 32445456048..795e4f754ff 100644 --- a/core/src/main/java/feast/core/config/FeatureStreamConfig.java +++ b/core/src/main/java/feast/core/config/FeatureStreamConfig.java @@ -45,7 +45,7 @@ public class FeatureStreamConfig { String DEFAULT_KAFKA_REQUEST_TIMEOUT_MS_CONFIG = "15000"; int DEFAULT_SPECS_TOPIC_PARTITIONING = 1; - short DEFAULT_SPECS_TOPIC_REPLICATION = 3; + short DEFAULT_SPECS_TOPIC_REPLICATION = 1; @Bean public KafkaAdmin admin(FeastProperties feastProperties) { diff --git a/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/writer/RedisFeatureSink.java b/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/writer/RedisFeatureSink.java index 6997cbcb877..4d2ebf14d6d 100644 --- a/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/writer/RedisFeatureSink.java +++ b/storage/connectors/redis/src/main/java/feast/storage/connectors/redis/writer/RedisFeatureSink.java @@ -110,7 +110,20 @@ public PCollection prepareWrite( "At least one RedisConfig or RedisClusterConfig must be provided to Redis Sink"); } specsView = featureSetSpecs.apply(ParDo.of(new ReferenceToString())).apply(View.asMultimap()); - return featureSetSpecs.apply(Keys.create()); + return featureSetSpecs + .apply( + "DummyDelay", + ParDo.of( + new DoFn< + KV, + KV>() { + @ProcessElement + public void process(ProcessContext c) throws InterruptedException { + Thread.sleep(1000); + c.output(c.element()); + } + })) + .apply(Keys.create()); } @Override