我需要一个可以从多个线程访问的HashMap.
有两个简单的选项,使用普通的HashMap并在其上进行同步或使用ConcurrentHashMap.
由于ConcurrentHashMap不会阻止读取操作,因此它似乎更适合我的需求(几乎完全是读取,几乎从不更新).另一方面,我希望无论如何都要求非常低的并发性,因此应该没有阻塞(只是管理锁的成本).
地图也将非常小(十个条目下),如果这有所不同.
与常规HashMap相比,读写操作的成本要高得多(我假设它们是多少)?或者,当可能存在中等级别的并发访问时,ConcurrentHashMap总是更好,无论读取/更新比率和大小如何?
java multithreading hashmap concurrenthashmap java.util.concurrent
我并不精通多线程.我试图通过一个生产者线程重复截取屏幕截图,该线程将BufferedImage对象添加到,ConcurrentLinkedQueue并且消费者线程将为对象poll排队BufferedImage以将它们保存在文件中.我可以通过重复轮询(while循环)来消耗它们,但我不知道如何使用notify()和消耗它们wait().我曾尝试使用wait(),并notify在较小的项目,但不能在这里实现它.
我有以下代码:
class StartPeriodicTask implements Runnable {
public synchronized void run() {
Robot robot = null;
try {
robot = new Robot();
} catch (AWTException e1) {
e1.printStackTrace();
}
Rectangle screenRect = new Rectangle(Toolkit.getDefaultToolkit()
.getScreenSize());
BufferedImage image = robot.createScreenCapture(screenRect);
if(null!=queue.peek()){
try {
System.out.println("Empty queue, so waiting....");
wait();
} catch (InterruptedException e) {
e.printStackTrace();
}
}else{
queue.add(image);
notify();
}
}
}
public class ImageConsumer implements …Run Code Online (Sandbox Code Playgroud) 我写了一个使用AbstractQueuedSynchronizer的简单类.我写了一个代表"门"的类,如果打开则可以传递,如果关闭则可以阻塞.这是代码:
public class GateBlocking {
final class Sync extends AbstractQueuedSynchronizer {
public Sync() {
setState(0);
}
@Override
protected int tryAcquireShared(int ignored) {
return getState() == 1 ? 1 : -1;
}
public void reset(int newState) {
setState(newState);
}
};
private Sync sync = new Sync();
public void open() {
sync.reset(1);
}
public void close() {
sync.reset(0);
}
public void pass() throws InterruptedException {
sync.acquireShared(1);
}
};
Run Code Online (Sandbox Code Playgroud)
不幸的是,如果一个线程阻塞传递方法,因为门被关闭而另一些线程同时打开门,被阻塞的线程不会被中断 - 它无限地阻塞.这是一个测试,显示它:
public class GateBlockingTest {
@Test
public void parallelPassClosedAndOpenGate() throws Exception{
final …Run Code Online (Sandbox Code Playgroud) 在高容量多线程java项目中,我需要实现一个非阻塞缓冲区.
在我的场景中,我有一个每秒接收约20,000个请求的Web层.我需要在某些数据结构(也就是所需的缓冲区)中累积其中一些请求,当它满了时(让我们假设它在包含1000个对象时已满)这些对象应该被序列化为一个文件,该文件将被发送到另一个服务器进一步处理.
实施应该是一个非阻塞的实施.我检查了ConcurrentLinkedQueue,但我不确定它是否适合这项工作.
我认为我需要使用2个队列,一旦第一个被填充,它将被一个新队列替换,并且完整队列("第一个")被交付以进行进一步处理.这是我现在想到的基本想法,但我仍然不知道它是否可行,因为我不确定我是否可以在java中切换指针(为了切换完整队列).
有什么建议?
谢谢
JDK7中ConcurrentHashMap中的scanAndLockForPut方法的源代码说:
private HashEntry<K,V> scanAndLockForPut(K key, int hash, V value) {
HashEntry<K,V> first = entryForHash(this, hash);
HashEntry<K,V> e = first;
HashEntry<K,V> node = null;
int retries = -1; // negative while locating node
while (!tryLock()) {
HashEntry<K,V> f; // to recheck first below
if (retries < 0) {
if (e == null) {
if (node == null) // speculatively create node
node = new HashEntry<K,V>(hash, key, value, null);
retries = 0;
}
else if (key.equals(e.key))
retries = 0;
else …Run Code Online (Sandbox Code Playgroud) 我一直在使用CompletableFuture.allOf(...)帮助程序来创建汇总期货,这些期货只有在其复合期货标记为已完成时才会变为“完成”,即:
CompletableFuture<?> future1 = new CompletableFuture<>();
CompletableFuture<?> future2 = new CompletableFuture<>();
CompletableFuture<?> future3 = new CompletableFuture<>();
CompletableFuture<?> future = CompletableFuture.allOf(future1, future2, future3);
Run Code Online (Sandbox Code Playgroud)
我希望对此功能稍作更改,在以下情况下,总的未来市场是完整的:
在后一种情况下,合计期货应立即(例外)完成,而不必等待其他期货完成(即快速失败)。
为了说明这一点,请CompletableFuture.allOf(...)考虑一下:
// First future completed, gotta wait for the rest of them...
future1.complete(null);
System.out.println("Future1 Complete, aggregate status: " + future.isDone());
// Second feature was erroneous! I'd like the aggregate to now be completed with failure
future2.completeExceptionally(new Exception());
System.out.println("Future2 Complete, aggregate status: " + future.isDone());
// Finally complete …Run Code Online (Sandbox Code Playgroud) 我不明白他们之间的区别.我认为来自锁定界面的锁也是可重入的,那么它们之间的区别是什么?你什么时候用?
java concurrency multithreading java.util.concurrent reentrantlock
我有一个运行多个线程的 tomcat - spring4.2 应用程序。每个线程仅从一个队列中出队,但是分配给一个队列的线程不止一个。
事情开始很好,但是在几个小时/大约 50 万个出队操作之后,我发现线程出队的速度非常慢。
在 jvisualvm 中看到橙色的线程 ie park 线程转储如下:
"EMLT_2" - Thread t@64
java.lang.Thread.State: WAITING
at sun.misc.Unsafe.park(Native Method)
- parking to wait for <2cf42d7> (a java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject)
at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2039)
at org.apache.commons.pool2.impl.LinkedBlockingDeque.takeFirst(LinkedBlockingDeque.java:583)
at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:442)
at org.apache.commons.pool2.impl.GenericObjectPool.borrowObject(GenericObjectPool.java:363)
at redis.clients.util.Pool.getResource(Pool.java:48)
at redis.clients.jedis.JedisPool.getResource(JedisPool.java:86)
at com.mycomp.sam.processors.SimpleDequeuer.dequeue(SimpleDequeuer.java:25)
at com.mycomp.sam.processors.EMLT.run(EMLT.java:29)
at java.lang.Thread.run(Thread.java:745)
Locked ownable synchronizers:
- None
"EMLT_1" - Thread t@63
java.lang.Thread.State: WAITING
at sun.misc.Unsafe.park(Native Method)
- parking to wait for <2cf42d7> (a java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject)
at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175)
at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(AbstractQueuedSynchronizer.java:2039)
at org.apache.commons.pool2.impl.LinkedBlockingDeque.takeFirst(LinkedBlockingDeque.java:583)
at …Run Code Online (Sandbox Code Playgroud) concurrency locking java.util.concurrent jedis apache-commons-pool
在Elasticsearch服务器中索引(存储)文档时,我们的应用程序encoutner超时异常。它不经常发生,但大约每天一次。这是详细信息。
pom.xml
<dependency>
<groupId>org.elasticsearch.client</groupId>
<artifactId>elasticsearch-rest-high-level-client</artifactId>
<version>6.0.0</version>
</dependency>
Run Code Online (Sandbox Code Playgroud)
初始化
@Bean
public RestHighLevelClient buildHighLevelClient() {
RestHighLevelClient client = new RestHighLevelClient(RestClient.builder(httplist.toArray(new HttpHost[]{})));
return client;
}
Run Code Online (Sandbox Code Playgroud)
异常信息
java.lang.RuntimeException: error while performing request
at org.elasticsearch.client.RestClient$SyncResponseListener.get(RestClient.java:682)
at org.elasticsearch.client.RestClient.performRequest(RestClient.java:220)
at org.elasticsearch.client.RestClient.performRequest(RestClient.java:192)
at org.elasticsearch.client.RestHighLevelClient.performRequest(RestHighLevelClient.java:428)
at org.elasticsearch.client.RestHighLevelClient.performRequestAndParseEntity(RestHighLevelClient.java:414)
at org.elasticsearch.client.RestHighLevelClient.bulk(RestHighLevelClient.java:229)
......
Caused by: java.util.concurrent.TimeoutException
at org.apache.http.nio.pool.AbstractNIOConnPool.processPendingRequest(AbstractNIOConnPool.java:364)
at org.apache.http.nio.pool.AbstractNIOConnPool.processNextPendingRequest(AbstractNIOConnPool.java:344)
at org.apache.http.nio.pool.AbstractNIOConnPool.release(AbstractNIOConnPool.java:318)
at org.apache.http.impl.nio.conn.PoolingNHttpClientConnectionManager.releaseConnection(PoolingNHttpClientConnectionManager.java:303)
at org.apache.http.impl.nio.client.AbstractClientExchangeHandler.releaseConnection(AbstractClientExchangeHandler.java:239)
Run Code Online (Sandbox Code Playgroud) timeoutexception java.util.concurrent rest-client elasticsearch
我在学习Paul Chiusano和Runar Bjanarson的著作“ Scala中的函数编程”(第7章-纯函数并行性)时遇到了以下情况。
package fpinscala.parallelism
import java.util.concurrent._
import language.implicitConversions
object Par {
type Par[A] = ExecutorService => Future[A]
def run[A](s: ExecutorService)(a: Par[A]): Future[A] = a(s)
def unit[A](a: A): Par[A] = (es: ExecutorService) => UnitFuture(a) // `unit` is represented as a function that returns a `UnitFuture`, which is a simple implementation of `Future` that just wraps a constant value. It doesn't use the `ExecutorService` at all. It's always done and can't be cancelled. Its `get` method simply returns the value …Run Code Online (Sandbox Code Playgroud) java parallel-processing scala java.util.concurrent scala-repl
java ×7
concurrency ×4
buffer ×1
hashmap ×1
jedis ×1
locking ×1
notify ×1
rest-client ×1
scala ×1
scala-repl ×1
wait ×1