标签: java.util.concurrent

如何了解Threads,尤其是Java

我总是对线程感到困惑,而我的班级现在大量使用它们.我们正在使用java.util.concurrent,但我甚至没有真正了解基础知识.UpDownLatch,Futures,Executors; 这些话只是飞过我的脑海.你们可以建议任何资源来帮助我们从头开始学习我需要的东西吗?

非常感谢提前!

java multithreading java.util.concurrent

6
推荐指数
2
解决办法
4162
查看次数

为什么iterator.hasNext不能与BlockingQueue一起使用?

我试图在BlockingQueue上使用迭代器方法并发现hasNext()是非阻塞的 - 即它不会等到添加更多元素,而是在没有元素时返回false.

所以这里是问题:

  1. 这是糟糕的设计还是错误的期望?
  2. 有没有办法使用BLockingQueue的阻塞方法及其父类Collection方法(例如,如果某些方法需要一个集合,我可以传递一个阻塞队列,并希望它的处理将等到Queue有更多元素)

这是一个示例代码块

public class SomeContainer{
     public static void main(String[] args){
        BlockingQueue bq = new LinkedBlockingQueue();
        SomeContainer h = new SomeContainer();
        Producer p = new Producer(bq);
        Consumer c = new Consumer(bq);
        p.produce();
        c.consume();
    }

    static class Producer{
        BlockingQueue q;
        public Producer(BlockingQueue q) {
            this.q = q;
        }

