Skip to content

Commit bc722c8

Browse files
authored
Merge pull request #304 from jmartisk/redis-batching
Use batching in Redis to avoid overflowing the connection pool
2 parents 62374e9 + 44e8a8c commit bc722c8

File tree

1 file changed

+4
-5
lines changed

1 file changed

+4
-5
lines changed

redis/runtime/src/main/java/io/quarkiverse/langchain4j/redis/RedisEmbeddingStore.java

Lines changed: 4 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -5,6 +5,7 @@
55
import static java.util.Collections.singletonList;
66
import static java.util.stream.Collectors.toList;
77

8+
import java.util.ArrayList;
89
import java.util.HashMap;
910
import java.util.List;
1011
import java.util.Map;
@@ -25,7 +26,6 @@
2526
import io.quarkiverse.langchain4j.QuarkusJsonCodecFactory;
2627
import io.quarkiverse.langchain4j.redis.runtime.RedisSchema;
2728
import io.quarkus.redis.datasource.ReactiveRedisDataSource;
28-
import io.quarkus.redis.datasource.json.ReactiveJsonCommands;
2929
import io.quarkus.redis.datasource.keys.KeyScanArgs;
3030
import io.quarkus.redis.datasource.search.CreateArgs;
3131
import io.quarkus.redis.datasource.search.Document;
@@ -125,9 +125,8 @@ private void addAllInternal(List<String> ids, List<Embedding> embeddings, List<T
125125
if (ids.isEmpty() || ids.size() != embeddings.size() || (embedded != null && embedded.size() != embeddings.size())) {
126126
throw new IllegalArgumentException("ids, embeddings and embedded must be non-empty and of the same size");
127127
}
128-
ReactiveJsonCommands<String> json = ds.json();
129128
int size = ids.size();
130-
Uni[] unis = new Uni[size];
129+
List<Request> commands = new ArrayList<>();
131130
for (int i = 0; i < size; i++) {
132131
String id = ids.get(i);
133132
Embedding embedding = embeddings.get(i);
@@ -147,9 +146,9 @@ private void addAllInternal(List<String> ids, List<Embedding> embeddings, List<T
147146
fields.putAll(textSegment.metadata().asMap());
148147
}
149148
String key = schema.getPrefix() + id;
150-
unis[i] = json.jsonSet(key, "$", fields);
149+
commands.add(Request.cmd(Command.JSON_SET).arg(key).arg("$").arg(Json.toJson(fields)));
151150
}
152-
Uni.join().all(unis).andFailFast().await().indefinitely();
151+
ds.getRedis().batchAndAwait(commands);
153152
}
154153

155154
@Override

0 commit comments

Comments
 (0)