davsclaus commented on code in PR #25273:
URL: https://github.com/apache/camel/pull/25273#discussion_r3729185957
##########
components/camel-ai/camel-langchain4j-embeddingstore/src/main/java/org/apache/camel/component/langchain4j/embeddingstore/LangChain4jEmbeddingStoreProducer.java:
##########
@@ -109,63 +110,150 @@ public void process(Exchange exchange) throws Exception {
}
/**
- * Adds an embedding to the store with optional text segment.
+ * Adds embeddings to the store with optional text segments and
caller-supplied IDs.
*
* <p>
- * Expects the following headers:
+ * Supports both single and batch operations:
+ * </p>
+ *
+ * <p>
+ * <b>Single operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbedding} header contains a single
+ * {@link Embedding}:
* </p>
* <ul>
- * <li>{@code CamelLangchain4jEmbeddingEmbedding} - The embedding vector
(required)</li>
- * <li>{@code CamelLangchain4jEmbeddingTextSegment} - Associated text
segment (optional)</li>
+ * <li>With caller-supplied ID header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingId}): calls
+ * {@code add(id, embedding)}</li>
+ * <li>With text segment header: calls {@code add(embedding,
textSegment)}</li>
+ * <li>Without text segment: calls {@code add(embedding)}</li>
* </ul>
*
* <p>
- * Returns the generated embedding ID in the message body.
+ * <b>Batch operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbeddings} header contains a
+ * {@code List<Embedding>}:
* </p>
+ * <ul>
+ * <li>With IDs header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingIds}) and text segments body: calls
+ * {@code addAll(ids, embeddings, textSegments)}</li>
+ * <li>With text segments body: calls {@code addAll(embeddings,
textSegments)}</li>
+ * <li>Without text segments: calls {@code addAll(embeddings)}</li>
+ * </ul>
*
* @param exchange the Camel exchange containing the embedding data
* @throws Exception if the add operation fails
*/
+ @SuppressWarnings("unchecked")
private void add(Exchange exchange) throws Exception {
final Message in = exchange.getMessage();
+ EmbeddingStore<TextSegment> store =
getEndpoint().getConfiguration().getEmbeddingStore();
+
+ // Check for batch embeddings header first
+ List<Embedding> embeddings =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDINGS, List.class);
+ if (embeddings != null) {
+ addBatch(in, store, embeddings);
+ return;
+ }
+ // Single embedding path
if (in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING) == null) {
throw new NoSuchHeaderException(
"The embedding is a required header for ADD operations",
exchange,
LangChain4jEmbeddingsHeaders.EMBEDDING);
}
Embedding embedding =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING, Embedding.class);
+
+ // Check for caller-supplied ID
+ String callerId =
in.getHeader(LangChain4jEmbeddingStoreHeaders.EMBEDDING_ID, String.class);
String id;
- if (in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT) != null) {
+ if (callerId != null) {
+ store.add(callerId, embedding);
+ id = callerId;
+ } else if (in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT) !=
null) {
TextSegment text =
in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT, TextSegment.class);
- id =
getEndpoint().getConfiguration().getEmbeddingStore().add(embedding, text);
+ id = store.add(embedding, text);
} else {
- id =
getEndpoint().getConfiguration().getEmbeddingStore().add(embedding);
+ id = store.add(embedding);
}
- Message out = exchange.getMessage();
- out.setBody(id);
+ in.setBody(id);
+ }
+
+ @SuppressWarnings("unchecked")
+ private void addBatch(Message in, EmbeddingStore<TextSegment> store,
List<Embedding> embeddings) {
+ List<String> callerIds =
in.getHeader(LangChain4jEmbeddingStoreHeaders.EMBEDDING_IDS, List.class);
+ Object body = in.getBody();
+ List<TextSegment> textSegments = null;
+
+ if (body instanceof List && !((List<?>) body).isEmpty() && ((List<?>)
body).get(0) instanceof TextSegment) {
+ textSegments = (List<TextSegment>) body;
+ }
+
+ List<String> ids;
+ if (callerIds != null && textSegments != null) {
+ store.addAll(callerIds, embeddings, textSegments);
+ ids = callerIds;
+ } else if (callerIds != null) {
+ // No addAll(ids, embeddings) overload in langchain4j, so loop
with add(id, embedding)
+ for (int i = 0; i < embeddings.size(); i++) {
Review Comment:
**[MEDIUM] Size mismatch between callerIds and embeddings not validated**
This loop iterates `embeddings.size()` times and accesses `callerIds.get(i)`
without checking `callerIds.size() == embeddings.size()`. If the caller
provides fewer IDs than embeddings, this throws an unhelpful
`IndexOutOfBoundsException`. If more IDs than embeddings, trailing IDs are
silently ignored but returned in the response body.
Consider adding a size-equality check before the loop:
```suggestion
if (callerIds.size() != embeddings.size()) {
throw new IllegalArgumentException(
"EMBEDDING_IDS size (" + callerIds.size() + ") must
match EMBEDDINGS size (" + embeddings.size() + ")");
}
for (int i = 0; i < embeddings.size(); i++) {
```
The same validation would also be useful before the `addAll(callerIds,
embeddings, textSegments)` call above, since langchain4j's error message for
mismatched sizes is less informative.
##########
components/camel-ai/camel-langchain4j-embeddings/src/main/java/org/apache/camel/component/langchain4j/embeddings/LangChain4jEmbeddingsProducer.java:
##########
@@ -55,4 +69,30 @@ public void process(Exchange exchange) throws Exception {
message.setHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT, in);
message.setHeader(LangChain4jEmbeddingsHeaders.EMBEDDING,
result.content());
}
+
+ private void processBatch(Exchange exchange, EmbeddingModel model, Message
message, List<Object> bodyList)
+ throws Exception {
+ // Convert each element to TextSegment using the type converter
+ List<TextSegment> segments = new ArrayList<>(bodyList.size());
+ for (Object item : bodyList) {
+ TextSegment segment =
exchange.getContext().getTypeConverter().mandatoryConvertTo(TextSegment.class,
item);
+ segments.add(segment);
+ }
+
+ final Response<List<Embedding>> result = model.embedAll(segments);
+
+ if (result.finishReason() != null) {
+ message.setHeader(LangChain4jEmbeddingsHeaders.FINISH_REASON,
result.finishReason());
+ }
+
+ if (result.tokenUsage() != null) {
+ message.setHeader(LangChain4jEmbeddingsHeaders.INPUT_TOKEN_COUNT,
result.tokenUsage().inputTokenCount());
+ message.setHeader(LangChain4jEmbeddingsHeaders.OUTPUT_TOKEN_COUNT,
result.tokenUsage().outputTokenCount());
+ message.setHeader(LangChain4jEmbeddingsHeaders.TOTAL_TOKEN_COUNT,
result.tokenUsage().totalTokenCount());
+ }
+
+ List<Embedding> embeddings = result.content();
+ message.setHeader(LangChain4jEmbeddingsHeaders.EMBEDDINGS, embeddings);
+ message.setBody(embeddings);
Review Comment:
**[MEDIUM] Text segments lost in batch embed-to-store pipeline**
The body is overwritten with `List<Embedding>` and no `TEXT_SEGMENTS` header
is set to preserve the input text segments. The single path preserves the
original text via the `TEXT_SEGMENT` header, but the batch path has no
equivalent.
When chained with the embedding store producer (the exact pipeline
documented in this PR's own `langchain4j-embeddingstore-component.adoc`), text
segments are lost — the store sees `List<Embedding>` as body, not
`List<TextSegment>`, so `addBatch()` falls through to `addAll(embeddings)`
without text.
This means search results will return embeddings but not the original text,
which is a data-loss concern for RAG pipelines.
Consider preserving the segments via a `TEXT_SEGMENTS` header (or as the
body, with embeddings only in the `EMBEDDINGS` header) to mirror the
single-item behavior.
##########
components/camel-ai/camel-langchain4j-embeddingstore/src/main/java/org/apache/camel/component/langchain4j/embeddingstore/LangChain4jEmbeddingStoreProducer.java:
##########
@@ -109,63 +110,150 @@ public void process(Exchange exchange) throws Exception {
}
/**
- * Adds an embedding to the store with optional text segment.
+ * Adds embeddings to the store with optional text segments and
caller-supplied IDs.
*
* <p>
- * Expects the following headers:
+ * Supports both single and batch operations:
+ * </p>
+ *
+ * <p>
+ * <b>Single operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbedding} header contains a single
+ * {@link Embedding}:
* </p>
* <ul>
- * <li>{@code CamelLangchain4jEmbeddingEmbedding} - The embedding vector
(required)</li>
- * <li>{@code CamelLangchain4jEmbeddingTextSegment} - Associated text
segment (optional)</li>
+ * <li>With caller-supplied ID header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingId}): calls
+ * {@code add(id, embedding)}</li>
+ * <li>With text segment header: calls {@code add(embedding,
textSegment)}</li>
+ * <li>Without text segment: calls {@code add(embedding)}</li>
* </ul>
*
* <p>
- * Returns the generated embedding ID in the message body.
+ * <b>Batch operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbeddings} header contains a
+ * {@code List<Embedding>}:
* </p>
+ * <ul>
+ * <li>With IDs header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingIds}) and text segments body: calls
+ * {@code addAll(ids, embeddings, textSegments)}</li>
+ * <li>With text segments body: calls {@code addAll(embeddings,
textSegments)}</li>
+ * <li>Without text segments: calls {@code addAll(embeddings)}</li>
+ * </ul>
*
* @param exchange the Camel exchange containing the embedding data
* @throws Exception if the add operation fails
*/
+ @SuppressWarnings("unchecked")
private void add(Exchange exchange) throws Exception {
final Message in = exchange.getMessage();
+ EmbeddingStore<TextSegment> store =
getEndpoint().getConfiguration().getEmbeddingStore();
+
+ // Check for batch embeddings header first
+ List<Embedding> embeddings =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDINGS, List.class);
+ if (embeddings != null) {
+ addBatch(in, store, embeddings);
+ return;
+ }
+ // Single embedding path
if (in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING) == null) {
throw new NoSuchHeaderException(
"The embedding is a required header for ADD operations",
exchange,
LangChain4jEmbeddingsHeaders.EMBEDDING);
}
Embedding embedding =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING, Embedding.class);
+
+ // Check for caller-supplied ID
+ String callerId =
in.getHeader(LangChain4jEmbeddingStoreHeaders.EMBEDDING_ID, String.class);
String id;
- if (in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT) != null) {
+ if (callerId != null) {
+ store.add(callerId, embedding);
Review Comment:
**[LOW] Caller-supplied ID path silently drops text segment on single add**
When `callerId` is set, `store.add(callerId, embedding)` is called and the
text-segment check is skipped entirely. If a `TEXT_SEGMENT` header is also
present, it's silently lost.
The langchain4j API has no `add(String id, Embedding, TextSegment)`
overload, but `addAll(List<String>, List<Embedding>, List<TextSegment>)` exists
and could be used with singleton lists to preserve the text segment. Worth
documenting this limitation at minimum.
##########
components/camel-ai/camel-langchain4j-embeddingstore/src/main/java/org/apache/camel/component/langchain4j/embeddingstore/LangChain4jEmbeddingStoreProducer.java:
##########
@@ -109,63 +110,150 @@ public void process(Exchange exchange) throws Exception {
}
/**
- * Adds an embedding to the store with optional text segment.
+ * Adds embeddings to the store with optional text segments and
caller-supplied IDs.
*
* <p>
- * Expects the following headers:
+ * Supports both single and batch operations:
+ * </p>
+ *
+ * <p>
+ * <b>Single operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbedding} header contains a single
+ * {@link Embedding}:
* </p>
* <ul>
- * <li>{@code CamelLangchain4jEmbeddingEmbedding} - The embedding vector
(required)</li>
- * <li>{@code CamelLangchain4jEmbeddingTextSegment} - Associated text
segment (optional)</li>
+ * <li>With caller-supplied ID header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingId}): calls
+ * {@code add(id, embedding)}</li>
+ * <li>With text segment header: calls {@code add(embedding,
textSegment)}</li>
+ * <li>Without text segment: calls {@code add(embedding)}</li>
* </ul>
*
* <p>
- * Returns the generated embedding ID in the message body.
+ * <b>Batch operation</b> - when the {@code
CamelLangChain4jEmbeddingsEmbeddings} header contains a
+ * {@code List<Embedding>}:
* </p>
+ * <ul>
+ * <li>With IDs header ({@code
CamelLangchain4jEmbeddingStoreEmbeddingIds}) and text segments body: calls
+ * {@code addAll(ids, embeddings, textSegments)}</li>
+ * <li>With text segments body: calls {@code addAll(embeddings,
textSegments)}</li>
+ * <li>Without text segments: calls {@code addAll(embeddings)}</li>
+ * </ul>
*
* @param exchange the Camel exchange containing the embedding data
* @throws Exception if the add operation fails
*/
+ @SuppressWarnings("unchecked")
private void add(Exchange exchange) throws Exception {
final Message in = exchange.getMessage();
+ EmbeddingStore<TextSegment> store =
getEndpoint().getConfiguration().getEmbeddingStore();
+
+ // Check for batch embeddings header first
+ List<Embedding> embeddings =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDINGS, List.class);
+ if (embeddings != null) {
+ addBatch(in, store, embeddings);
+ return;
+ }
+ // Single embedding path
if (in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING) == null) {
throw new NoSuchHeaderException(
"The embedding is a required header for ADD operations",
exchange,
LangChain4jEmbeddingsHeaders.EMBEDDING);
}
Embedding embedding =
in.getHeader(LangChain4jEmbeddingsHeaders.EMBEDDING, Embedding.class);
+
+ // Check for caller-supplied ID
+ String callerId =
in.getHeader(LangChain4jEmbeddingStoreHeaders.EMBEDDING_ID, String.class);
String id;
- if (in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT) != null) {
+ if (callerId != null) {
+ store.add(callerId, embedding);
+ id = callerId;
+ } else if (in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT) !=
null) {
TextSegment text =
in.getHeader(LangChain4jEmbeddingsHeaders.TEXT_SEGMENT, TextSegment.class);
- id =
getEndpoint().getConfiguration().getEmbeddingStore().add(embedding, text);
+ id = store.add(embedding, text);
} else {
- id =
getEndpoint().getConfiguration().getEmbeddingStore().add(embedding);
+ id = store.add(embedding);
}
- Message out = exchange.getMessage();
- out.setBody(id);
+ in.setBody(id);
+ }
+
+ @SuppressWarnings("unchecked")
+ private void addBatch(Message in, EmbeddingStore<TextSegment> store,
List<Embedding> embeddings) {
+ List<String> callerIds =
in.getHeader(LangChain4jEmbeddingStoreHeaders.EMBEDDING_IDS, List.class);
+ Object body = in.getBody();
+ List<TextSegment> textSegments = null;
+
+ if (body instanceof List && !((List<?>) body).isEmpty() && ((List<?>)
body).get(0) instanceof TextSegment) {
+ textSegments = (List<TextSegment>) body;
+ }
+
+ List<String> ids;
+ if (callerIds != null && textSegments != null) {
+ store.addAll(callerIds, embeddings, textSegments);
+ ids = callerIds;
+ } else if (callerIds != null) {
+ // No addAll(ids, embeddings) overload in langchain4j, so loop
with add(id, embedding)
+ for (int i = 0; i < embeddings.size(); i++) {
+ store.add(callerIds.get(i), embeddings.get(i));
+ }
+ ids = callerIds;
+ } else if (textSegments != null) {
+ ids = store.addAll(embeddings, textSegments);
+ } else {
+ ids = store.addAll(embeddings);
+ }
+
+ in.setBody(ids);
}
/**
- * Removes an embedding from the store by its ID.
+ * Removes embeddings from the store. Supports multiple removal strategies:
*
- * <p>
- * Expects the embedding ID as the message body (String).
- * </p>
+ * <ul>
+ * <li><b>By filter</b>: when the {@code
CamelLangchain4jEmbeddingStoreFilter} header is set, removes all embeddings
+ * matching the filter via {@code removeAll(Filter)}</li>
+ * <li><b>By ID list</b>: when the body is a {@code Collection<String>},
removes all specified embeddings via
+ * {@code removeAll(Collection)}</li>
+ * <li><b>By single ID</b>: when the body is a single {@code String},
removes that embedding via
+ * {@code remove(id)}</li>
+ * </ul>
*
- * @param exchange the Camel exchange containing the embedding ID to
remove
+ * @param exchange the Camel exchange containing removal parameters
* @throws Exception if the remove operation fails
*/
+ @SuppressWarnings("unchecked")
private void remove(Exchange exchange) throws Exception {
final Message in = exchange.getMessage();
- String id = in.getBody(String.class);
+ EmbeddingStore<TextSegment> store =
getEndpoint().getConfiguration().getEmbeddingStore();
- getEndpoint().getConfiguration().getEmbeddingStore().remove(id);
+ // Check for filter-based removal first
+ Filter filter = in.getHeader(LangChain4jEmbeddingStoreHeaders.FILTER,
Filter.class);
+ if (filter != null) {
+ store.removeAll(filter);
+ return;
+ }
- Message out = exchange.getMessage();
+ Object body = in.getBody();
+
+ // Batch removal by collection of IDs
+ if (body instanceof Collection) {
+ Collection<String> ids = (Collection<String>) body;
+ store.removeAll(ids);
+ return;
+ }
+
+ // Single ID removal
+ String id = in.getBody(String.class);
+ if (id != null && !id.isEmpty()) {
+ store.remove(id);
+ return;
+ }
+
+ throw new IllegalArgumentException(
Review Comment:
**[LOW] Missing upgrade guide entry for REMOVE behavior change**
This `IllegalArgumentException` is an improvement over the previous behavior
(passing null to `store.remove(null)`), but it's a behavioral change that could
affect existing routes.
Per project conventions, behavioral changes should be documented in
`docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc`.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]