【发布时间】:2016-09-12 05:40:13
【问题描述】:
我必须并行化现有的后台任务,这样它就不会连续消耗“x”资源,而是仅使用“y”线程(y
代码结构如下:
class BaseBackground implements Runnable {
@Override
public void run() {
int[] resources = findResources(...);
for (int resource : resources) {
processResource(resource);
}
stopProcessing();
}
public abstract void processResource(final int resource);
public void void stopProcessing() {
// Override by subclass as needed
}
}
class ChildBackground extends BaseBackground {
@Override
public abstract void processResource(final int resource) {
// does some work here
}
public void void stopProcessing() {
// reset some counts and emit metrics
}
}
我已按以下方式修改了ChildBackground:
class ChildBackground extends BaseBackground {
private final BlockingQueue<Integer> resourcesToBeProcessed;
public ChildBackground() {
ExecutorService executorService = Executors.newFixedThreadPool(2);
for (int i = 0; i < 2; ++i) {
executorService.submit(new ResourceProcessor());
}
}
@Override
public abstract void processResource(final int resource) {
resourcesToBeProcessed.add(resource);
}
public void void stopProcessing() {
// reset some counts and emit metrics
}
public class ResourceProcessor implements Runnable {
@Override
public void run() {
while (true) {
int nextResource = resourcesToBeProcessed.take();
// does some work
}
}
}
}
我不会每次都创建和拆除 ExecutorService,因为垃圾收集在我的服务中有点问题。虽然,我不明白这会有多糟糕,因为我不会在每次迭代中产生超过 10 个线程。
我无法理解如何等待所有ResourceProcessors 完成一次迭代的处理资源,以便我可以重置一些计数并在stopProcessing 中发出指标。我考虑了以下选项:
1) executorService.awaitTermination(超时)。这不会真正起作用,因为它会一直阻塞直到超时,因为ResourceProcessor 线程永远不会真正完成它们的工作
2) 我可以找出findResources 之后的资源数量,并将其提供给子类,并让每个ResourceProcessor 增加处理的资源数量。在重置计数之前,我将不得不等待 stopProcessing 中处理完所有资源。我需要类似 CountDownLatch 的东西,但它应该算 UP。这个选项会有很多状态管理,我不是特别喜欢。
3) 我可以更新public abstract void processResource(final int resource) 以包含总资源计数,并让子进程等到所有线程都处理完总资源。在这种情况下也会有一些状态管理,但仅限于子类。
在这两种情况中的任何一种情况下,我都必须添加 wait() 和 notify() 逻辑,但我对我的方法没有信心。这就是我所拥有的:
class ChildBackground extends BaseBackground {
private static final int UNSET_TOTAL_RESOURCES = -1;
private final BlockingQueue<Integer> resourcesToBeProcessed;
private int totalResources = UNSET_TOTAL_RESOURCES;
private final AtomicInteger resourcesProcessed = new AtomicInteger(0);
public ChildBackground() {
ExecutorService executorService = Executors.newFixedThreadPool(2);
for (int i = 0; i < 2; ++i) {
executorService.submit(new ResourceProcessor());
}
}
@Override
public abstract void processResource(final int resource, final int totalResources) {
if (this.totalResources == UNSET_TOTAL_RESOURCES) {
this.totalResources = totalResources;
} else {
Preconditions.checkState(this.totalResources == totalResources, "Consecutive poll requests are using different total resources count, previous=%s, new=%s", this.totalResources, totalResources);
}
resourcesToBeProcessed.add(resource);
}
public void void stopProcessing() {
try {
waitForAllResourcesToBeProcessed();
} catch (InterruptedException e) {
e.printStackTrace();
}
resourcesProcessed.set(0);
totalResources = UNSET_TOTAL_RESOURCES;
// reset some counts and emit metrics
}
private void incrementProcessedResources() {
synchronized (resourcesProcessed) {
resourcesProcessed.getAndIncrement();
resourcesProcessed.notify();
}
}
private void waitForAllResourcesToBeProcessed() throws InterruptedException {
synchronized (resourcesProcessed) {
while (resourcesProcessed.get() != totalResources) {
resourcesProcessed.wait();
}
}
}
public class ResourceProcessor implements Runnable {
@Override
public void run() {
while (true) {
int nextResource = resourcesToBeProcessed.take();
try {
// does some work
} finally {
incrementProcessedResources();
}
}
}
}
}
我不确定使用AtomicInteger 是否是正确的方法,如果是,我是否需要调用wait() 和notify()。如果我不使用 wait() 和 notify(),我什至不必在同步块中执行所有内容。
如果我应该为每次迭代简单地创建和关闭 ExecutorService 或者我应该采用第四种方法,请让我知道您对这种方法的看法。
【问题讨论】:
-
看看这个页面上的
Callables and Futureswinterbe.com/posts/2015/04/07/… -
Future 不是一个选项,因为执行程序线程在一个紧密循环中运行,只有在服务停用时才会停止。
-
在分配给执行者的任务完成时得到通知是未来思想的目的。确定不能使用,例如 docs.oracle.com/javase/8/docs/api/java/util/concurrent/… ?您将向执行者提交一批资源(按资源累积 1 个未来),创建一个等待所有这些资源组合的未来,并让最后一个计算您的指标。在提交最后一个未来时,您可能会忘记所有这些。采用反应式风格。
标签: java multithreading synchronized executorservice atomicinteger