【问题标题】:Wait for consumer to finish task before others can start等待消费者完成任务,然后其他人才能开始
【发布时间】:2017-03-22 01:53:41
【问题描述】:

我有一个 Java 应用程序,它遍历树状文件夹结构并最终删除整个文件夹结构。为此,我将Blocking Queue 与一个生产者(遍历树并放置需要删除的文件的路径)和一些实际执行删除作业的消费者一起使用。

文件夹必须为空才能被删除,因此,考虑具有以下结构:

/folder1/file1
/folder1/file2
/folder1/file3
/folder1/

这是BlockingQueue在任何给定点可能包含的内容。考虑到有 4 个消费者轮询队列:

Consumer1 会拿起并删除/folder1/file1

Consumer2 会拿起并删除/folder1/file2

Consumer3 会拿起并删除/folder1/file3

Consumer4 会拿起并删除/folder1


如果 Consumer3 没有完成删除 /folder1/file3,Consumer4 将无法删除 /folder1/,因为它将被标记为非空。

有没有办法让消费者线程等待其他消费者线程完成某些任务?

【问题讨论】:

  • 您可以为每个文件夹使用CountDownLatch,它将使用文件夹中包含的文件数进行初始化。 BlockingQueue 将包含一个元组 (file, countDownLatch)。如果File是文件,一旦删除,消费者调用countDown(),如果File是文件夹,消费者调用await()然后删除文件夹。
  • 这是您的真实用例还是类比?
  • @NicolasLabrot 这是一个真实的案例。每个文件夹可以包含数十万(如果不是数百万)文件/子文件夹
  • @AdrianDanielCulea 由于您使用的是树,因此请专注于删除树的叶子。
  • 出于好奇,您是否有任何证据支持 N 个线程能够比一个线程更快地删除文件和文件夹的想法?您的计算机可能有多个 CPU,但您的硬盘驱动器可能只有一个接口。

标签: java multithreading producer-consumer blockingqueue


【解决方案1】:

有很多方法可以解决业务问题。

方法 1:您的问题是当消费者 4 实际进入文件夹时,它需要等待删除所有文件。如果消费者 4 有权访问文件夹 1,我认为它不必这样做。它可以转到文件夹(操作系统路径)并检查它是否为空。如果为空,则将其删除,否则请等待。

方法 2:您的生产者线程可以做更多的工作。如果它发现文件夹1中的所有文件都需要删除。它不必放置所有文件名,然后是文件夹名。它应该只输入文件夹名称。只有一个消费者线程会获取文件夹名称并将其删除。

【讨论】:

  • 非常好的点,但是,方法 1 不可行,因为所有删除都不是“在我的机器上”发生的,而是在分布式数据库中发生的,因此删除返回文件夹中文件数量的元数据被其他进程删除。
  • 方法 2 可能有效,但是,如果我的理解完全错误,您建议使用消费者和生产者遍历树,我试图避免构建此系统
  • 在方法 2 中,我假设如果消费者线程删除了一个文件夹,那么包含的文件也会被删除。我还假设当Producer线程实际放入/folder1/file1 /folder1/file2 /folder1/file3 /folder1/最后一个Entry只有在Producer确保folder1中包含的所有文件都被删除时才会进行。如果是这样,生产者不应该放 file1、file2、file3.. 它应该只放 folder1
  • 如果文件夹不为空,则不能永远删除(来自与填充程序交互的 java 驱动程序的底层删除代码不允许这样做)
  • 如果在删除所有包含文件之前无法删除文件夹。您删除文件夹的消费者线程可以按以下方式工作。假设消费者 4 找到 Folder1。现在它将产生 3 个任务。即删除 Folder1 中包含的所有 3 个文件。现在,由于您正在从消费者线程 4 提交所有这些任务。所以您可以在 Future 上等待消费者 4 以完成所有子任务。
【解决方案2】:

这是另一种方法。我假设 Folder1 将在来自 Folder1 的所有文件都排队之后排队。这里的诀窍是每当消费者线程拾取文件时,将线程名称更改为文件名称。现在,当消费者线程(消费者 4)来删除目录 Folder1 时,它将检查所有正在执行的线程,其名称为“Folder1”。如果它找到任何东西,那么它将等待它们的完成。我有以下代码示例。如果有任何问题,请告诉我

import java.util.ArrayList;
import java.util.List;
import java.util.Set;
import java.util.concurrent.Callable;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.stream.Collectors;

public class DeleteFilesAndFolders
{

public void executeSomeTasks() throws InterruptedException
{
    ExecutorService executor = Executors.newFixedThreadPool(5);
    List<DeleteTask> tasks = createSomeDummyTasks();

    for (DeleteTask t : tasks) {
        executor.submit(t);
    }

    executor.shutdown();
}

private List<DeleteTask> createSomeDummyTasks()
{

    List<DeleteTask> tasks = new ArrayList<>();
    tasks.add(new DeleteTask("/good/folder/file1"));
    tasks.add(new DeleteTask("/good/folder/file2"));
    tasks.add(new DeleteTask("/good/folder/file3"));
    tasks.add(new DeleteTask("/good/folder/file4"));
    tasks.add(new DeleteTask("/good/folder/file6"));
    tasks.add(new DeleteTask("/good/folder/file7"));
    tasks.add(new DeleteTask("/good/folder/file9"));
    tasks.add(new DeleteTask("/good/folder/file8"));
    tasks.add(new DeleteTask("/good/folder"));

    return tasks;
}

public static class DeleteTask implements Callable<String>
{

    volatile String fileNameToDelete;

    public DeleteTask(String fileNameToDelete)
    {
        this.fileNameToDelete = fileNameToDelete;
    }

    @Override
    public String call() throws Exception
    {

        // Just checking if it is a directory
        if (fileNameToDelete.equalsIgnoreCase("/good/folder")) {
            waitForDeletionOfChildFiles(fileNameToDelete);
        }

        String originalName = Thread.currentThread().getName();
        Thread.currentThread().setName(fileNameToDelete);
        Thread.sleep(20000);
        Thread.currentThread().setName(originalName);

        return null;
    }

    private void waitForDeletionOfChildFiles(String directoryName) throws InterruptedException
    {
        while (true) {
            Set<Thread> threadSet = Thread.getAllStackTraces().keySet();

            List<Thread> anyChildThreads = threadSet.stream()
                    .filter(p -> p.getName().contains(directoryName))
                    .collect(Collectors.toList());

            if (!anyChildThreads.isEmpty()) {
                System.out.println(
                        " Some Child threads still Running. I should be waiting !!");
                anyChildThreads.stream().forEach(x -> {
                    System.out.println(x.getName());
                });

            } else {
                System.out.println(
                        " All Done.. no wait clean this directory !!");
                break;
            }

            Thread.sleep(5000);

        }
    }
}

public static void main(String args[]) throws InterruptedException
{
    DeleteFilesAndFolders et = new DeleteFilesAndFolders();
    et.executeSomeTasks();
}

}

【讨论】:

    猜你喜欢
    • 2020-03-15
    • 2018-10-23
    • 2013-05-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多