From c9c1467878a1add4b22cb2deb2cedd773ab04ffa Mon Sep 17 00:00:00 2001 From: Test Date: Wed, 30 Sep 2026 14:40:43 +0000 Subject: [PATCH 1/3] fix(ENGKNOW-3933): MAP/MULTIMAP cache key, single load per request and lower memory - Cache key did not include caseInsensitive and had no separator between ic and oc, so MAP -cis and plain MAP of the same file shared one cached map, and ic=1,oc=[23] collided with ic=12,oc=[3]. Use a delimited key with all parameters. - Concurrent pipelines of a request (e.g. pgor partitions) sharing the session cache each built their own copy of the same map. Load under a per request and key lock, re-checking the cache after acquiring it. - Share repeated value strings while loading a map (bounded pool), map value columns repeat heavily. - MULTIMAP: move entries into a presized output map instead of copying, so both maps are not fully held at the same time. Measured heap retained after GC (1M line files, values from 50 distinct strings): single map 160.8 MB -> 98.4 MB, multimap 90.4 MB -> 28.0 MB, multimap peak during load 350.1 MB -> 305.9 MB. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Utilities/UTestMapAndListUtilities.scala | 157 ++++++++++++++++++ .../MapAndListUtilities.scala | 90 +++++++--- 2 files changed, 228 insertions(+), 19 deletions(-) create mode 100644 gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala diff --git a/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala new file mode 100644 index 000000000..f51d93e50 --- /dev/null +++ b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala @@ -0,0 +1,157 @@ +/* + * BEGIN_COPYRIGHT + * + * Copyright (C) 2011-2013 deCODE genetics Inc. + * Copyright (C) 2013-2019 WuXi NextCode Inc. + * All Rights Reserved. + * + * GORpipe is free software: you can redistribute it and/or modify + * it under the terms of the AFFERO GNU General Public License as published by + * the Free Software Foundation. + * + * GORpipe is distributed "AS-IS" AND WITHOUT ANY WARRANTY OF ANY KIND, + * INCLUDING ANY IMPLIED WARRANTY OF MERCHANTABILITY, + * NON-INFRINGEMENT, OR FITNESS FOR A PARTICULAR PURPOSE. See + * the AFFERO GNU General Public License for the complete license terms. + * + * You should have received a copy of the AFFERO GNU General Public License + * along with GORpipe. If not, see + * + * END_COPYRIGHT + */ + +package gorsat.Utilities + +import java.io.{File, PrintWriter} +import java.util.concurrent.atomic.AtomicInteger +import java.util.concurrent.{Callable, CountDownLatch, Executors, TimeUnit} +import gorsat.gorsatGorIterator.MapAndListUtilities +import gorsat.process.GenericSessionFactory +import org.gorpipe.gor.model.Row +import org.gorpipe.model.gor.iterators.LineIterator +import org.junit.runner.RunWith +import org.scalatest.funsuite.AnyFunSuite +import org.scalatestplus.junit.JUnitRunner + +/** + * Tests for the map/multimap loading and caching in MapAndListUtilities (ENGKNOW-3933). + */ +@RunWith(classOf[JUnitRunner]) +class UTestMapAndListUtilities extends AnyFunSuite { + + private def file(lines: String*): File = { + val file = File.createTempFile("UTestMapAndListUtilities", ".tsv") + file.deleteOnExit() + val writer = new PrintWriter(file) + try lines.foreach(writer.println) finally writer.close() + file + } + + /** Line iterator over in memory lines that counts the lines read, optionally waiting before the first line. */ + private class CountingLineIterator(lines: Seq[String], counter: AtomicInteger, beforeFirst: () => Unit = () => ()) + extends LineIterator { + private val it = lines.iterator + private var first = true + + override def nextLine: String = { + if (first) { + first = false + beforeFirst() + } + counter.incrementAndGet() + it.next() + } + + override def hasNext: Boolean = it.hasNext + + override def next(): Row = throw new UnsupportedOperationException + + override def close(): Unit = {} + } + + test("Case insensitive and case sensitive map of the same file are cached separately") { + val input = file("Abc\tv1") + val session = new GenericSessionFactory().create() + + val sensitive = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = false, 1, Array(1), + asSet = false, skipEmpty = false, session) + val insensitive = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = true, 1, Array(1), + asSet = false, skipEmpty = false, session) + + assert(sensitive.containsKey("Abc")) + assert(insensitive.containsKey("ABC")) + } + + test("Case insensitive and case sensitive multimap of the same file are cached separately") { + val input = file("Abc\tv1") + val session = new GenericSessionFactory().create() + + val sensitive = MapAndListUtilities.getMultiHashMap(input.getAbsolutePath, caseInsensitive = false, session) + val insensitive = MapAndListUtilities.getMultiHashMap(input.getAbsolutePath, caseInsensitive = true, session) + + assert(sensitive.containsKey("Abc")) + assert(insensitive.containsKey("ABC")) + } + + test("Maps with different input and output columns do not share a cache entry") { + // ic=1,oc=[23] and ic=12,oc=[3] used to produce the same cache key. + val input = file((0 until 24).map(i => s"c$i").mkString("\t")) + val session = new GenericSessionFactory().create() + + val first = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = false, 1, Array(23), + asSet = false, skipEmpty = false, session) + val second = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = false, 12, Array(3), + asSet = false, skipEmpty = false, session) + + assert(first.get("c0") == "c23") + assert(second.get((0 until 12).map(i => s"c$i").mkString("\t")) == "c3") + } + + test("Concurrent loads of the same map in a session read the file once") { + val lines = (1 to 1000).map(i => s"k$i\tv${i % 10}") + val session = new GenericSessionFactory().create() + val linesRead = new AtomicInteger() + val threads = 4 + val allStarted = new CountDownLatch(threads) + val executor = Executors.newFixedThreadPool(threads) + try { + val results = (1 to threads).map { _ => + executor.submit(new Callable[java.util.Map[String, String]] { + override def call(): java.util.Map[String, String] = { + allStarted.countDown() + val iterator = new CountingLineIterator(lines, linesRead, () => Thread.sleep(200)) + MapAndListUtilities.getSingleHashMap("concurrent.tsv", iterator, caseInsensitive = false, 1, Array(1), + asSet = false, skipEmpty = false, session) + } + }) + }.map(_.get(30, TimeUnit.SECONDS)) + + assert(linesRead.get() == lines.size) + results.foreach(map => assert(map eq results.head)) + } finally { + executor.shutdownNow() + } + } + + test("Repeated map values share one string instance") { + val input = file("k1\tsame", "k2\tsame", "k3\tother") + val session = new GenericSessionFactory().create() + + val map = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, asSet = false, skipEmpty = false, session) + + assert(map.get("k1") == "same") + assert(map.get("k1") eq map.get("k2")) + assert(map.get("k3") == "other") + } + + test("Repeated multimap values share one string instance and keep their order") { + val input = file("k1\tsame", "k1\tother", "k2\tsame") + val session = new GenericSessionFactory().create() + + val map = MapAndListUtilities.getMultiHashMap(input.getAbsolutePath, caseInsensitive = false, session) + + assert(map.get("k1").toSeq == Seq("same", "other")) + assert(map.get("k2").toSeq == Seq("same")) + assert(map.get("k1")(0) eq map.get("k2")(0)) + } +} diff --git a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala index 90ba2df1a..7557c966c 100644 --- a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala +++ b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala @@ -23,6 +23,7 @@ package gorsat.gorsatGorIterator import java.nio.file.Files +import java.util.concurrent.ConcurrentHashMap import java.util.stream.Collectors import org.gorpipe.gor.model.{DriverBackedFileReader, FileReader} import org.gorpipe.gor.session.GorSession @@ -112,17 +113,67 @@ object MapAndListUtilities { } } + // Locks so that concurrent pipelines of a request, which share the session cache, load a given map only once. + private val loadLocks = new ConcurrentHashMap[String, Object]() + + private def cacheKey(kind: String, filename: String, ic: Int, oc: Array[Int], asSet: Boolean, + caseInsensitive: Boolean): String = + s"$kind|$filename|$ic|${oc.mkString(",")}|$asSet|$caseInsensitive" + + /** + * Returns the cached value if present, otherwise loads it while holding a lock for the request and key, so + * concurrent callers wait for the first load instead of each building their own copy. + */ + private def loadOnce[T](extFilename: String, iterator: LineIterator, session: GorSession) + (cached: => Option[T])(load: => T): T = { + cached match { + case Some(value) => + iterator.close() + value + case None => + val lockKey = String.valueOf(session.getRequestId) + "|" + extFilename + val lock = loadLocks.computeIfAbsent(lockKey, _ => new Object) + try { + lock.synchronized { + cached match { + case Some(value) => + iterator.close() + value + case None => load + } + } + } finally { + loadLocks.remove(lockKey, lock) + } + } + } + + /** + * Shares equal strings while one map is loaded, map value columns tend to repeat a lot. Bounded so that maps with + * mostly unique values do not pay for a large pool. + */ + private class StringPool(maxSize: Int = 100000) { + private val pool = new java.util.HashMap[String, String]() + + def apply(s: String): String = { + val existing = pool.get(s) + if (existing != null) { + existing + } else { + if (pool.size < maxSize) pool.put(s, s) + s + } + } + } + def getSingleHashMap(filename: String, iterator: LineIterator, caseInsensitive: Boolean, ic: Int, oc: Array[Int], asSet: Boolean, skipEmpty: Boolean, session: GorSession): singleHashMap = { - val extFilename = "map" + filename + ic + oc.mkString(",") + asSet + val extFilename = cacheKey("map", filename, ic, oc, asSet, caseInsensitive) val ocl = oc.length - syncGetSingleHashMap(extFilename, session) match { - case Some(theMap) => - iterator.close() - theMap - case None => + loadOnce(extFilename, iterator, session)(syncGetSingleHashMap(extFilename, session)) { try { val colMap = new java.util.HashMap[String, String]() + val values = new StringPool() val mmu: MemoryMonitorUtil = new MemoryMonitorUtil(MemoryMonitorUtil.basicOutOfMemoryHandler) @@ -142,11 +193,11 @@ object MapAndListUtilities { if (caseInsensitive) cols.slice(0, ic).mkString("\t").toUpperCase else cols.slice(0, ic).mkString("\t") if (colMap.getOrDefault(lookupString,null) == null) { - colMap.put(lookupString, oc.tail.map(c => cols(c)).foldLeft(cols(oc.head))(_ + "\t" + _)) + colMap.put(lookupString, values(oc.tail.map(c => cols(c)).foldLeft(cols(oc.head))(_ + "\t" + _))) } else { val existingValues = colMap.get(lookupString).split("\t",-1) val newValues = if( skipEmpty ) existingValues.zip(oc.map(c => cols(c))).map(_.productIterator.filter(_.toString.nonEmpty).mkString(",")) else existingValues.zip(oc.map(c => cols(c))).map(x => x._1 + "," + x._2 ) - colMap.put(lookupString, newValues.tail.foldLeft(newValues.head)(_ + "\t" + _)) + colMap.put(lookupString, values(newValues.tail.foldLeft(newValues.head)(_ + "\t" + _))) } } } @@ -161,22 +212,19 @@ object MapAndListUtilities { def getMultiHashMap(filename: String, iterator: LineIterator, caseInsensitive: Boolean, ic: Int, oc: Array[Int], session: GorSession): multiHashMap = { - val extFilename = "multimap" + filename + ic + oc.mkString(",") + val extFilename = cacheKey("multimap", filename, ic, oc, asSet = false, caseInsensitive) val ocl = oc.length - syncGetMultiHashMap(extFilename, session) match { - case Some(theMap) => - iterator.close() - theMap - case None => + loadOnce(extFilename, iterator, session)(syncGetMultiHashMap(extFilename, session)) { val multiMap = new java.util.HashMap[String, ListBuffer[String]]() try { + val values = new StringPool() val mmu: MemoryMonitorUtil = new MemoryMonitorUtil(MemoryMonitorUtil.basicOutOfMemoryHandler) while (iterator.hasNext) { val x = iterator.nextLine val cols = x.split("\t", -1) mmu.check("getMultiHashMap", mmu.lineNum, x) if (cols.length >= ic + ocl) { - val (a, b) = (cols.slice(0, ic).mkString("\t"), oc.tail.map(c => cols(c)).foldLeft(cols(oc.head))(_ + "\t" + _)) + val (a, b) = (cols.slice(0, ic).mkString("\t"), values(oc.tail.map(c => cols(c)).foldLeft(cols(oc.head))(_ + "\t" + _))) val cisa = if (caseInsensitive) a.toUpperCase else a if (multiMap.containsKey(cisa)) { multiMap.put(cisa, multiMap.get(cisa) += b) @@ -186,10 +234,14 @@ object MapAndListUtilities { } } } - val multiOutputMap = new java.util.HashMap[String, Array[String]]() - multiMap.forEach((k, v) => { - multiOutputMap.put(k, v.toArray) - }) + // Move the entries over rather than copying them, so both maps are not fully held at the same time. + val multiOutputMap = new java.util.HashMap[String, Array[String]]((multiMap.size / 0.75f).toInt + 1) + val entries = multiMap.entrySet().iterator() + while (entries.hasNext) { + val entry = entries.next() + multiOutputMap.put(entry.getKey, entry.getValue.toArray) + entries.remove() + } syncAddMultiHashMap(extFilename, multiOutputMap, session) multiOutputMap } finally { From 660965e4480030cb23c4f1e5f196b740a32179ff Mon Sep 17 00:00:00 2001 From: Test Date: Thu, 1 Oct 2026 11:45:53 +0000 Subject: [PATCH 2/3] fix(ENGKNOW-3933): include skipEmpty in MAP cache key MAP -e and plain MAP of the same file and columns shared one cache entry, although skipEmpty changes how values of duplicate keys are merged. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../gorsat/Utilities/UTestMapAndListUtilities.scala | 13 +++++++++++++ .../gorsatGorIterator/MapAndListUtilities.scala | 8 ++++---- 2 files changed, 17 insertions(+), 4 deletions(-) diff --git a/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala index f51d93e50..b67a9e101 100644 --- a/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala +++ b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala @@ -107,6 +107,19 @@ class UTestMapAndListUtilities extends AnyFunSuite { assert(second.get((0 until 12).map(i => s"c$i").mkString("\t")) == "c3") } + test("Maps with and without skipEmpty of the same file are cached separately") { + val input = file("k1\ta", "k1\t") + val session = new GenericSessionFactory().create() + + val keepEmpty = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = false, 1, Array(1), + asSet = false, skipEmpty = false, session) + val skipEmpty = MapAndListUtilities.getSingleHashMap(input.getAbsolutePath, caseInsensitive = false, 1, Array(1), + asSet = false, skipEmpty = true, session) + + assert(keepEmpty.get("k1") == "a,") + assert(skipEmpty.get("k1") == "a") + } + test("Concurrent loads of the same map in a session read the file once") { val lines = (1 to 1000).map(i => s"k$i\tv${i % 10}") val session = new GenericSessionFactory().create() diff --git a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala index 7557c966c..56a5cd308 100644 --- a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala +++ b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala @@ -117,8 +117,8 @@ object MapAndListUtilities { private val loadLocks = new ConcurrentHashMap[String, Object]() private def cacheKey(kind: String, filename: String, ic: Int, oc: Array[Int], asSet: Boolean, - caseInsensitive: Boolean): String = - s"$kind|$filename|$ic|${oc.mkString(",")}|$asSet|$caseInsensitive" + caseInsensitive: Boolean, skipEmpty: Boolean): String = + s"$kind|$filename|$ic|${oc.mkString(",")}|$asSet|$caseInsensitive|$skipEmpty" /** * Returns the cached value if present, otherwise loads it while holding a lock for the request and key, so @@ -168,7 +168,7 @@ object MapAndListUtilities { def getSingleHashMap(filename: String, iterator: LineIterator, caseInsensitive: Boolean, ic: Int, oc: Array[Int], asSet: Boolean, skipEmpty: Boolean, session: GorSession): singleHashMap = { - val extFilename = cacheKey("map", filename, ic, oc, asSet, caseInsensitive) + val extFilename = cacheKey("map", filename, ic, oc, asSet, caseInsensitive, skipEmpty) val ocl = oc.length loadOnce(extFilename, iterator, session)(syncGetSingleHashMap(extFilename, session)) { try { @@ -212,7 +212,7 @@ object MapAndListUtilities { def getMultiHashMap(filename: String, iterator: LineIterator, caseInsensitive: Boolean, ic: Int, oc: Array[Int], session: GorSession): multiHashMap = { - val extFilename = cacheKey("multimap", filename, ic, oc, asSet = false, caseInsensitive) + val extFilename = cacheKey("multimap", filename, ic, oc, asSet = false, caseInsensitive, skipEmpty = false) val ocl = oc.length loadOnce(extFilename, iterator, session)(syncGetMultiHashMap(extFilename, session)) { val multiMap = new java.util.HashMap[String, ListBuffer[String]]() From 7aaba85c5ee898b839a4e46634f0afa07676e92b Mon Sep 17 00:00:00 2001 From: Test Date: Thu, 1 Oct 2026 11:51:21 +0000 Subject: [PATCH 3/3] fix(ENGKNOW-3933): share one in-progress map load per session cache Replace the per-request lock in loadOnce with a CompletableFuture per (session cache, key), so callers waiting for a load: - can be interrupted, and close their own iterator instead of holding it - get the result or the failure of that load, rather than each retrying a failed load in turn - only wait for loads into the same session cache, not for unrelated sessions that happen to share a request id Also stop pooling merged values of duplicate map keys, they are replaced by later duplicates, and make the concurrent load test wait for all threads to start. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../Utilities/UTestMapAndListUtilities.scala | 106 +++++++++++++++++- .../MapAndListUtilities.scala | 48 +++++--- 2 files changed, 138 insertions(+), 16 deletions(-) diff --git a/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala index b67a9e101..419aca206 100644 --- a/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala +++ b/gortools/src/test/scala/gorsat/Utilities/UTestMapAndListUtilities.scala @@ -24,15 +24,18 @@ package gorsat.Utilities import java.io.{File, PrintWriter} import java.util.concurrent.atomic.AtomicInteger -import java.util.concurrent.{Callable, CountDownLatch, Executors, TimeUnit} +import java.util.concurrent.{Callable, CountDownLatch, ExecutionException, ExecutorService, Executors, Future, TimeUnit} import gorsat.gorsatGorIterator.MapAndListUtilities import gorsat.process.GenericSessionFactory import org.gorpipe.gor.model.Row +import org.gorpipe.gor.session.{GorSession, GorSessionCache} import org.gorpipe.model.gor.iterators.LineIterator import org.junit.runner.RunWith import org.scalatest.funsuite.AnyFunSuite import org.scalatestplus.junit.JUnitRunner +import scala.util.Try + /** * Tests for the map/multimap loading and caching in MapAndListUtilities (ENGKNOW-3933). */ @@ -69,6 +72,15 @@ class UTestMapAndListUtilities extends AnyFunSuite { override def close(): Unit = {} } + private def async[T](executor: ExecutorService)(body: => T): Future[T] = + executor.submit(new Callable[T] { + override def call(): T = body + }) + + private def loadMap(name: String, iterator: LineIterator, session: GorSession): java.util.Map[String, String] = + MapAndListUtilities.getSingleHashMap(name, iterator, caseInsensitive = false, 1, Array(1), asSet = false, + skipEmpty = false, session) + test("Case insensitive and case sensitive map of the same file are cached separately") { val input = file("Abc\tv1") val session = new GenericSessionFactory().create() @@ -132,6 +144,7 @@ class UTestMapAndListUtilities extends AnyFunSuite { executor.submit(new Callable[java.util.Map[String, String]] { override def call(): java.util.Map[String, String] = { allStarted.countDown() + allStarted.await() val iterator = new CountingLineIterator(lines, linesRead, () => Thread.sleep(200)) MapAndListUtilities.getSingleHashMap("concurrent.tsv", iterator, caseInsensitive = false, 1, Array(1), asSet = false, skipEmpty = false, session) @@ -146,6 +159,97 @@ class UTestMapAndListUtilities extends AnyFunSuite { } } + test("A caller waiting for another load of the same map can be interrupted") { + val session = new GenericSessionFactory().create() + val loading = new CountDownLatch(1) + val release = new CountDownLatch(1) + val waiterDone = new CountDownLatch(1) + val executor = Executors.newFixedThreadPool(2) + try { + val loader = async(executor) { + loadMap("interrupt.tsv", new CountingLineIterator(Seq("k\tv"), new AtomicInteger(), () => { + loading.countDown() + release.await() + }), session) + } + loading.await() + val waiter = async(executor) { + try loadMap("interrupt.tsv", new CountingLineIterator(Seq("k\tv"), new AtomicInteger()), session) + finally waiterDone.countDown() + } + Thread.sleep(200) + waiter.cancel(true) + + assert(waiterDone.await(5, TimeUnit.SECONDS), "waiter did not stop when interrupted") + release.countDown() + assert(loader.get(30, TimeUnit.SECONDS).get("k") == "v") + } finally { + release.countDown() + executor.shutdownNow() + } + } + + test("A failed load fails the callers waiting for it instead of each of them retrying") { + val session = new GenericSessionFactory().create() + val loads = new AtomicInteger() + val threads = 4 + val allStarted = new CountDownLatch(threads) + val executor = Executors.newFixedThreadPool(threads) + try { + val results = (1 to threads).map { _ => + async(executor) { + allStarted.countDown() + allStarted.await() + loadMap("failing.tsv", new CountingLineIterator(Seq("k\tv"), new AtomicInteger(), () => { + loads.incrementAndGet() + Thread.sleep(200) + throw new IllegalStateException("load failed") + }), session) + } + }.map(future => Try(future.get(30, TimeUnit.SECONDS))) + + assert(loads.get() == 1) + results.foreach { result => + val cause = result.failed.get.asInstanceOf[ExecutionException].getCause + assert(cause.getMessage == "load failed") + } + + // A later call loads again. + val map = loadMap("failing.tsv", new CountingLineIterator(Seq("k\tv"), new AtomicInteger()), session) + assert(map.get("k") == "v") + } finally { + executor.shutdownNow() + } + } + + test("Sessions with the same request id but separate caches do not wait for each other") { + val first = new GenericSessionFactory().create() + val second = new GorSession(first.getRequestId) + second.init(first.getProjectContext, first.getSystemContext, new GorSessionCache()) + val loading = new CountDownLatch(1) + val release = new CountDownLatch(1) + val executor = Executors.newFixedThreadPool(2) + try { + val firstLoad = async(executor) { + loadMap("shared-request.tsv", new CountingLineIterator(Seq("k\tv1"), new AtomicInteger(), () => { + loading.countDown() + release.await() + }), first) + } + loading.await() + val secondLoad = async(executor) { + loadMap("shared-request.tsv", new CountingLineIterator(Seq("k\tv2"), new AtomicInteger()), second) + } + + assert(secondLoad.get(5, TimeUnit.SECONDS).get("k") == "v2") + release.countDown() + assert(firstLoad.get(30, TimeUnit.SECONDS).get("k") == "v1") + } finally { + release.countDown() + executor.shutdownNow() + } + } + test("Repeated map values share one string instance") { val input = file("k1\tsame", "k2\tsame", "k3\tother") val session = new GenericSessionFactory().create() diff --git a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala index 56a5cd308..2b3365715 100644 --- a/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala +++ b/model/src/main/scala/gorsat/gorsatGorIterator/MapAndListUtilities.scala @@ -23,7 +23,7 @@ package gorsat.gorsatGorIterator import java.nio.file.Files -import java.util.concurrent.ConcurrentHashMap +import java.util.concurrent.{CompletableFuture, ConcurrentHashMap, ExecutionException} import java.util.stream.Collectors import org.gorpipe.gor.model.{DriverBackedFileReader, FileReader} import org.gorpipe.gor.session.GorSession @@ -113,37 +113,54 @@ object MapAndListUtilities { } } - // Locks so that concurrent pipelines of a request, which share the session cache, load a given map only once. - private val loadLocks = new ConcurrentHashMap[String, Object]() + // Loads in progress, so that concurrent pipelines sharing a session cache load a given map only once. Keyed on the + // cache instance (GorSessionCache uses identity equality), not the request id, as unrelated sessions can share one. + private case class LoadKey(cache: AnyRef, name: String) + private val loadsInProgress = new ConcurrentHashMap[LoadKey, CompletableFuture[AnyRef]]() private def cacheKey(kind: String, filename: String, ic: Int, oc: Array[Int], asSet: Boolean, caseInsensitive: Boolean, skipEmpty: Boolean): String = s"$kind|$filename|$ic|${oc.mkString(",")}|$asSet|$caseInsensitive|$skipEmpty" /** - * Returns the cached value if present, otherwise loads it while holding a lock for the request and key, so - * concurrent callers wait for the first load instead of each building their own copy. + * Returns the cached value if present, otherwise loads it once for the session cache. Concurrent callers close + * their own iterator and wait (interruptibly) for the first load, sharing its result or its failure. */ - private def loadOnce[T](extFilename: String, iterator: LineIterator, session: GorSession) - (cached: => Option[T])(load: => T): T = { + private def loadOnce[T <: AnyRef](extFilename: String, iterator: LineIterator, session: GorSession) + (cached: => Option[T])(load: => T): T = { cached match { case Some(value) => iterator.close() value case None => - val lockKey = String.valueOf(session.getRequestId) + "|" + extFilename - val lock = loadLocks.computeIfAbsent(lockKey, _ => new Object) - try { - lock.synchronized { - cached match { + val key = LoadKey(session.getCache, extFilename) + val ours = new CompletableFuture[AnyRef]() + val inProgress = loadsInProgress.putIfAbsent(key, ours) + if (inProgress != null) { + iterator.close() + try { + inProgress.get().asInstanceOf[T] + } catch { + case e: ExecutionException => throw e.getCause + } + } else { + try { + // A load may have finished between the first check and registering ours. + val value = cached match { case Some(value) => iterator.close() value case None => load } + ours.complete(value) + value + } catch { + case e: Throwable => + ours.completeExceptionally(e) + throw e + } finally { + loadsInProgress.remove(key, ours) } - } finally { - loadLocks.remove(lockKey, lock) } } } @@ -197,7 +214,8 @@ object MapAndListUtilities { } else { val existingValues = colMap.get(lookupString).split("\t",-1) val newValues = if( skipEmpty ) existingValues.zip(oc.map(c => cols(c))).map(_.productIterator.filter(_.toString.nonEmpty).mkString(",")) else existingValues.zip(oc.map(c => cols(c))).map(x => x._1 + "," + x._2 ) - colMap.put(lookupString, values(newValues.tail.foldLeft(newValues.head)(_ + "\t" + _))) + // Not pooled, the merged value is replaced again by any further duplicates of the key. + colMap.put(lookupString, newValues.tail.foldLeft(newValues.head)(_ + "\t" + _)) } } }