【问题标题】:Parallel writers : How to stop threads with Poison pill pattern?并行作者:如何用毒丸模式停止线程?
【发布时间】:2021-09-27 13:07:27
【问题描述】:

在我的情况下读取必须是顺序的,瓶颈被定义为处理和写入数据库。

按照@Mahmoudhere 的建议,我已经设法通过使用阻塞队列将两个进程(读/写)分开,因此写入步骤能够在不影响读取的情况下扩展。

为了在没有更多要阅读的项目时停止侦听队列,我引入了毒丸模式。在这种情况下,我的队列阅读器如下所示:

@RequiredArgsConstructor
class BlockingQueueItemReader<T> implements ItemReader<T> {

    private final BlockingQueue<T> queue;
    private final T poisonPill;
    private final int timeoutSeconds;

    @Nullable
    @Override
    public T read() throws Exception {
        T taken = queue.poll(timeoutSeconds, TimeUnit.SECONDS);

        if (poisonPill.equals(taken)) {
            return null;
        }

        return taken;
    }
}

为了同时运行多个写入,我在 step2 中添加了一个执行器,例如:


    @Bean
    public TaskExecutor taskExecutor() {
        ThreadPoolTaskExecutor executor = new ThreadPoolTaskExecutor();
        executor.setCorePoolSize(10);
        executor.setMaxPoolSize(10);
        executor.setThreadNamePrefix("MyExe-");
        return executor;
    }

    @Bean
    public Step step2() {
        return steps.get("step2")
                .<Person, Person>chunk(10)
                .reader(new BlockingQueueItemReader<>(queue(), POISON))
                .writer(items -> {
                    for (Person item : items) {
                        System.out.println("item = " + item);
                    }
                })
                .taskExecutor(taskExecutor())
                .throttleLimit(8)
                .build();
    }

现在,多个线程同时处理多个块,这正是我所寻找的。​​p>

我现在的问题是BlockingQueueItemReader。一些读者在poll 行被阻止。发生这种情况是因为最后读取的元素不是 POISON 对象,同时另一个线程找到它并返回 null(因此该线程将停止而不是其他线程)。

为了解决这个问题,我再次将实现更改为:

@RequiredArgsConstructor
public class BlockingQueueItemReader<T> implements ItemReader<T> {

    private final BlockingQueue<T> queue;
    private final T poisonPill;
    private final int timeoutSeconds;
    private boolean exhausted;


    @Nullable
    @Override
    public T read() throws Exception {
        if (exhausted) {
            return null;
        }

        T taken = queue.poll(timeoutSeconds, TimeUnit.SECONDS);
        exhausted = poisonPill.equals(taken);
        
        if (exhausted) {
            return null;
        }

        return taken;
    }
}

这样所有线程都正常退出。

我的问题是:我对这个版本不满意,exhausted 变量的仔细检查看起来很难看!

当至少一个线程找到POISON 对象时,是否有另一种方法告诉所有相关线程停止?

【问题讨论】:

    标签: java spring spring-batch java.util.concurrent


    【解决方案1】:

    当至少一个线程找到 POISON 对象时,是否有另一种方法告诉所有相关线程停止?

    我认为“毒丸”的想法不适用于这个线程同步问题(或者至少不容易干净地实现)。在我看来,基于超时的方法更好,因为它不需要任何额外的代码来在队列中注入有毒项目 + 在阅读器中检测 + 同步线程。

    【讨论】:

    • 感谢您的回答。在我的情况下,超时很难选择正确的值。这就是我使用毒丸的原因。我必须在最后一个实现中同步线程吗?
    猜你喜欢
    • 1970-01-01
    • 2012-11-30
    • 1970-01-01
    • 2016-06-30
    • 2020-01-28
    • 1970-01-01
    • 2019-12-29
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多