Skip to content

Commit 191f0ae

Browse files
committed
pull config to RedisSink
1 parent bda8fb4 commit 191f0ae

2 files changed

Lines changed: 17 additions & 20 deletions

File tree

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

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -46,9 +46,6 @@
4646

4747
public class RedisCustomIO {
4848

49-
private static final int DEFAULT_BATCH_SIZE = 10000;
50-
private static final int DEFAULT_FREQUENCY_SECONDS = 30;
51-
5249
private static TupleTag<FeatureRow> successfulInsertsTag =
5350
new TupleTag<FeatureRow>("successfulInserts") {};
5451
private static TupleTag<FailedElement> failedInsertsTupleTag =
@@ -69,8 +66,8 @@ public static class Write extends PTransform<PCollection<FeatureRow>, WriteResul
6966

7067
private PCollectionView<Map<String, Iterable<FeatureSetSpec>>> featureSetSpecs;
7168
private RedisIngestionClient redisIngestionClient;
72-
private int batchSize = DEFAULT_BATCH_SIZE;
73-
private Duration flushFrequency = Duration.standardSeconds(DEFAULT_FREQUENCY_SECONDS);
69+
private int batchSize;
70+
private Duration flushFrequency;
7471

7572
public Write(
7673
RedisIngestionClient redisIngestionClient,

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

Lines changed: 15 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,8 @@
3939

4040
@AutoValue
4141
public abstract class RedisFeatureSink implements FeatureSink {
42+
private static final int DEFAULT_BATCH_SIZE = 10000;
43+
private static final int DEFAULT_FREQUENCY_SECONDS = 30;
4244

4345
/**
4446
* Initialize a {@link RedisFeatureSink.Builder} from a {@link StoreProto.Store.RedisConfig}.
@@ -113,30 +115,28 @@ public PCollection<FeatureSetReference> prepareWrite(
113115

114116
@Override
115117
public PTransform<PCollection<FeatureRow>, WriteResult> writer() {
118+
int flushFrequencySeconds = DEFAULT_FREQUENCY_SECONDS;
119+
116120
if (getRedisClusterConfig() != null) {
117-
RedisCustomIO.Write writer =
118-
new RedisCustomIO.Write(
119-
new RedisClusterIngestionClient(getRedisClusterConfig()), getSpecsView());
120121

121122
if (getRedisClusterConfig().getFlushFrequencySeconds() > 0) {
122-
writer =
123-
writer.withFlushFrequency(
124-
Duration.standardSeconds(getRedisClusterConfig().getFlushFrequencySeconds()));
123+
flushFrequencySeconds = getRedisClusterConfig().getFlushFrequencySeconds();
125124
}
126125

127-
return writer;
128-
} else if (getRedisConfig() != null) {
129-
RedisCustomIO.Write writer =
130-
new RedisCustomIO.Write(
131-
new RedisStandaloneIngestionClient(getRedisConfig()), getSpecsView());
126+
return new RedisCustomIO.Write(
127+
new RedisClusterIngestionClient(getRedisClusterConfig()), getSpecsView())
128+
.withFlushFrequency(Duration.standardSeconds(flushFrequencySeconds))
129+
.withBatchSize(DEFAULT_BATCH_SIZE);
132130

131+
} else if (getRedisConfig() != null) {
133132
if (getRedisConfig().getFlushFrequencySeconds() > 0) {
134-
writer =
135-
writer.withFlushFrequency(
136-
Duration.standardSeconds(getRedisConfig().getFlushFrequencySeconds()));
133+
flushFrequencySeconds = getRedisConfig().getFlushFrequencySeconds();
137134
}
138135

139-
return writer;
136+
return new RedisCustomIO.Write(
137+
new RedisStandaloneIngestionClient(getRedisConfig()), getSpecsView())
138+
.withFlushFrequency(Duration.standardSeconds(flushFrequencySeconds))
139+
.withBatchSize(DEFAULT_BATCH_SIZE);
140140
} else {
141141
throw new RuntimeException(
142142
"At least one RedisConfig or RedisClusterConfig must be provided to Redis Sink");

0 commit comments

Comments
 (0)