jnioche commented on issue #2212:
URL: https://github.com/apache/stormcrawler/issues/2212#issuecomment-6016982902
Using the new StatusUpdaterBenchmark class on both modules shows unexpected
results. I was expecting the legacy implementation to be slower compared to the
new one but it turned out to be the other way round.
```
docker run -d --rm --name opensearch-bench -p 127.0.0.1:9201:9200 -e
discovery.type=single-node -e DISABLE_SECURITY_PLUGIN=true -e
bootstrap.memory_lock=true -e "OPENSEARCH_JAVA_OPTS=-Xms4g -Xmx4g" --ulimit
memlock=-1:-1 --ulimit nofile=65536:65536 opensearchproject/opensearch:3.9.0
storm local target/benchmark-0.1-SNAPSHOT.jar
org.apache.stormcrawler.persistence.StatusUpdaterBenchmark
org.apache.stormcrawler.opensearch.persistence.StatusUpdaterBolt 4
opensearch-conf.yaml outlinks.ndjson
```
where outlinks.ndjson is a massive file at the URLFRontier import format.
with
```
opensearch.status.bulkActions: 500
opensearch.status.flushInterval: "2s"
opensearch.status.concurrentRequests: 8
```
gave
`opensearch` module : 109K OPS
`opensearch-java` module : 95K OPS
in case you were wondering `URLFrontier`: 174K OPS
+ DigitalPebble's closed-source URLFrontier-compatible backend: 385K OPS
Claude suggested that the
```
the most likely cause is how the new AsyncBulkProcessor frees up a slot for
the next bulk request. It makes the bolt's own thread send bulk requests
itself, which the legacy processor never does.
**Main suspect**: the bolt thread ends up making the bulk HTTP call
In AsyncBulkProcessor.java:
The thread pool is ThreadPoolExecutor(1, concurrentRequests, …, new
SynchronousQueue<>(), …, CallerRunsPolicy). With the default
concurrentRequests=1, it has one worker thread and no queue.
client.bulk(request) is a blocking call that runs on that worker thread.
In whenComplete, the permit is released first, and only then does
listener.afterBulk(...) run on the same worker thread. afterBulk is the
expensive part: it calls WaitAckCache.processBulkResponse, takes the fair lock,
and acks or fails up to bulkActions tuples.
Here is what happens:
Bulk N comes back. The worker releases the permit and starts running
afterBulk for bulk N.
Meanwhile the bolt thread has filled the next buffer. It acquires the permit
straight away and submits bulk N+1.
The only worker is still busy in afterBulk, and a SynchronousQueue has
nowhere to park the task, so the pool rejects it. CallerRunsPolicy then runs
client.bulk() on the bolt's executor thread.
The bolt is now blocked for the whole round trip of bulk N+1, so it can't
call store() or buffer the next batch.
The legacy BulkProcessor behaves differently:
It uses bulkAsync, so the bolt thread only hands off the request and never
waits for the HTTP call.
It releases the semaphore in a finally block after the listener has run, so
this race can't happen.
Under load, the legacy processor fills batch N+1 while batch N is in flight.
The new one often fills a batch and then waits a full round trip before it can
start the next one. That roughly costs you one round trip per batch, and it
gets worse when afterBulk is slow. It would also match a result where running
more updater instances makes the gap look smaller.
Fix options:
Simplest: release the permit after the listener has run, or make the pool's
maximum size concurrentRequests + 1. Either way the bolt never runs a bulk call
itself.
Cleaner: use OpenSearchAsyncClient.bulk(), which returns a
CompletableFuture, and drop the executor entirely. The behaviour would then
match the legacy bulkAsync.
```
--
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]