- html - 出于某种原因,IE8 对我的 Sass 文件中继承的 html5 CSS 不友好?
- JMeter 在响应断言中使用 span 标签的问题
- html - 在 :hover and :active? 上具有不同效果的 CSS 动画
- html - 相对于居中的 html 内容固定的 CSS 重复背景?
我的代码基本上遵循官方教程,主要目的是收集一个订阅(Constants.UNFINISHEDSUBID)中的所有消息并在另一个订阅上重新发布它们。但目前我面临着一个我无法解决的问题。在我的实现中,调用subscriber.stopAsync() 会导致以下异常:
Mai 04, 2017 4:59:25 PM com.google.common.util.concurrent.AbstractFuture executeListener
SCHWERWIEGEND: RuntimeException while executing runnable com.google.common.util.concurrent.Futures$6@6e13e898 with executor java.util.concurrent.Executors$DelegatedScheduledExecutorService@2f3c6ac4
java.util.concurrent.RejectedExecutionException: Task java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask@60d40af2 rejected from java.util.concurrent.ScheduledThreadPoolExecutor@d55b6e[Terminated, pool size = 0, active threads = 0, queued tasks = 0, completed tasks = 320]
at java.util.concurrent.ThreadPoolExecutor$AbortPolicy.rejectedExecution(ThreadPoolExecutor.java:2047)
at java.util.concurrent.ThreadPoolExecutor.reject(ThreadPoolExecutor.java:823)
at java.util.concurrent.ScheduledThreadPoolExecutor.delayedExecute(ScheduledThreadPoolExecutor.java:326)
at java.util.concurrent.ScheduledThreadPoolExecutor.schedule(ScheduledThreadPoolExecutor.java:533)
at java.util.concurrent.ScheduledThreadPoolExecutor.execute(ScheduledThreadPoolExecutor.java:622)
at java.util.concurrent.Executors$DelegatedExecutorService.execute(Executors.java:668)
at com.google.common.util.concurrent.AbstractFuture.executeListener(AbstractFuture.java:817)
at com.google.common.util.concurrent.AbstractFuture.complete(AbstractFuture.java:753)
at com.google.common.util.concurrent.AbstractFuture.set(AbstractFuture.java:613)
at io.grpc.stub.ClientCalls$GrpcFuture.set(ClientCalls.java:458)
at io.grpc.stub.ClientCalls$UnaryStreamToFuture.onClose(ClientCalls.java:437)
at io.grpc.internal.ClientCallImpl.closeObserver(ClientCallImpl.java:428)
at io.grpc.internal.ClientCallImpl.access$100(ClientCallImpl.java:76)
at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.close(ClientCallImpl.java:514)
at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl.access$700(ClientCallImpl.java:431)
at io.grpc.internal.ClientCallImpl$ClientStreamListenerImpl$1StreamClosed.runInContext(ClientCallImpl.java:546)
at io.grpc.internal.ContextRunnable.run(ContextRunnable.java:52)
at io.grpc.internal.SerializingExecutor$TaskRunner.run(SerializingExecutor.java:152)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)
我还注意到这种随机的情况,有时是所有消息,有时只是少数或没有一条消息被收集。调用subscriber.stopAsync()不是正确的方法吗?
我当前的实现:
protected void pullUnfinished() throws Exception {
List<PubsubMessage> jobsToRepublish = new ArrayList<>();
SubscriptionName subscription =
SubscriptionName.create(Constants.PROJECTID, Constants.UNFINISHEDSUBID);
MessageReceiver receiver = new MessageReceiver() {
@Override
public void receiveMessage(PubsubMessage message, AckReplyConsumer consumer) {
synchronized(jobsToRepublish){
jobsToRepublish.add(message);
}
String unfinishedJob = message.getData().toStringUtf8();
LOG.info("got message: {}", unfinishedJob);
consumer.ack();
}
};
Subscriber subscriber = null;
try {
ChannelProvider channelProvider = new PlainTextChannelProvider();
subscriber = Subscriber.defaultBuilder(subscription, receiver)
.setChannelProvider(channelProvider)
.build();
subscriber.addListener(new Subscriber.Listener() {
@Override
public void failed(Subscriber.State from, Throwable failure) {
System.err.println(failure);
}
}, MoreExecutors.directExecutor());
subscriber.startAsync().awaitRunning();
Thread.sleep(60000);
} finally {
if (subscriber != null) {
subscriber.stopAsync(); //Causes the exception
}
}
publishJobs(jobsToRepublish);
}
public class PlainTextChannelProvider implements ChannelProvider {
@Override
public boolean shouldAutoClose() {
// TODO Auto-generated method stub
return false;
}
@Override
public boolean needsExecutor() {
// TODO Auto-generated method stub
return false;
}
@Override
public ManagedChannel getChannel() throws IOException {
return NettyChannelBuilder.forAddress("localhost", 8085)
.negotiationType(NegotiationType.PLAINTEXT)
.build();
}
@Override
public ManagedChannel getChannel(Executor executor) throws IOException {
return getChannel();
}
}
最佳答案
当我从 JUnit 测试运行类似代码时遇到了完全相同的问题,并发现了这个 related answer一般而言,在多线程上,表明线程池已关闭,而监听器仍在引用它。我还查看了 Subscriber.java 的代码在 GitHub 上,在 JavaDoc 中关于 startAsync() 找到了一个接收大量消息的示例,建议等待 stopAsync() 终止。
尝试改变
subscriber.stopAsync();
至
subscriber.stopAsync().awaitTerminated();
为我工作。
关于java - Subscriber.stopAsync() 导致 RejectedExecutionException,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/43786716/
在我们的一项服务中,有人添加了这样(简化)的一段代码: public class DeleteMe { public static void main(String[] args) {
下面是我的方法,其中我有单线程执行器在 run 方法中执行一些任务。 private void trigger(final Packet packet) { // this line is
是什么导致了此 RejectedExecutionException? [Running, pool size = 40, active threads = 3, queued tasks = 20,
为什么当来自线程池的线程之一抛出 RejectedExecutionException 时主线程没有停止?我在这里做错了吗?线程池中的第三个线程抛出 RejectedExecutionExceptio
除了先前在 Executor 上调用的 shutdown() 之外,是否还有其他原因导致 RejectedExecutionException 被抛出(我使用的是 singleThreadExecut
谁能给我提供一个获得 RejectedExecutionException 的例子可能是一个现实生活中的例子。提前致谢。 最佳答案 Anybody able to provide me with an
间歇性头痛需要帮助。代码调用 com.google.api.client.http.HttpRequest#executeAsync() 基本上具有以下逻辑, @Beta public Fut
我正在开发一款在 Android NDK 中运行大部分原生代码的社交游戏。游戏有 3 个主要的 ndk pthreads: 一个游戏线程 服务器通信线程 主渲染线程(通过 Renderer.onRen
我在我的 tomcat 服务器 (+liferay) 上遇到此异常 java.util.concurrent.RejectedExecutionException 我的课是这样的: public cl
我想实现以下行为: 从文件中读取 n 个事件 在线程中处理它们 如果仍有任何事件,请返回步骤 1 我编写了以下应用程序来测试解决方案,但它在随机时刻失败,例如。 java.lang.IllegalSt
我正在使用 AsyncTask 从远程服务器获取大量缩略图并在 GridView 中显示它们。问题是,我的 GridView 一次显示 20 个缩略图,因此创建 20 个 AsyncTasks 并启动
我正在用 java 编写一个多线程程序。我写过这样的东西 exec.execute(p) // where p is a runnable task working on an array prin
我正在从远程服务器获取大量缩略图,并使用 AsyncTask 在 GridView 中显示它们。问题是,我的 GridView 一次显示 20 个缩略图,因此创建 20 个 AsyncTask 并启动
我有这个客户: OkHttpClient okHttpClient = new OkHttpClient.Builder() .pingInterval(Duration.of
我正在尝试将行批量放入 HBase(0.90.0)中,大小约为 1000(行)我有多个生产者线程将数据写入队列,还有一个消费者线程每几分钟唤醒一次,并写入所有内容在队列中作为批处理到 HBase。但是
我希望atomicInteger的值为100,然后程序终止 public static void main(String[] args) throws InterruptedException {
我创建调度程序来测试处理 RejectedExecutionException: @Component public class TestScheduler { private final T
我正在尝试实现生产者-消费者模式,并且我希望能够阻止消费者。到目前为止我写道: import java.util.concurrent.BlockingQueue; import java.util.
我的代码基本上遵循官方教程,主要目的是收集一个订阅(Constants.UNFINISHEDSUBID)中的所有消息并在另一个订阅上重新发布它们。但目前我面临着一个我无法解决的问题。在我的实现中,调用
这是一个例子: // max 1 pending task in queue val queue = LinkedBlockingQueue(1) // max 1 thread / 1 active
我是一名优秀的程序员,十分优秀!