davsclaus commented on code in PR #27638:
URL: https://github.com/apache/camel/pull/27638#discussion_r4236604156
##########
components/camel-lucene/src/main/java/org/apache/camel/component/lucene/LuceneIndexer.java:
##########
@@ -70,18 +70,25 @@ public LuceneIndexer(File sourceDirectory, File
indexDirectory, Analyzer analyze
}
}
- public void index(Exchange exchange) throws Exception {
+ public synchronized void index(Exchange exchange) throws Exception {
+ // synchronized: the index writer is a field and Lucene allows one
open writer per index directory
LOG.debug("Indexing {}", exchange);
openIndexWriter();
- Map<String, Object> headers = exchange.getIn().getHeaders();
- add("exchangeId", exchange.getExchangeId(), true);
- for (Entry<String, Object> entry : headers.entrySet()) {
- String field = entry.getKey();
- String value =
exchange.getContext().getTypeConverter().mandatoryConvertTo(String.class,
entry.getValue());
- add(field, value, true);
- }
+ try {
+ Map<String, Object> headers = exchange.getIn().getHeaders();
+ add("exchangeId", exchange.getExchangeId(), true);
+ for (Entry<String, Object> entry : headers.entrySet()) {
+ String field = entry.getKey();
+ String value =
exchange.getContext().getTypeConverter().mandatoryConvertTo(String.class,
entry.getValue());
+ add(field, value, true);
+ }
- add("contents", exchange.getIn().getMandatoryBody(String.class), true);
+ add("contents", exchange.getIn().getMandatoryBody(String.class),
true);
+ } catch (Exception e) {
+ // discard the documents of this exchange and release the index
write lock
+ indexWriter.rollback();
+ throw e;
+ }
closeIndexWriter();
Review Comment:
`closeIndexWriter()` is still outside the `try`, so if
`indexWriter.commit()` throws (for example an IO error on the index files) the
writer stays open, and the write lock stays held, the same as before.
`IndexWriter.close()` always closes, even on failure, but `commit()` does not.
Moving the commit into the `try` covers that path too. A `rollback()` on a
writer that is already closed is a no-op. Also consider not letting a failing
rollback hide the original exception:
```java
add("contents", exchange.getIn().getMandatoryBody(String.class),
true);
indexWriter.commit();
} catch (Exception e) {
try {
indexWriter.rollback();
} catch (Exception re) {
e.addSuppressed(re);
}
throw e;
}
indexWriter.close();
```
Non-blocking.
##########
components/camel-lucene/src/main/java/org/apache/camel/component/lucene/LuceneEndpoint.java:
##########
@@ -81,6 +81,16 @@ public Producer createProducer() throws Exception {
return new LuceneIndexProducer(this, this.config, indexer);
}
+ @Override
+ protected void doShutdown() throws Exception {
+ // the index directory is shared by all producers of this endpoint, so
it is closed with the endpoint
+ // and not when a producer (route) stops
+ if (indexer != null) {
+ indexer.getNiofsDirectory().close();
Review Comment:
FYI, no change needed: `removeEndpoint()` (used when a route is removed)
only *stops* the endpoint, so this `doShutdown()` does not run there. That is
harmless, because in Lucene 9 `FSDirectory.close()` only flips `isOpen` and
deletes pending files. The file handles belong to the writer and readers. After
a `CamelContext` stop/start in place, the routes keep this endpoint instance
with a closed directory, as on main today, so this is not a regression. If you
want to cover that later, open the directory in `doStart()` and close it in
`doStop()` of the endpoint, since route restarts do not stop the endpoint.
--
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]