1616 */
1717package feast .storage .connectors .redis .retriever ;
1818
19+ import com .google .common .collect .ImmutableMap ;
20+ import feast .proto .core .StoreProto ;
21+ import feast .proto .core .StoreProto .Store .RedisClusterConfig ;
1922import feast .storage .connectors .redis .serializer .RedisKeyPrefixSerializerV2 ;
2023import feast .storage .connectors .redis .serializer .RedisKeySerializerV2 ;
2124import io .lettuce .core .KeyValue ;
25+ import io .lettuce .core .ReadFrom ;
2226import io .lettuce .core .RedisFuture ;
2327import io .lettuce .core .RedisURI ;
2428import io .lettuce .core .cluster .api .StatefulRedisClusterConnection ;
@@ -36,6 +40,13 @@ public class RedisClusterClient implements RedisClientAdapter {
3640 private final RedisKeySerializerV2 serializer ;
3741 @ Nullable private final RedisKeySerializerV2 fallbackSerializer ;
3842
43+ private static final Map <RedisClusterConfig .ReadFrom , ReadFrom > PROTO_TO_LETTUCE_TYPES =
44+ ImmutableMap .of (
45+ RedisClusterConfig .ReadFrom .MASTER , ReadFrom .MASTER ,
46+ RedisClusterConfig .ReadFrom .MASTER_PREFERRED , ReadFrom .MASTER_PREFERRED ,
47+ RedisClusterConfig .ReadFrom .REPLICA , ReadFrom .REPLICA ,
48+ RedisClusterConfig .ReadFrom .REPLICA_PREFERRED , ReadFrom .REPLICA_PREFERRED );
49+
3950 @ Override
4051 public RedisFuture <List <KeyValue <byte [], byte []>>> hmget (byte [] key , byte []... fields ) {
4152 return asyncCommands .hmget (key , fields );
@@ -73,13 +84,16 @@ private RedisClusterClient(Builder builder) {
7384 this .serializer = builder .serializer ;
7485 this .fallbackSerializer = builder .fallbackSerializer ;
7586
87+ // allows reading from replicas
88+ this .asyncCommands .readOnly ();
89+
7690 // Disable auto-flushing
7791 this .asyncCommands .setAutoFlushCommands (false );
7892 }
7993
80- public static RedisClientAdapter create (Map < String , String > config ) {
94+ public static RedisClientAdapter create (StoreProto . Store . RedisClusterConfig config ) {
8195 List <RedisURI > redisURIList =
82- Arrays .stream (config .get ( "connection_string" ).split ("," ))
96+ Arrays .stream (config .getConnectionString ( ).split ("," ))
8397 .map (
8498 hostPort -> {
8599 String [] hostPortSplit = hostPort .trim ().split (":" );
@@ -90,14 +104,15 @@ public static RedisClientAdapter create(Map<String, String> config) {
90104 io .lettuce .core .cluster .RedisClusterClient .create (redisURIList )
91105 .connect (new ByteArrayCodec ());
92106
93- RedisKeySerializerV2 serializer =
94- new RedisKeyPrefixSerializerV2 (config .getOrDefault ("key_prefix" , "" ));
107+ connection .setReadFrom (PROTO_TO_LETTUCE_TYPES .get (config .getReadFrom ()));
108+
109+ RedisKeySerializerV2 serializer = new RedisKeyPrefixSerializerV2 (config .getKeyPrefix ());
95110
96111 Builder builder = new Builder (connection , serializer );
97112
98- if (Boolean . parseBoolean ( config .getOrDefault ( "enable_fallback" , "false" ) )) {
113+ if (config .getEnableFallback ( )) {
99114 RedisKeySerializerV2 fallbackSerializer =
100- new RedisKeyPrefixSerializerV2 (config .getOrDefault ( "fallback_prefix" , "" ));
115+ new RedisKeyPrefixSerializerV2 (config .getKeyPrefix ( ));
101116 builder = builder .withFallbackSerializer (fallbackSerializer );
102117 }
103118
0 commit comments