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]
