【问题标题】:Multithread: Dynamically object locking based on resource in hava多线程:基于java中资源的动态对象锁定
【发布时间】:2018-03-26 18:21:00
【问题描述】:

这里需要线程专家的眼睛......

我在一个 poc 应用程序上,我在 FTP 服务器上上传文件

在 FTP 服务器中有多个文件夹。根据输入响应,我从文件夹中读取文件并移动到另一个文件夹

应用可以同时被多个线程访问。

所以问题是这样的:

假设 FTP 有一个文件夹 Folder_A 和 A_A_FOLDER 现在 Folder_A 有 10 个文件。 一个线程来了,从 FTP 读取了 10 个文件并开始对其进行一些计算, 它一一计算,然后移动到 A_A_FOLDER 它处于过程的中间(假设它成功地将 5 个文件从 Folder_A 移动到 A_A_FOLDER) 然后另一个线程来了,它选择了剩余的 5 个文件,因为它们被线程 1 处理不足,所以线程 2 也开始处理这 5 个文件

所以这里有重复文件的问题

void m1(String folderName) {
// FTP related code
}

我已经通过使用同步关键字解决了这个问题

现在一切都在同步,所有处理工作正常

synchronized void m1(String folderName) {
// code
}

文件夹名称决定需要处理哪个文件夹

现在我开始面临性能问题

因为方法是同步的,所以所有线程都会等到处理线程未完成其任务。

我可以通过以下步骤改进:

(在找到解决方案之前,这里有一些关于这个问题的故事)

正如我提到的 m1 方法的 folderName 参数决定将处理哪个文件夹, 所以假设我在 Ftp 服务器中有 4 个文件夹(A、B、A_T、B_T),2 个文件夹是需要从(A 和 B)读取数据的文件夹, 2个文件夹是数据将移动的文件夹(A_T和B_T)

A_T 和 B_T 不是这里的问题,因为它们对于每个文件夹 A 和 B 都是唯一的 因此,如果该方法将从 A 读取,那么它会将其移至 A_T,与 B 相同(移至 B_T)

现在:

假设 m1 方法有 4 个线程,文件夹 A 有 3 个线程,文件夹 B 有 1 个线程 如果基于 fileName 参数以某种方式同步请求以便我可以提高性能,则意味着 1 个线程将在 A 上工作,另外 2 个线程将阻塞,因为它们的 fileName 相同,所以他们将等到第一个线程未完成它的任务,线程 4 将并行无需任何锁定过程即可工作,因为它的文件名不同

那么我怎样才能在代码级别实现这一点(在文件名上同步)?

注意:我知道我可以使用资源的静态锁定列表然后锁定文件名资源来打破这个逻辑 例如:

private final Object A = new Object();
private final Object B = new Object();

但这种方法的问题是文件夹可以动态添加,所以我不能这样做。

需要你们的帮助。

【问题讨论】:

  • 您可以使用ConcurrentHashMap<String, CountDownLatch>,其中键是文件名,值是ne CountDownLatch(1) 的实例。一个执行移动的线程控制闩锁并将其倒计时,其他线程仅在闩锁未倒计时且仅等待其文件名时才等待闩锁。试一试,如果你不成功,我会给你写一个例子
  • 您对 FTP 的使用是一个根本问题,因为在 FTP 下无法知道文件传输是否已完成。传输可能中途失败,然后重新启动或从中断的地方继续,并且在传输“完成”时没有可用的信号或符号。您永远无法 100% 确定您在服务器上传目录中看到的内容是否可用。您应该改为迁移到rsync,可以将其配置为写入临时文件并仅在收到所有数据时创建最终文件。
  • @JimGarrison 感谢您的回复,但这里的主要问题不是 FTP,我已经实现了 FTP 的流程,这也不是我关心的问题,这里的重点是多线程控制这种场景。因为在多线程的情况下,很多情况下都会出现这种问题。所以我很好奇在多线程环境中处理这个问题
  • @OlegSklyar 我在下面添加了完整的解决方案,请您查看一下。您是否有任何建议(我是否需要任何额外的工作),将是可观的。谢谢!

标签: java multithreading java.util.concurrent


【解决方案1】:

一种方法是为每个目录维护一个锁:

public class DirectoryTaskManager {
    public static void main(String[] args) throws IOException {
        DirectoryTaskManager manager = new DirectoryTaskManager();
        manager.withDirLock(new File("Folder_A"), () -> System.out.println("Doing something..."));
    }

    public void withDirLock(File dir, Runnable task) throws IOException {
        ReentrantLock lock = getDirLock(dir);
        lock.lock();
        try {
            task.run();
        } finally {
            lock.unlock();
        }
    }

    private Map<File, ReentrantLock> dirLocks = Collections.synchronizedMap(new HashMap<>());