        void produce(){
        new Thread(){
            public void run() {
            for(int i=0; i<10; i++){
                for(int j=0;j<10; j++){
                    q.add(i+" - "+j);
                }
                try {
                    Thread.sleep(30000);
                } catch (InterruptedException e) {
                    e.printStackTrace();
                } …
Run Code Online (Sandbox Code Playgroud)

java collections multithreading java.util.concurrent

6
推荐指数
2
解决办法
9920
查看次数

java重用执行器

我在模拟系统上工作,在每个时间步,我必须模拟许多模型.我使用FixedThreadPool来加速计算:

ExecutorService executor = Executors.newFixedThreadPool(nThread);
for (Model m : models) {
  executor.execute( m.simulationTask() );
}
executor.shutdown();
while ( ! executor.awaitTermination(10, TimeUnit.MINUTES) ) { 
  System.out.println("wait"); 
}
Run Code Online (Sandbox Code Playgroud)

现在,执行程序execute()在调用后不能用于新任务shutdown().有没有办法重置执行程序,所以我可以在下一个模拟步骤中重用现有的执行程序(及其线程)?

java multithreading java.util.concurrent concurrent-programming

6
推荐指数
1
解决办法
6540
查看次数

如何停止未来的超时

我正在计算等待串行事件发生超时的未来:

Future<Response> future = executor.submit(new CommunicationTask(this, request));
response = new Response("timeout");
try {
  response = future.get(timeoutMilliseconds, TimeUnit.MILLISECONDS);
} catch (InterruptedException | TimeoutException e) {
  future.cancel(true);
  log.info("Execution time out." + e);
} catch (ExecutionException e) {
  future.cancel(true);
  log.error("Encountered problem communicating with device: " + e);
}
Run Code Online (Sandbox Code Playgroud)

CommunicationTask班实施了Observer监听来自串行端口的变化接口.

问题是从串口读取速度相对较慢,即使发生串行事件,时间耗尽TimeoutException也会抛出a.当串行事件发生时,我该怎么做才能停止未来的超时时钟?

我试了一下,AtomicReference但没有改变任何东西:

public class CommunicationTask implements Callable<Response>, Observer {
  private AtomicReference atomicResponse = new AtomicReference(new Response("timeout"));
  private CountDownLatch latch = new CountDownLatch(1);
  private SerialPort port;

  CommunicationTask(SerialCommunicator …
Run Code Online (Sandbox Code Playgroud)

java java.util.concurrent

6
推荐指数
1
解决办法
1150
查看次数

以并发方式逐行处理文件

现在我正在从事有关数据格式转换的工作.有一个大文件,比如10GB,我实现的当前解决方案是逐行读取这个文件,转换每行的格式,然后输出到输出文件.我发现变换过程是一个瓶颈.所以我试图以同时的方式做到这一点.

每条线都是一个完整的单元,与其他线路无关.由于线路中的某些特定值不满足需求,因此可能会丢弃某些线路.

现在我有两个计划:

  1. 一个线程从输入文件中逐行读取数据,然后将该行放入队列,几个线程从队列中获取行,转换格式,然后将行放入输出队列,最后输出线程从输出队列中读取行并写入输出文件.

  2. 几个线程当前从输入文件的不同部分读取数据,然后处理该行并通过输出队列或文件锁输出到文件.

你们能给我一些建议吗?对此,我真的非常感激.

提前致谢!

java concurrency java.util.concurrent

6
推荐指数
1
解决办法
508
查看次数

ForkJoinPool在invokeAll/join期间停止

我尝试使用ForkJoinPool 来并行化我的CPU密集型计算.我对ForkJoinPool的理解是,只要任何任务可以执行,它就会继续工作.不幸的是,我经常观察工作线程空闲/等待,因此并非所有CPU都保持忙碌状态.有时我甚至观察到额外的工作线程.

我没想到这一点,因为我严格尝试使用非阻塞任务.我的观察非常类似于ForkJoinPool似乎浪费了一个线程.在对ForkJoinPool进行了大量调试之后我猜了一下:

我使用invokeAll()在子任务列表上分配工作.在invokeAll()完成后执行第一个任务本身,它开始加入其他任务.这很好,直到下一个要连接的任务位于执行队列之上.不幸的是,我提交了异步的其他任务而没有加入它们.我期望ForkJoin框架首先继续执行这些任务,然后再转回加入任何剩余的任务.

但它似乎不是这样工作的.相反,工作线程停止调用wait()直到等待的任务准备好(可能是由其他工作线程执行).我没有验证这一点,但似乎是调用join()的一般缺陷.

ForkJoinPool提供了一个asyncMode,但这是一个全局参数,不能用于单个提交.但我喜欢看到我的异步分叉任务很快就会被执行.

那么,为什么ForkJoinTask.doJoin()不是简单地在其队列之上执行任何可用任务,直到它准备好(由自己执行或被其他人窃取)?

java lock-free java.util.concurrent fork-join

6
推荐指数
1
解决办法
3072
查看次数

为什么在ReentrantReadAndWriteLock中,readLock()应该在writeLock.lock()之前解锁?

ReentrantReadAndWriteLock类javadoc:

void processCachedData() {
  rwl.readLock().lock();
  if (!cacheValid) {
     // Must release read lock before acquiring write lock
5:   rwl.readLock().unlock();
6:   rwl.writeLock().lock();
     // Recheck state because another thread might have acquired
     //   write lock and changed state before we did.
     if (!cacheValid) {
       data = ...
       cacheValid = true;
     }
     // Downgrade by acquiring read lock before releasing write lock
14:  rwl.readLock().lock();
15:  rwl.writeLock().unlock(); // Unlock write, still hold read
  }

  use(data);
  rwl.readLock().unlock();
}
Run Code Online (Sandbox Code Playgroud)

为什么我们必须在获取注释中写入的写锁之前释放读锁?如果当前线程持有读锁定,那么当其他线程不再读取时,应该允许它获取写锁定,无论当前线程是否还保持读锁定.这是我期望的行为.最终在第4行和第5行锁定升级并在第14行和第15行锁定降级我希望在ReentrantReadAndWriteLock类内部完成.为什么那是不可能的?

换句话说,我希望代码能够正常工作:

void processCachedData() { …
Run Code Online (Sandbox Code Playgroud)

java concurrency multithreading locking java.util.concurrent

6
推荐指数
2
解决办法
892
查看次数

如何在两个返回布尔值的并行线程上用Java进行短路评估?

我正在寻找逻辑上等同于以下问题的指导:

public boolean parallelOR() {
    ExecutorService executor = Executors.newFixedThreadPool(2);
    Future<Boolean> taskA = executor.submit( new SlowTaskA() );
    Future<Boolean> taskB = executor.submit( new SlowTaskB() );

    return taskA.get() || taskB.get(); // This is not what I want
    // Exception handling omitted for clarity
 }
Run Code Online (Sandbox Code Playgroud)

上述结构给出了正确的结果,但是即使从taskB完成后结果已知,总是等待taskA完成.

有没有更好的结构,如果任何一个线程返回true而不等待第二个线程完成,它将允许返回一个值?

(涉及的平台是Android,如果这会影响结果).

java lazy-evaluation short-circuiting java.util.concurrent threadpool

6
推荐指数
1
解决办法
247
查看次数

并行执行竞争计算并丢弃除完成的第一个之外的所有计算

我写了一个基于随机性生成迷宫的函数.大多数时候,这个功能非常快.但每隔一段时间,由于随机数运气不好,需要几秒钟.

我想多次并行启动这个功能,让最快的功能"赢".

Scala标准库(或Java标准库)是否为此作业提供了合适的工具?

java concurrency multithreading scala java.util.concurrent

6
推荐指数
2
解决办法
201
查看次数

如何使列表线程安全的序列化?

我正在使用ThreadSafeList,因为它将数据从进程流式传输到Web服务器,然后在将数据传入客户端时将数据流回流,因此我获得了很大的收益.在内存中我使用Spring Caching(引擎盖下的ehcache)将数据保存在JVM中,一切都很顺利.当我开始达到我的堆限制并且Spring Caching在我使用它时将我的ThreadSafeList序列化到磁盘时,麻烦就开始了,导致了ConcurrentModificationExceptions.我可以覆盖Serialization接口的私有writeObject和readObject方法来解决问题吗?我不确定如何做到这一点或我是否应该放弃我的ThreadSafeList.

回到我开始这个程序的时候,我使用的是BlockingDeque,但这还不够,因为当我放置并采用结构时,我记不起用于缓存的数据......我不能使用ConcurrentMap因为我需要订购在我的列表中...我应该去ConcurrentNavigableMap吗?我想用ThreadSafeList滚动自己,自定义私有序列化功能可能是浪费?

Java Code Geeks ThreadSafeList

java multithreading java.util.concurrent concurrentmodification spring-cache

6
推荐指数
1
解决办法
387
查看次数