From 79aabbf1fe9fc85a54836fa2e944e7f7cdcdd771 Mon Sep 17 00:00:00 2001 From: pyalex Date: Mon, 13 Jul 2020 14:25:35 +0300 Subject: [PATCH] ingestion spec update: should take some time before ack --- .../feast/core/config/FeatureStreamConfig.java | 2 +- .../connectors/redis/writer/RedisFeatureSink.java | 15 ++++++++++++++- 2 files changed, 15 insertions(+), 2 deletions(-) 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