    public ReentrantLock getDirLock(File dir) throws IOException {
        // Resolve the canonical file here so that different paths 
        // to the same file use the same lock
        File canonicalDir = dir.getCanonicalFile();
        if (!canonicalDir.exists() || !canonicalDir.isDirectory()) {
            throw new FileNotFoundException(canonicalDir.getName());
        }
        return dirLocks.computeIfAbsent(canonicalDir, d -> new ReentrantLock());
    }
}

【讨论】:

  • 顶一下:@teppic 谢谢你的指导。我在下面添加了完整的工作示例,请您查看一下
  • 不错的答案。有什么理由使用 Collections.synchronizedMap 而不是 ConcurrentHashMap?
  • @Sneh:在大多数情况下,我更喜欢Collections.synchronizedMapConcurrentHashMap 为读取速度牺牲了一致性,并且具有更复杂的锁定语义。
  • 我明白你的意思,但由于对象级锁定,syncronizedMap 的性能确实受到影响。它对这个答案没有太大的影响,因为只调用了 computeIfAbsenr。这是一本好书ibm.com/developerworks/java/library/j-jtp08223/index.html
【解决方案2】:

感谢@teppic 和@OlegSklyar 的指导 最后这里是完整的工作示例,

FolderImpl -> 有方法名调用,可以被多个线程访问

我使用了 ConcurrentHashMap(读取可以非常快,而写入是通过锁定完成的。)它比 synchronizedMap 快 它将保存文件夹名称和 ReentrantLock,因此锁定将对文件夹名称起作用

public class FolderImpl {
    private FolderImpl(){
        System.out.println("init................");
    }

    private ConcurrentHashMap<String, ReentrantLock> concurrentHashMap= new ConcurrentHashMap();
    private static final FolderImpl singleTon = new FolderImpl();
    public static FolderImpl getSingleTon() {
        return singleTon;
    }

    public void call(String name) throws Exception{
        ReentrantLock getDirLock = getDirLock(name);
        getDirLock.lock();
        try {
        for (int i = 0; i < 100; i ++) {
            System.out.println(i+":"+name+":"+Thread.currentThread().getName());
            try {
                Thread.sleep(30);
            } catch (Exception e) {
                e.printStackTrace();
            }
        }}finally {
            getDirLock.unlock();
        }

    }

    public ReentrantLock getDirLock(String site)  {
        return concurrentHashMap.computeIfAbsent(site, d -> new ReentrantLock());
    }
}

TaskCaller 线程调用 call 方法,这里是 sleep 风格,所以另一个线程可以 git 执行时间

public class TaskCaller extends Thread{
    public FolderImpl singleTon = FolderImpl.getSingleTon();
    public TaskCaller(String name) {
        super();
        this.name = name;
    }

    private String name;
    @Override
    public void run() {
        for (int i = 0; i < 5; i++) {
            System.out.println("Name:"+Thread.currentThread().getName());

            try {
                singleTon.call(name);
                sleep(10);
            } catch (Exception e) {
                e.printStackTrace();
            }
        }
    }
}

TestExecution 类将执行 10 个线程进行测试

public class TestExecution {

    public static void main(String[] args) {
        TaskCaller testThreadCC = new TaskCaller("A_FOLDER");
        TaskCaller testThreadCC2 = new TaskCaller("A_FOLDER");
        TaskCaller testThreadCC3 = new TaskCaller("B_FOLDER");
        TaskCaller testThreadCC4 = new TaskCaller("C_FOLDER");
        TaskCaller testThreadCC5 = new TaskCaller("C_FOLDER");
        TaskCaller testThreadCC6 = new TaskCaller("C_FOLDER");
        TaskCaller testThreadCC7 = new TaskCaller("A_FOLDER");
        TaskCaller testThreadCC8 = new TaskCaller("A_FOLDER");
        TaskCaller testThreadCC9 = new TaskCaller("B_FOLDER");
        TaskCaller testThreadCC10 = new TaskCaller("B_FOLDER");

        testThreadCC.start();
        testThreadCC2.start();
        testThreadCC3.start();
        testThreadCC4.start();
        testThreadCC5.start();
        testThreadCC6.start();
        testThreadCC7.start();
        testThreadCC8.start();
        testThreadCC9.start();
        testThreadCC10.start();

    }

}

【讨论】:

  • 顺便说一句,锁定机制的选择会产生很大的不同。阅读此blog.takipi.com/…
  • @Sneh 谢谢,博客很丰富,但这里 StampedLock 或 ReentrantReadWriteLock 没有满足我的要求,它们可以,但它们对于这种情况效率不高。在这些锁定策略中,最适合允许阅读锁被多个读取线程同时持有,只要没有写入者,并且写入锁是独占的。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-06-10
相关资源
最近更新 更多