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]

Reply via email to