mengw15 commented on code in PR #7539:
URL: https://github.com/apache/texera/pull/7539#discussion_r3778034439


##########
common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstanceSpec.scala:
##########
@@ -0,0 +1,220 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.texera.amber.core.storage
+
+import com.google.common.base.Ticker
+import org.apache.texera.amber.core.storage.result.iceberg.IcebergDocument
+import org.apache.texera.amber.core.tuple.{AttributeType, Schema, Tuple}
+import org.apache.texera.amber.util.IcebergUtil
+import org.apache.iceberg.Table
+import org.apache.iceberg.catalog.{Catalog, Namespace, TableIdentifier}
+import org.scalatest.flatspec.AnyFlatSpec
+import org.scalatest.matchers.should.Matchers
+
+import java.time.Duration
+
+/**
+  * Spec for the bounded catalog cache (#7290): only *idle-expired* entries 
are closed.
+  * A size-evicted catalog is dropped un-closed (it may be mid-operation under 
load),
+  * a catalog displaced by replaceInstance stays open (the replacing caller 
may still
+  * hold and later restore it -- wrap-and-restore, as amber's integration
+  * IcebergDocumentSpec does), and holders resolve their catalog per operation 
so a
+  * replacement is visible immediately.
+  *
+  * Size and idle eviction are exercised on isolated caches built through the
+  * package-private factory (with a manual ticker), never on the JVM-wide cache
+  * that parallel suites share. Tests that do touch the shared cache use their
+  * own spec-unique warehouse keys.
+  */
+class IcebergCatalogInstanceSpec extends AnyFlatSpec with Matchers {
+
+  /** A closable catalog stub; the cache only ever needs `close()` on 
eviction. */
+  private class FakeCatalog(catalogName: String) extends Catalog with 
AutoCloseable {
+    @volatile var closed = false
+    override def close(): Unit = closed = true
+    override def name(): String = catalogName
+    override def listTables(namespace: Namespace): 
java.util.List[TableIdentifier] =
+      throw new UnsupportedOperationException
+    override def dropTable(identifier: TableIdentifier, purge: Boolean): 
Boolean =
+      throw new UnsupportedOperationException
+    override def renameTable(from: TableIdentifier, to: TableIdentifier): Unit 
=
+      throw new UnsupportedOperationException
+    override def loadTable(identifier: TableIdentifier): Table =
+      throw new UnsupportedOperationException
+  }
+
+  /** A ticker the tests advance by hand, making idle expiry deterministic. */
+  private class ManualTicker extends Ticker {
+    @volatile private var nanos = 0L
+    def advance(duration: Duration): Unit = nanos += duration.toNanos
+    override def read(): Long = nanos
+  }
+
+  "the catalog cache" should "drop entries beyond the size bound without 
closing them" in {
+    // Size pressure means more simultaneously hot warehouses than the bound; 
the
+    // evicted catalog may be mid-operation, so it must decay via GC, never be 
closed.
+    val cache =
+      IcebergCatalogInstance.buildCatalogCache(2, Duration.ofMinutes(60), new 
ManualTicker)
+    val fakes = (1 to 3).map(i => new FakeCatalog(s"size-$i"))
+
+    fakes.zipWithIndex.foreach { case (fake, i) => cache.put(s"warehouse-$i", 
fake) }
+
+    cache.size() should be <= 2L
+    fakes.count(_.closed) shouldBe 0
+  }
+
+  it should "close an entry left idle beyond the expiry window" in {
+    val ticker = new ManualTicker
+    val cache = IcebergCatalogInstance.buildCatalogCache(64, 
Duration.ofMinutes(60), ticker)
+    val idle = new FakeCatalog("idle")
+    cache.put("idle", idle)
+
+    ticker.advance(Duration.ofMinutes(61))
+    // Reads alone may defer removal processing; cleanUp() drains it 
deterministically.
+    cache.cleanUp()
+
+    cache.getIfPresent("idle") shouldBe null
+    idle.closed shouldBe true
+  }
+
+  it should "surface loader failures with their original exception type" in {
+    val cache =
+      IcebergCatalogInstance.buildCatalogCache(64, Duration.ofMinutes(60), new 
ManualTicker)
+
+    // Guava wraps a runtime failure in UncheckedExecutionException and a 
checked one
+    // in ExecutionException; getOrLoad must rethrow the original in both 
cases.
+    val runtimeFailure = intercept[IllegalArgumentException] {
+      IcebergCatalogInstance.getOrLoad(
+        cache,
+        "unsupported",
+        () => throw new IllegalArgumentException("Unsupported catalog type")
+      )
+    }
+    runtimeFailure.getMessage should include("Unsupported catalog type")
+
+    an[java.io.IOException] should be thrownBy
+      IcebergCatalogInstance.getOrLoad(
+        cache,
+        "unreachable",
+        () => throw new java.io.IOException("connection refused")
+      )
+  }
+
+  "getInstance" should "return the catalog installed for its warehouse" in {
+    val installed = new FakeCatalog("installed")
+    IcebergCatalogInstance.replaceInstance(installed, 
Some("catalog-cache-spec-get"))
+
+    IcebergCatalogInstance.getInstance(Some("catalog-cache-spec-get")) should 
be theSameInstanceAs
+      installed
+  }
+
+  "replaceInstance" should "leave the displaced catalog open for its owner 
(wrap-and-restore)" in {
+    // Integration tests wrap the shared catalog in a spy and restore it 
afterwards;
+    // closing the displaced instance would hand back a dead catalog (#7290 
review).
+    val original = new FakeCatalog("original")
+    val wrapper = new FakeCatalog("wrapper")
+    IcebergCatalogInstance.replaceInstance(original, 
Some("catalog-cache-spec-replace"))
+
+    IcebergCatalogInstance.replaceInstance(wrapper, 
Some("catalog-cache-spec-replace"))
+    original.closed shouldBe false
+
+    IcebergCatalogInstance.replaceInstance(original, 
Some("catalog-cache-spec-replace"))
+    wrapper.closed shouldBe false
+    IcebergCatalogInstance.getInstance(
+      Some("catalog-cache-spec-replace")
+    ) should be theSameInstanceAs
+      original
+  }
+
+  it should "keep a re-registered shared instance open" in {
+    // LocalHadoopIcebergCatalog.ensure re-puts one shared instance from every 
suite
+    // (and under several warehouse names); none of that may close it.
+    val shared = new FakeCatalog("shared")
+    IcebergCatalogInstance.replaceInstance(shared, 
Some("catalog-cache-spec-idempotent"))
+
+    IcebergCatalogInstance.replaceInstance(shared, 
Some("catalog-cache-spec-idempotent"))
+
+    shared.closed shouldBe false
+    IcebergCatalogInstance.getInstance(Some("catalog-cache-spec-idempotent")) 
should
+      be theSameInstanceAs shared
+  }
+
+  "IcebergDocument.clear" should "address one catalog for the whole 
check-then-drop" in {
+    // Per-use resolution means per logical operation, not per call: the fake 
below
+    // swaps the cache entry from INSIDE the existence check, and the drop 
must still
+    // land on the catalog the operation started with (#7290 review, round 2).
+    class ImpostorCatalog extends FakeCatalog("impostor") {
+      @volatile var dropCalls = 0
+      override def tableExists(identifier: TableIdentifier): Boolean = true
+      override def dropTable(identifier: TableIdentifier, purge: Boolean): 
Boolean = {
+        dropCalls += 1; true
+      }
+    }
+    class SwappingCatalog(impostor: ImpostorCatalog) extends 
FakeCatalog("swapping") {
+      @volatile var dropCalls = 0
+      override def tableExists(identifier: TableIdentifier): Boolean = {
+        IcebergCatalogInstance.replaceInstance(impostor, 
Some("catalog-cache-spec-clear"))
+        true
+      }
+      override def dropTable(identifier: TableIdentifier, purge: Boolean): 
Boolean = {
+        dropCalls += 1; true
+      }
+    }
+    val impostor = new ImpostorCatalog
+    val swapping = new SwappingCatalog(impostor)
+    IcebergCatalogInstance.replaceInstance(swapping, 
Some("catalog-cache-spec-clear"))
+    val amberSchema = Schema().add("id", AttributeType.INTEGER)
+    val document = new IcebergDocument[Tuple](
+      "catalog_cache_spec",
+      "clear_probe",
+      IcebergUtil.toIcebergSchema(amberSchema),
+      IcebergUtil.toGenericRecord,
+      (schema, record) => IcebergUtil.fromRecord(record, 
IcebergUtil.fromIcebergSchema(schema)),
+      Some("catalog-cache-spec-clear")
+    )
+
+    document.clear()
+
+    swapping.dropCalls shouldBe 1
+    impostor.dropCalls shouldBe 0
+  }
+
+  "IcebergDocument" should "resolve its catalog per use, seeing a replacement 
immediately" in {
+    // Pins the per-use `def` (#7290): a `lazy val` would keep returning the 
catalog
+    // that was current at first access, i.e. a reference the cache may have 
closed.
+    val amberSchema = Schema().add("id", AttributeType.INTEGER)
+    val document = new IcebergDocument[Tuple](
+      "catalog_cache_spec",
+      "swap_probe",
+      IcebergUtil.toIcebergSchema(amberSchema),
+      IcebergUtil.toGenericRecord,
+      (schema, record) => IcebergUtil.fromRecord(record, 
IcebergUtil.fromIcebergSchema(schema)),

Review Comment:
   Done in 22afc3e6b — a local helper installs the catalog and reads it back, 
so each case is one line.



##########
amber/src/main/python/core/storage/iceberg/iceberg_catalog_instance.py:
##########
@@ -35,6 +35,11 @@ class IcebergCatalogInstance:
       testing or reconfiguration.
     """
 
+    # Deliberately unbounded, unlike the Scala side's bounded cache (#7290): a 
Python
+    # worker process is spawned per worker and destroyed when the execution 
ends
+    # (PythonWorkflowWorker.postStop), so this dict holds the single warehouse 
that
+    # execution used and dies with the process. If PVMs ever become long-lived 
(e.g.
+    # pooled across executions), this needs the same bounding the JVM 
singleton has.

Review Comment:
   It documents a deliberate asymmetry this PR creates: the Scala registry 
becomes bounded while the Python one stays a plain dict, and the natural 
reading is that Python was missed. It wasn't — a PVM is spawned per worker and 
destroyed when the execution ends, so its dict holds one warehouse and dies 
with the process. The last sentence is the load-bearing one: it names the 
assumption ("PVMs are short-lived") that makes "unbounded" safe, right above 
the code someone would change if PVMs were ever pooled across executions.



##########
amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala:
##########
@@ -114,32 +113,33 @@ class IcebergDocumentSpec extends 
VirtualDocumentSpec[Tuple] with BeforeAndAfter
     val (batch1, batch2) = items.splitAt(batchSize)
 
     // Write two separate batches to produce two committed data files.
-    // This also initialises `document.catalog` (lazy val) with the real 
catalog, which
-    // is why we open a fresh reader document below after injecting the spy.
     val writer1 = document.writer(UUID.randomUUID().toString)
     writer1.open(); batch1.foreach(writer1.putOne); writer1.close()
 
     val writer2 = document.writer(UUID.randomUUID().toString)
     writer2.open(); batch2.foreach(writer2.putOne); writer2.close()
 
-    val refreshCount = new AtomicInteger(0)
+    val metadataLoadCount = new AtomicInteger(0)
     val realCatalog = IcebergCatalogInstance.getInstance()
-    IcebergCatalogInstance.replaceInstance(catalogWithRefreshSpy(realCatalog, 
refreshCount))
-    // Open a fresh reader: its `catalog` lazy val hasn't been initialised 
yet, so it
-    // will pick up the spy catalog on first access inside seekToUsableFile.
+    IcebergCatalogInstance.replaceInstance(
+      catalogWithMetadataLoadSpy(realCatalog, metadataLoadCount)
+    )
+    // Open a fresh reader; it resolves its catalog per use (#7290), so every 
metadata
+    // load inside seekToUsableFile goes through the spy installed above.

Review Comment:
   The wrapping was scalafmt's, not a style choice — the renamed helper pushed 
the line past `maxColumn = 100`. Shortened the names in 22afc3e6b so it fits on 
one line again and the diff stays minimal. The rename itself is load-bearing 
though: the spy counted `table.refresh()`, and after this PR the reader never 
calls `refresh()` — it re-resolves the table instead — so the assertion would 
have been vacuously true (0 <= 4) and could no longer catch a lazy-advancement 
regression. It now counts `loadTable`, where those metadata loads surface.



##########
amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala:
##########
@@ -114,32 +113,33 @@ class IcebergDocumentSpec extends 
VirtualDocumentSpec[Tuple] with BeforeAndAfter
     val (batch1, batch2) = items.splitAt(batchSize)
 
     // Write two separate batches to produce two committed data files.
-    // This also initialises `document.catalog` (lazy val) with the real 
catalog, which
-    // is why we open a fresh reader document below after injecting the spy.
     val writer1 = document.writer(UUID.randomUUID().toString)
     writer1.open(); batch1.foreach(writer1.putOne); writer1.close()
 
     val writer2 = document.writer(UUID.randomUUID().toString)
     writer2.open(); batch2.foreach(writer2.putOne); writer2.close()
 
-    val refreshCount = new AtomicInteger(0)
+    val metadataLoadCount = new AtomicInteger(0)
     val realCatalog = IcebergCatalogInstance.getInstance()
-    IcebergCatalogInstance.replaceInstance(catalogWithRefreshSpy(realCatalog, 
refreshCount))
-    // Open a fresh reader: its `catalog` lazy val hasn't been initialised 
yet, so it
-    // will pick up the spy catalog on first access inside seekToUsableFile.
+    IcebergCatalogInstance.replaceInstance(
+      catalogWithMetadataLoadSpy(realCatalog, metadataLoadCount)
+    )
+    // Open a fresh reader; it resolves its catalog per use (#7290), so every 
metadata
+    // load inside seekToUsableFile goes through the spy installed above.
     val readerDoc = getDocument
     try {
       val retrieved = readerDoc.get().toList
       assert(
         retrieved.toSet == items.toSet,
         "All records from both files should be read correctly"
       )
-      // With lazy file advancement seekToUsableFile() (and therefore 
table.refresh()) is called:
-      //   once on iterator creation, once when the last file is exhausted → 2 
total.
-      // Without the fix it would be called once per hasNext() on the last 
file → O(batchSize).
+      // With lazy file advancement the table is (re-)resolved once at iterator
+      // construction plus once per seekToUsableFile — construction seek and 
the
+      // final exhausted-files seek → 3 total. Without lazy advancement it 
would be
+      // once per hasNext() on the last file → O(batchSize).
       assert(
-        refreshCount.get() <= 4,
-        s"table.refresh() should be called at most 4 times (lazy advancement), 
but was ${refreshCount.get()}"
+        metadataLoadCount.get() <= 4,
+        s"the table should be loaded at most 4 times (lazy advancement), but 
was ${metadataLoadCount.get()}"

Review Comment:
   Same rename as above (22afc3e6b). The comment did need correcting on its own 
account: it said 3 total, which went stale when the iterator's eager 
constructor-time load was dropped — it is 2 now (construction seek + final 
exhausted-files seek).



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to