我总是对线程感到困惑,而我的班级现在大量使用它们.我们正在使用java.util.concurrent,但我甚至没有真正了解基础知识.UpDownLatch,Futures,Executors; 这些话只是飞过我的脑海.你们可以建议任何资源来帮助我们从头开始学习我需要的东西吗?
非常感谢提前!
我试图在BlockingQueue上使用迭代器方法并发现hasNext()是非阻塞的 - 即它不会等到添加更多元素,而是在没有元素时返回false.
所以这里是问题:
这是一个示例代码块
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) 我在模拟系统上工作,在每个时间步,我必须模拟许多模型.我使用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
我正在计算等待串行事件发生超时的未来:
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) 现在我正在从事有关数据格式转换的工作.有一个大文件,比如10GB,我实现的当前解决方案是逐行读取这个文件,转换每行的格式,然后输出到输出文件.我发现变换过程是一个瓶颈.所以我试图以同时的方式做到这一点.
每条线都是一个完整的单元,与其他线路无关.由于线路中的某些特定值不满足需求,因此可能会丢弃某些线路.
现在我有两个计划:
一个线程从输入文件中逐行读取数据,然后将该行放入队列,几个线程从队列中获取行,转换格式,然后将行放入输出队列,最后输出线程从输出队列中读取行并写入输出文件.
几个线程当前从输入文件的不同部分读取数据,然后处理该行并通过输出队列或文件锁输出到文件.
你们能给我一些建议吗?对此,我真的非常感激.
提前致谢!
我尝试使用ForkJoinPool 来并行化我的CPU密集型计算.我对ForkJoinPool的理解是,只要任何任务可以执行,它就会继续工作.不幸的是,我经常观察工作线程空闲/等待,因此并非所有CPU都保持忙碌状态.有时我甚至观察到额外的工作线程.
我没想到这一点,因为我严格尝试使用非阻塞任务.我的观察非常类似于ForkJoinPool似乎浪费了一个线程.在对ForkJoinPool进行了大量调试之后我猜了一下:
我使用invokeAll()在子任务列表上分配工作.在invokeAll()完成后执行第一个任务本身,它开始加入其他任务.这很好,直到下一个要连接的任务位于执行队列之上.不幸的是,我提交了异步的其他任务而没有加入它们.我期望ForkJoin框架首先继续执行这些任务,然后再转回加入任何剩余的任务.
但它似乎不是这样工作的.相反,工作线程停止调用wait()直到等待的任务准备好(可能是由其他工作线程执行).我没有验证这一点,但似乎是调用join()的一般缺陷.
ForkJoinPool提供了一个asyncMode,但这是一个全局参数,不能用于单个提交.但我喜欢看到我的异步分叉任务很快就会被执行.
那么,为什么ForkJoinTask.doJoin()不是简单地在其队列之上执行任何可用任务,直到它准备好(由自己执行或被其他人窃取)?
从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
我正在寻找逻辑上等同于以下问题的指导:
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
我写了一个基于随机性生成迷宫的函数.大多数时候,这个功能非常快.但每隔一段时间,由于随机数运气不好,需要几秒钟.
我想多次并行启动这个功能,让最快的功能"赢".
Scala标准库(或Java标准库)是否为此作业提供了合适的工具?
我正在使用ThreadSafeList,因为它将数据从进程流式传输到Web服务器,然后在将数据传入客户端时将数据流回流,因此我获得了很大的收益.在内存中我使用Spring Caching(引擎盖下的ehcache)将数据保存在JVM中,一切都很顺利.当我开始达到我的堆限制并且Spring Caching在我使用它时将我的ThreadSafeList序列化到磁盘时,麻烦就开始了,导致了ConcurrentModificationExceptions.我可以覆盖Serialization接口的私有writeObject和readObject方法来解决问题吗?我不确定如何做到这一点或我是否应该放弃我的ThreadSafeList.
回到我开始这个程序的时候,我使用的是BlockingDeque,但这还不够,因为当我放置并采用结构时,我记不起用于缓存的数据......我不能使用ConcurrentMap因为我需要订购在我的列表中...我应该去ConcurrentNavigableMap吗?我想用ThreadSafeList滚动自己,自定义私有序列化功能可能是浪费?
java multithreading java.util.concurrent concurrentmodification spring-cache
java ×10
concurrency ×3
collections ×1
fork-join ×1
lock-free ×1
locking ×1
scala ×1
spring-cache ×1
threadpool ×1