This is an automated email from the ASF dual-hosted git repository.
chengpan pushed a commit to branch branch-0.3
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git
The following commit(s) were added to refs/heads/branch-0.3 by this push:
new 7566bd20e [CELEBORN-721] Fix concurrent bug in ChangePartitionManager
7566bd20e is described below
commit 7566bd20e98973f0f68b97ea9a7e2659bf84d41d
Author: zky.zhoukeyong <[email protected]>
AuthorDate: Tue Jun 27 21:30:47 2023 +0800
[CELEBORN-721] Fix concurrent bug in ChangePartitionManager
### What changes were proposed in this pull request?
Fixes concurrent bug in ChangePartitionManager.
### Why are the changes needed?
Before this PR, ```ChangePartitionManager.start``` tries to synchronize on
```requests``` in the body
of ```run()```, but the synchronized keyword was put outside of the
```batchHandleChangePartitionExecutors.submit```,
which has no effect.
When I was testing https://github.com/apache/incubator-celeborn/pull/1588 ,
I encountered unexpected situations that
when all ```rss-lifecycle-manager-change-partition-executor``` threads are
idle, the ```inBatchPartitions``` is still not
empty:
```
23/06/27 20:35:55 INFO ChangePartitionManager: Inside run, shuffleId 0
inBatchPartitions size 834
```
### Does this PR introduce _any_ user-facing change?
No.
### How was this patch tested?
Manual test.
Closes #1634 from waitinfuture/721.
Authored-by: zky.zhoukeyong <[email protected]>
Signed-off-by: Cheng Pan <[email protected]>
(cherry picked from commit ebff17ec3c5565a8f7454d9b5b3de725594af444)
Signed-off-by: Cheng Pan <[email protected]>
---
.../scala/org/apache/celeborn/client/ChangePartitionManager.scala | 8 ++++----
1 file changed, 4 insertions(+), 4 deletions(-)
diff --git
a/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
b/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
index beb9cb70c..557cc3239 100644
---
a/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
+++
b/client/src/main/scala/org/apache/celeborn/client/ChangePartitionManager.scala
@@ -75,10 +75,10 @@ class ChangePartitionManager(
override def run(): Unit = {
try {
changePartitionRequests.asScala.foreach { case (shuffleId,
requests) =>
- requests.synchronized {
- batchHandleChangePartitionExecutors.submit {
- new Runnable {
- override def run(): Unit = {
+ batchHandleChangePartitionExecutors.submit {
+ new Runnable {
+ override def run(): Unit = {
+ requests.synchronized {
// For each partition only need handle one request
val distinctPartitions = requests.asScala.filter {
case (partitionId, _) =>
!inBatchPartitions.get(shuffleId).contains(partitionId)