标签: java.util.concurrent

ConcurrentHashMap有什么缺点吗?

我需要一个可以从多个线程访问的HashMap.

有两个简单的选项,使用普通的HashMap并在其上进行同步或使用ConcurrentHashMap.

由于ConcurrentHashMap不会阻止读取操作,因此它似乎更适合我的需求(几乎完全是读取,几乎从不更新).另一方面,我希望无论如何都要求非常低的并发性,因此应该没有阻塞(只是管理锁的成本).

地图也将非常小(十个条目下),如果这有所不同.

与常规HashMap相比,读写操作的成本要高得多(我假设它们是多少)?或者,当可能存在中等级别的并发访问时,ConcurrentHashMap总是更好,无论读取/更新比率和大小如何?

java multithreading hashmap concurrenthashmap java.util.concurrent

5
推荐指数
1
解决办法
8575
查看次数

带有wait()和notify()的ConcurrentLinkedQueue

我并不精通多线程.我试图通过一个生产者线程重复截取屏幕截图,该线程将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)

java concurrency notify wait java.util.concurrent

5
推荐指数
2
解决办法
5340
查看次数

即使等待条件已经改变,AbstractQueuedSynchronizer.acquireShared也会无限期地等待

我写了一个使用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 java.util.concurrent

5
推荐指数
1
解决办法
143
查看次数

java中的非阻塞缓冲区

在高容量多线程java项目中,我需要实现一个非阻塞缓冲区.

在我的场景中,我有一个每秒接收约20,000个请求的Web层.我需要在某些数据结构(也就是所需的缓冲区)中累积其中一些请求,当它满了时(让我们假设它在包含1000个对象时已满)这些对象应该被序列化为一个文件,该文件将被发送到另一个服务器进一步处理.

实施应该是一个非阻塞的实施.我检查了ConcurrentLinkedQueue,但我不确定它是否适合这项工作.

我认为我需要使用2个队列,一旦第一个被填充,它将被一个新队列替换,并且完整队列("第一个")被交付以进行进一步处理.这是我现在想到的基本想法,但我仍然不知道它是否可行,因为我不确定我是否可以在java中切换指针(为了切换完整队列).

有什么建议?

谢谢

java multithreading buffer java.util.concurrent

5
推荐指数
1
解决办法
686
查看次数

JDK7中的ConcurrentHashMap代码说明(scanAndLockForPut)

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)

concurrency concurrenthashmap java.util.concurrent

5
推荐指数
1
解决办法
227
查看次数

CompletableFuture-汇总未来的快速失败

我一直在使用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 java.util.concurrent completable-future

5
推荐指数
1
解决办法
1748
查看次数

Java 5中的Lock和ReentrantLock有什么区别?

我不明白他们之间的区别.我认为来自锁定界面的锁也是可重入的,那么它们之间的区别是什么?你什么时候用?

java concurrency multithreading java.util.concurrent reentrantlock

5
推荐指数
1
解决办法
2888
查看次数

Redis 出队时处于“停车等待”状态的线程

我有一个运行多个线程的 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

5
推荐指数
1
解决办法
1420
查看次数

Elasticsearch RestHighLevelClient(6.0.0)-运行时发生TimeoutException

在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

5
推荐指数
1
解决办法
550
查看次数

在REPL中的Scala中具有java.util.concurrent._的死锁

我在学习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

5
推荐指数
1
解决办法
148
查看次数