He-Pin opened a new pull request, #3439:
URL: https://github.com/apache/pekko/pull/3439

   ### Motivation
   `drainQueue()` is called on every future completion and performs an O(n) 
`buffer.toList` snapshot plus iteration. When all parallelism slots are 
occupied (`partitionsInProgress.size >= parallelism`), `canStartNextElement` 
returns false for every element, making the entire snapshot and iteration 
wasted work. This is the common steady-state case under backpressure.
   
   ### Modification
   Add a guard `partitionsInProgress.size < parallelism` to the `drainQueue` 
condition, skipping the O(n) allocation and traversal when no new element can 
possibly be started.
   
   ### Result
   In steady state (all parallelism slots busy), each future completion avoids 
an unnecessary O(n) buffer snapshot. The benefit scales with buffer size and 
parallelism.
   
   ### Tests
   - `sbt "stream-tests / Test / testOnly 
org.apache.pekko.stream.scaladsl.FlowMapAsyncPartitionedSpec"`
   - `sbt "stream-typed-tests / Test / testOnly 
org.apache.pekko.stream.MapAsyncPartitionedSpec"`
   
   ### References
   Refs akka/akka-core#32031


-- 
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]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to