This is an automated email from the ASF dual-hosted git repository.

zhouky pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/incubator-celeborn.git


The following commit(s) were added to refs/heads/main by this push:
     new a1199a989 [CELEBORN-728] Celeborn won't clean remnant application 
directory on HDFS if worker is restarted
a1199a989 is described below

commit a1199a98954fa3c7fc513c57a8e7d91709550547
Author: Demon Liang <[email protected]>
AuthorDate: Wed Jun 28 17:53:50 2023 +0800

    [CELEBORN-728] Celeborn won't clean remnant application directory on HDFS 
if worker is restarted
    
    ### What changes were proposed in this pull request?
    To clean the remnant application directory after Celeborn Worker is 
restarted.
    
    ### Why are the changes needed?
    Remnant application directories will not be deleted, because 
`hadoopFs.listFiles(path,false)` will not list directories.
    
    ### Does this PR introduce _any_ user-facing change?
    No.
    
    Closes #1641 from Demon-Liang/0.3-dev.
    
    Authored-by: Demon Liang <[email protected]>
    Signed-off-by: zky.zhoukeyong <[email protected]>
    (cherry picked from commit 42a9160c8ceaf79bae514c54dafcb5b8e12d5251)
    Signed-off-by: zky.zhoukeyong <[email protected]>
---
 .../apache/celeborn/service/deploy/worker/storage/StorageManager.scala  | 2 +-
 1 file changed, 1 insertion(+), 1 deletion(-)

diff --git 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
index ae29d79df..09a0d832c 100644
--- 
a/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
+++ 
b/worker/src/main/scala/org/apache/celeborn/service/deploy/worker/storage/StorageManager.scala
@@ -531,7 +531,7 @@ final private[worker] class StorageManager(conf: 
CelebornConf, workerSource: Abs
     if (hadoopFs != null) {
       val hdfsWorkPath = new Path(hdfsDir, conf.workerWorkingDir)
       if (hadoopFs.exists(hdfsWorkPath)) {
-        val iter = hadoopFs.listFiles(hdfsWorkPath, false)
+        val iter = hadoopFs.listStatusIterator(hdfsWorkPath)
         while (iter.hasNext) {
           val fileStatus = iter.next()
           if (!appIds.contains(fileStatus.getPath.getName)) {

Reply via email to