【发布时间】: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