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]
