我有以下代码片段,它基本上扫描需要执行的任务列表,然后将每个任务提供给执行程序执行.
该JobExecutor反过来又创造了另一个执行人(做数据库的东西...读取和写入数据队列),并完成任务.
JobExecutor返回Future<Boolean>提交的任务.当其中一个任务失败时,我想优雅地中断所有线程并通过捕获所有异常来关闭执行程序.我需要做哪些改变?
public class DataMovingClass {
private static final AtomicInteger uniqueId = new AtomicInteger(0);
private static final ThreadLocal<Integer> uniqueNumber = new IDGenerator();
ThreadPoolExecutor threadPoolExecutor = null ;
private List<Source> sources = new ArrayList<Source>();
private static class IDGenerator extends ThreadLocal<Integer> {
@Override
public Integer get() {
return uniqueId.incrementAndGet();
}
}
public void init(){
// load sources list
}
public boolean execute() {
boolean succcess = true ;
threadPoolExecutor = new ThreadPoolExecutor(10,10,
10, TimeUnit.SECONDS, new ArrayBlockingQueue<Runnable>(1024), …Run Code Online (Sandbox Code Playgroud) java concurrency multithreading exception-handling threadpoolexecutor
当在方法中提交新任务
execute(java.lang.Runnable)且corePoolSize运行的线程少于正在运行时,即使其他工作线程处于空闲状态,也会创建一个新线程来处理该请求.
1)如果有空闲线程,为什么需要创建一个新线程来处理请求?
如果运行的线程数多于
corePoolSize但少于maximumPoolSize线程,则仅当队列已满时才会创建新线程.
2)我不明白corePoolSize和maximumPoolSize这里之间的区别.其次,当线程小于时,队列如何才能满maximumPoolSize?如果线程等于或大于,则队列只能是满的maximumPoolSize.不是吗?
除了Executor接口比普通线程(例如管理)具有一些优势之外,在执行以下操作之间是否存在任何真正的内部差异(大的性能差异,资源消耗......)
ExecutorService executor = Executors.newSingleThreadExecutor();
executor.submit(runnable);
Run Code Online (Sandbox Code Playgroud)
和:
Thread thread = new Thread(runnable);
thread.start();
Run Code Online (Sandbox Code Playgroud)
我这里只询问一个帖子.
我有一个链接的阻塞队列,我在其中执行插入和删除操作.
我需要知道哪一个更好put或者offer在链接阻塞队列的情况下.
性能参数是CPU利用率,内存和总吞吐量.
应用程序使用是实时系统,其中可以有多个传入请求和更少的线程来处理我们需要在队列中插入元素的位置.
我读了Java文件的put和offer在内部应用程序中没有太大区别.
在我的Android项目中,我有很多地方需要异步运行一些代码(Web请求,调用db等).这不是长时间运行的任务(最多几秒钟).到目前为止,我正在创建一个新线程,通过任务传递一个新的runnable.但是最近我读了一篇关于Java中的线程和并发的文章,并且理解为每个任务创建一个新的Thread并不是一个好的决定.
所以现在我ThreadPoolExecutor在我的Application类中创建了一个包含5个线程的类.这是代码:
public class App extends Application {
private ThreadPoolExecutor mPool;
@Override
public void onCreate() {
super.onCreate();
mPool = (ThreadPoolExecutor)Executors.newFixedThreadPool(5);
}
}
Run Code Online (Sandbox Code Playgroud)
我还有一种方法可以将Runnable任务提交给执行者:
public void submitRunnableTask(Runnable task){
if(!mPool.isShutdown() && mPool.getActiveCount() != mPool.getMaximumPoolSize()){
mPool.submit(task);
} else {
new Thread(task).start();
}
}
Run Code Online (Sandbox Code Playgroud)
因此,当我想在我的代码中运行异步任务时,我得到实例App并调用submitRunnableTask将runnable传递给它的方法.正如你所看到的,我也检查一下,如果线程池有自由线程来执行我的任务,如果没有,我创建一个新线程(我不认为这会发生,但无论如何......我不知道我希望我的任务在队列中等待并减慢应用程序).
在onTerminateApplication 的回调方法中,我关闭了池.
所以我的问题如下:这种模式比代码中创建新的线程更好吗?我的新方法有哪些优点和缺点?它会引起我不知道的问题吗?你能告诉我一些比这更好的东西来管理我的异步任务吗?
PS我在Android和Java方面有一些经验,但我远不是一个并发大师.所以可能有些方面我在这类问题中不太了解.任何建议将被认真考虑.
如何创建scala.concurrent.ExecutionContext?
文档通常给出一个总体摘要,并提到"默认"实现scala.concurrent.ExecutionContext.global.
尽管如此,有时你必须创建你的个人EC,而不使用akka和其他这样的工具.
我搜索了很多但找不到任何解决方案.我用这样的方式使用java线程池:
ExecutorService c = Executors.newFixedThreadPool(3);
for (int i = 0; i < 10; ++i) {
c.execute(new MyTask(i));
}
Run Code Online (Sandbox Code Playgroud)
以这种方式,任务以后续顺序执行(如在队列中).但我需要改变"选择下一个任务"策略.所以我希望为每个任务分配指定优先级(它不是线程优先级),并且执行任务对应于这些优先级.因此,当执行程序完成另一个任务时,它会将下一个任务选为具有最高优先级的任务.它描述了常见问题.也许有更简单的方法不考虑优先级.它选择最后添加的任务作为执行而不是第一次添加.粗略地讲,FixedThreadPool使用FIFO策略.我可以使用例如LIFO策略吗?
我需要在Java中构建一个工作池,每个工作者都有自己的连接套接字; 当工作线程运行时,它使用套接字但保持打开以便以后重用.我们决定使用这种方法,因为与临时创建,连接和销毁套接字相关的开销需要太多的开销,所以我们需要一种方法,通过这种方法,工作池预先初始化了它们的套接字连接,准备好在保持套接字资源不受其他线程影响的同时承担工作(套接字不是线程安全的),所以我们需要这些内容......
public class SocketTask implements Runnable {
Socket socket;
public SocketTask(){
//create + connect socket here
}
public void run(){
//use socket here
}
Run Code Online (Sandbox Code Playgroud)
}
在应用程序启动时,我们想要初始化工作程序,并希望套接字连接在某种程度上......
MyWorkerPool pool = new MyWorkerPool();
for( int i = 0; i < 100; i++)
pool.addWorker( new WorkerThread());
Run Code Online (Sandbox Code Playgroud)
当应用程序请求工作时,我们将任务发送到工作池以立即执行...
pool.queueWork( new SocketTask(..));
Run Code Online (Sandbox Code Playgroud)
更新了工作代码
根据Gray和jontejj的有用评论,我有以下代码工作...
SocketTask
public class SocketTask implements Runnable {
private String workDetails;
private static final ThreadLocal<Socket> threadLocal =
new ThreadLocal<Socket>(){
@Override
protected Socket initialValue(){
return new Socket();
}
};
public SocketTask(String details){ …Run Code Online (Sandbox Code Playgroud) CompletableFuture::supplyAsync(() -> IO bound queries)
如何为CompletableFuture :: supplyAsync选择Executor以避免污染ForkJoinPool.commonPool().
有许多选项Executors(newCachedThreadPool,newWorkStealingPool,newFixedThreadPool等)
我在这里阅读了关于新ForkJoinPool的内容
如何为我的用例选择合适的?
java executorservice java-8 threadpoolexecutor completable-future
在java-9 中,引入了类中的新方法completeOnTimeoutCompletableFuture:
public CompletableFuture<T> completeOnTimeout(T value, long timeout,
TimeUnit unit) {
if (unit == null)
throw new NullPointerException();
if (result == null)
whenComplete(new Canceller(Delayer.delay(
new DelayedCompleter<T>(this, value),
timeout, unit)));
return this;
}
Run Code Online (Sandbox Code Playgroud)
我不明白为什么它在其实现中使用静态 ScheduledThreadPoolExecutor:
static ScheduledFuture<?> delay(Runnable command, long delay,
TimeUnit unit) {
return delayer.schedule(command, delay, unit);
}
Run Code Online (Sandbox Code Playgroud)
哪里
static final ScheduledThreadPoolExecutor delayer;
static {
(delayer = new ScheduledThreadPoolExecutor(
1, new DaemonThreadFactory())).
setRemoveOnCancelPolicy(true);
}
Run Code Online (Sandbox Code Playgroud)
对我来说这是一种非常奇怪的方法,因为它可能成为整个应用程序的瓶颈:唯一一个ScheduledThreadPoolExecutor只有一个线程保留在池中以执行所有可能的CompletableFuture任务?
我在这里错过了什么?
PS它看起来像:
1)这段代码的作者不愿意提取这种逻辑,而是倾向于重用ScheduledThreadPoolExecutor …
java multithreading threadpoolexecutor java-9 completable-future
java ×9
concurrency ×3
threadpool ×3
android ×1
asynchronous ×1
java-8 ×1
java-9 ×1
scala ×1
thread-local ×1