gpt4 book ai didi

java - 如何终止从阻塞队列中检索

转载 作者:行者123 更新时间:2023-11-30 06:37:38 26 4
gpt4 key购买 nike

我有一些代码,我使用执行器和阻塞队列执行多个任务。结果必须作为迭代器返回,因为这是我工作的应用程序所期望的。但是,任务和添加到队列的结果之间存在 1:N 的关系,所以我不能使用 ExecutorCompletionService .在调用 hasNext() 时,我需要知道所有任务何时完成并将所有结果添加到队列中,以便我可以停止从队列中检索结果。请注意,一旦将项目放入队列,另一个线程应该准备好使用( Executor.invokeAll() ,阻塞直到所有任务完成,这不是我想要的,也不是超时)。这是我的第一次尝试,我使用 AtomicInteger 只是为了证明这一点,即使它不起作用。有人可以帮助我理解如何解决这个问题吗?

public class ResultExecutor<T> implements Iterable<T> {
private BlockingQueue<T> queue;
private Executor executor;
private AtomicInteger count;

public ResultExecutor(Executor executor) {
this.queue = new LinkedBlockingQueue<T>();
this.executor = executor;
count = new AtomicInteger();
}

public void execute(ExecutorTask task) {
executor.execute(task);
}

public Iterator<T> iterator() {
return new MyIterator();
}

public class MyIterator implements Iterator<T> {
private T current;
public boolean hasNext() {
if (count.get() > 0 && current == null)
{
try {
current = queue.take();
count.decrementAndGet();
} catch (InterruptedException e) {
e.printStackTrace();
}
}
return current != null;
}

public T next() {
final T ret = current;
current = null;
return ret;
}

public void remove() {
throw new UnsupportedOperationException();
}

}

public class ExecutorTask implements Runnable{
private String name;

public ExecutorTask(String name) {
this.name = name;

}

private int random(int n)
{
return (int) Math.round(n * Math.random());
}


@SuppressWarnings("unchecked")
public void run() {
try {
int random = random(500);
Thread.sleep(random);
queue.put((T) (name + ":" + random + ":1"));
queue.put((T) (name + ":" + random + ":2"));
queue.put((T) (name + ":" + random + ":3"));
queue.put((T) (name + ":" + random + ":4"));
queue.put((T) (name + ":" + random + ":5"));

count.addAndGet(5);
} catch (InterruptedException e) {
e.printStackTrace();
}
}
}

}

调用代码如下:

    Executor e = Executors.newFixedThreadPool(2);
ResultExecutor<Result> resultExecutor = new ResultExecutor<Result>(e);

resultExecutor.execute(resultExecutor.new ExecutorTask("A"));
resultExecutor.execute(resultExecutor.new ExecutorTask("B"));

Iterator<Result> iter = resultExecutor.iterator();
while (iter.hasNext()) {
System.out.println(iter.next());
}

最佳答案

Queue 中使用“毒药”对象来表示任务将不再提供结果。

class Client
{

public static void main(String... argv)
throws Exception
{
BlockingQueue<String> queue = new LinkedBlockingQueue<String>();
ExecutorService workers = Executors.newFixedThreadPool(2);
workers.execute(new ExecutorTask("A", queue));
workers.execute(new ExecutorTask("B", queue));
Iterator<String> results =
new QueueMarkersIterator<String>(queue, ExecutorTask.MARKER, 2);
while (results.hasNext())
System.out.println(results.next());
}

}

class QueueMarkersIterator<T>
implements Iterator<T>
{

private final BlockingQueue<? extends T> queue;

private final T marker;

private int count;

private T next;

QueueMarkersIterator(BlockingQueue<? extends T> queue, T marker, int count)
{
this.queue = queue;
this.marker = marker;
this.count = count;
this.next = marker;
}

public boolean hasNext()
{
if (next == marker)
next = nextImpl();
return (next != marker);
}

public T next()
{
if (next == marker)
next = nextImpl();
if (next == marker)
throw new NoSuchElementException();
T tmp = next;
next = marker;
return tmp;
}

/*
* Block until the status is known. Interrupting the current thread
* will cause iteration to cease prematurely, even if elements are
* subsequently queued.
*/
private T nextImpl()
{
while (count > 0) {
T o;
try {
o = queue.take();
}
catch (InterruptedException ex) {
count = 0;
Thread.currentThread().interrupt();
break;
}
if (o == marker) {
--count;
}
else {
return o;
}
}
return marker;
}

public void remove()
{
throw new UnsupportedOperationException();
}

}

class ExecutorTask
implements Runnable
{

static final String MARKER = new String();

private static final Random random = new Random();

private final String name;

private final BlockingQueue<String> results;

public ExecutorTask(String name, BlockingQueue<String> results)
{
this.name = name;
this.results = results;
}

public void run()
{
int random = ExecutorTask.random.nextInt(500);
try {
Thread.sleep(random);
}
catch (InterruptedException ignore) {
}
final int COUNT = 5;
for (int idx = 0; idx < COUNT; ++idx)
results.add(name + ':' + random + ':' + (idx + 1));
results.add(MARKER);
}

}

关于java - 如何终止从阻塞队列中检索,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/3214597/

26 4 0
Copyright 2021 - 2024 cfsdn All Rights Reserved 蜀ICP备2022000587号
广告合作:1813099741@qq.com 6ren.com