This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new a09260a46c feat(storage): bound the per-warehouse catalog cache and
release evicted catalogs (#7539)
a09260a46c is described below
commit a09260a46cbfa844537b8865ff0f79754bf3daa5
Author: Meng Wang <[email protected]>
AuthorDate: Thu Aug 13 21:07:17 2026 +0000
feat(storage): bound the per-warehouse catalog cache and release evicted
catalogs (#7539)
### What changes were proposed in this PR?
**TL;DR**: Texera keeps one Iceberg catalog client (an HTTP client + its
connection pool) per warehouse, in a process-wide registry that never
removes entries. That used to be harmless: there was effectively one
shared warehouse, so the registry held one entry forever. With per-user
warehouses (#6870) users create warehouses freely, and every warehouse a
long-lived JVM (web server, computing unit) ever touches adds one more
permanently-held client — the registry only grows for the life of the
process, accumulating connection pools that are never released. This PR
(1) bounds that registry and closes clients that have gone idle, and (2)
reworks every reader/writer that used to pin a client reference
long-term to re-resolve it per operation — which is what makes releasing
clients safe. Everything below is the detail of those two moves.
`IcebergCatalogInstance` kept one catalog client per warehouse name for
the life of the process; with per-user warehouses (#6870) that set is
unbounded, and each REST catalog holds an HTTP client. The map is now a
Guava cache (`maximumSize` 64 + `expireAfterAccess` 60 min, mirroring
`HuggingFaceModelResource`'s bounded-cache precedent). An entry idle for
the expiry window is closed — nothing can be using it, and idle entries
are exactly what a long-lived JVM accumulates. An entry evicted by size
is only dropped, never closed: size pressure means more simultaneously
hot warehouses than the bound, and closing a hot catalog would fail the
operations still using it. Load degrades into rebuild churn, not errors
— a dropped catalog lives only as long as its in-flight operations
(per-operation resolution bounds every borrow), after which GC reclaims
it while the server's keepalive timeout severs its idle connections.
For eviction to be safe, holders stop pinning a catalog — or anything
derived from one — across a logical operation: `IcebergDocument`'s `lazy
val` becomes a per-use `def`, its `clear()` resolves one catalog for the
whole check-then-drop, the reader re-resolves its table on every seek
instead of refreshing a pinned one (which also keeps a polling reader's
cache entry live), and `IcebergTableWriter` takes the warehouse rather
than a `Catalog` and loads its table per flush. The lookup rides Guava's
per-key locking, dropping the previous JVM-wide `synchronized` that held
one lock across a cache miss's REST config round trip, and unwraps
Guava's `ExecutionException` family so `createCatalog` failures keep the
types they had before.
Only idle-expired entries are closed. `replaceInstance` stays a plain
`put`: a caller that replaces an entry may still hold and later restore
the old reference — amber's integration spec wrap-and-restores the
shared catalog, and endpoint reconfiguration (#7358) will swap catalogs
the same way.
Also dedupes `DocumentFactory`'s three copies of the URI→(warehouse,
namespace, storage key) decode block into one resolver, as promised in
#6944 review. The Python side is untouched: a PVM is spawned per worker
and destroyed when the execution ends, so its catalog dict holds the
single warehouse that execution used and dies with the process — nothing
accumulates there to bound.
### Any related issues, documentation, discussions?
Closes #7290.
### How was this PR tested?
New `IcebergCatalogInstanceSpec` covers the cache contract: idle expiry
closes the catalog while size eviction drops it un-closed — on isolated
caches built through a package-private factory with a manual ticker, so
the JVM-wide cache that parallel suites share is never touched — a
catalog displaced by `replaceInstance` stays open for its owner, and
loader failures keep their original exception type.
Per-use resolution is pinned at every holder: `IcebergDocument` sees a
replacement immediately, `clear()` addresses one catalog across its
check-then-drop, the reader re-resolves per seek, and the writer loads
through the catalog installed at flush time. Existing iceberg suites
(`IcebergDocumentSpec`, `IcebergTableWriterSpec`, `OnIcebergSpec`,
`DocumentFactorySpec`) pass locally — 769 tests in `workflow-core` — and
amber's integration `IcebergDocumentSpec` is green in CI.
### Was this PR authored or co-authored using generative AI tooling?
Generated-by: Claude Code (claude-fable-5)
---
.../result/iceberg/IcebergDocumentSpec.scala | 50 ++--
.../amber/core/storage/DocumentFactory.scala | 58 ++---
.../core/storage/IcebergCatalogInstance.scala | 107 ++++++++-
.../storage/result/iceberg/IcebergDocument.scala | 30 ++-
.../result/iceberg/IcebergTableWriter.scala | 22 +-
.../core/storage/IcebergCatalogInstanceSpec.scala | 262 +++++++++++++++++++++
.../storage/result/iceberg/CountingCatalog.scala | 48 ++++
.../result/iceberg/IcebergDocumentSpec.scala | 33 +++
.../result/iceberg/IcebergTableWriterSpec.scala | 49 +++-
.../storage/result/iceberg/OnIcebergSpec.scala | 8 +-
10 files changed, 574 insertions(+), 93 deletions(-)
diff --git
a/amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala
b/amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala
index aa8011e553..fa5e84e4c9 100644
---
a/amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala
+++
b/amber/src/test/integration/org/apache/texera/amber/storage/result/iceberg/IcebergDocumentSpec.scala
@@ -24,7 +24,6 @@ import org.apache.texera.amber.core.state.State
import org.apache.texera.amber.core.storage.model.{VirtualDocument,
VirtualDocumentSpec}
import org.apache.texera.amber.core.storage.{DocumentFactory,
IcebergCatalogInstance, VFSURIFactory}
import org.apache.texera.amber.core.tuple.{Attribute, AttributeType, Schema,
Tuple}
-import org.apache.iceberg.Table
import org.apache.texera.amber.core.virtualidentity.{
ExecutionIdentity,
OperatorIdentity,
@@ -114,19 +113,17 @@ 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 loadCount = 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(catalogWithLoadSpy(realCatalog,
loadCount))
+ // 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
@@ -134,12 +131,12 @@ class IcebergDocumentSpec extends
VirtualDocumentSpec[Tuple] with BeforeAndAfter
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 resolved once per
seekToUsableFile —
+ // the construction seek and the final exhausted-files seek → 2 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()}"
+ loadCount.get() <= 4,
+ s"the table should be loaded at most 4 times (lazy advancement), but
was ${loadCount.get()}"
)
} finally {
IcebergCatalogInstance.replaceInstance(realCatalog)
@@ -309,35 +306,18 @@ class IcebergDocumentSpec extends
VirtualDocumentSpec[Tuple] with BeforeAndAfter
}
}
- /** Returns a dynamic proxy for `realTable` that increments `counter` on
every `refresh()` call. */
- private def tableWithRefreshSpy(realTable: Table, counter: AtomicInteger):
Table =
- Proxy
- .newProxyInstance(
- classOf[Table].getClassLoader,
- Array(classOf[Table]),
- new InvocationHandler {
- override def invoke(proxy: Object, method: Method, args:
Array[Object]): Object = {
- if (method.getName == "refresh") counter.incrementAndGet()
- if (args == null) method.invoke(realTable) else
method.invoke(realTable, args: _*)
- }
- }
- )
- .asInstanceOf[Table]
-
- /** Returns a dynamic proxy for `realCatalog` that wraps every loaded
`Table` with a refresh spy. */
- private def catalogWithRefreshSpy(realCatalog: Catalog, counter:
AtomicInteger): Catalog =
+ /** Returns a dynamic proxy for `realCatalog` that counts `loadTable` calls.
*/
+ private def catalogWithLoadSpy(realCatalog: Catalog, counter:
AtomicInteger): Catalog =
Proxy
.newProxyInstance(
classOf[Catalog].getClassLoader,
Array(classOf[Catalog]),
new InvocationHandler {
override def invoke(proxy: Object, method: Method, args:
Array[Object]): Object = {
- val result =
- if (args == null) method.invoke(realCatalog) else
method.invoke(realCatalog, args: _*)
- if (method.getName == "loadTable" && result != null)
- tableWithRefreshSpy(result.asInstanceOf[Table], counter)
- else
- result
+ // The reader re-resolves its table per seek (#7290) instead of
refreshing a
+ // pinned one, so metadata loads now surface as `loadTable` calls
here.
+ if (method.getName == "loadTable") counter.incrementAndGet()
+ if (args == null) method.invoke(realCatalog) else
method.invoke(realCatalog, args: _*)
}
}
)
diff --git
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/DocumentFactory.scala
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/DocumentFactory.scala
index f84158e0e3..16fa5c07cb 100644
---
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/DocumentFactory.scala
+++
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/DocumentFactory.scala
@@ -66,6 +66,30 @@ object DocumentFactory {
}
}
+ /**
+ * The iceberg coordinates a VFS URI resolves to: which warehouse's
catalog, which
+ * namespace, and which table (storage key). One resolver shared by every
VFS entry
+ * point below, so the decode steps cannot drift apart (promised in #6944
review).
+ */
+ private case class IcebergLocation(
+ warehouse: Option[String],
+ namespace: String,
+ storageKey: String
+ )
+
+ private def resolveIcebergLocation(uri: URI): IcebergLocation = {
+ val components = decodeURI(uri)
+ IcebergLocation(
+ components.warehouse,
+ resolveNamespace(components.resourceType),
+ sanitizeURIPath(uri)
+ )
+ }
+
+ private val tupleSerde: (IcebergSchema, Tuple) => Record =
IcebergUtil.toGenericRecord
+ private val tupleDeserde: (IcebergSchema, Record) => Tuple = (schema,
record) =>
+ IcebergUtil.fromRecord(record, IcebergUtil.fromIcebergSchema(schema))
+
/**
* Create a document for storage specified by the uri.
* This document is suitable for storing structural data, i.e. the schema
is required to create such document.
@@ -76,11 +100,7 @@ object DocumentFactory {
def createDocument(uri: URI, schema: Schema): VirtualDocument[_] = {
uri.getScheme match {
case VFS_FILE_URI_SCHEME =>
- val components = decodeURI(uri)
- val warehouse = components.warehouse
- val resourceType = components.resourceType
- val storageKey = sanitizeURIPath(uri)
- val namespace = resolveNamespace(resourceType)
+ val IcebergLocation(warehouse, namespace, storageKey) =
resolveIcebergLocation(uri)
val icebergSchema = IcebergUtil.toIcebergSchema(schema)
IcebergUtil.createTable(
@@ -90,16 +110,12 @@ object DocumentFactory {
icebergSchema,
overrideIfExists = true
)
- val serde: (IcebergSchema, Tuple) => Record =
IcebergUtil.toGenericRecord
- val deserde: (IcebergSchema, Record) => Tuple = (schema, record) =>
- IcebergUtil.fromRecord(record, IcebergUtil.fromIcebergSchema(schema))
-
new IcebergDocument[Tuple](
namespace,
storageKey,
icebergSchema,
- serde,
- deserde,
+ tupleSerde,
+ tupleDeserde,
warehouse
)
case unsupportedScheme =>
@@ -122,11 +138,7 @@ object DocumentFactory {
def documentExists(uri: URI): Boolean = {
uri.getScheme match {
case VFS_FILE_URI_SCHEME =>
- val components = decodeURI(uri)
- val warehouse = components.warehouse
- val resourceType = components.resourceType
- val storageKey = sanitizeURIPath(uri)
- val namespace = resolveNamespace(resourceType)
+ val IcebergLocation(warehouse, namespace, storageKey) =
resolveIcebergLocation(uri)
IcebergCatalogInstance
.getInstance(warehouse)
.tableExists(TableIdentifier.of(namespace, storageKey))
@@ -170,11 +182,7 @@ object DocumentFactory {
uri.getScheme match {
case DATASET_FILE_URI_SCHEME => (new DatasetFileDocument(uri), None)
case VFS_FILE_URI_SCHEME =>
- val components = decodeURI(uri)
- val warehouse = components.warehouse
- val resourceType = components.resourceType
- val storageKey = sanitizeURIPath(uri)
- val namespace = resolveNamespace(resourceType)
+ val IcebergLocation(warehouse, namespace, storageKey) =
resolveIcebergLocation(uri)
val table = IcebergUtil
.loadTableMetadata(
@@ -187,17 +195,13 @@ object DocumentFactory {
)
val amberSchema = IcebergUtil.fromIcebergSchema(table.schema())
- val serde: (IcebergSchema, Tuple) => Record =
IcebergUtil.toGenericRecord
- val deserde: (IcebergSchema, Record) => Tuple = (schema, record) =>
- IcebergUtil.fromRecord(record, IcebergUtil.fromIcebergSchema(schema))
-
(
new IcebergDocument[Tuple](
namespace,
storageKey,
table.schema(),
- serde,
- deserde,
+ tupleSerde,
+ tupleDeserde,
warehouse
),
Some(amberSchema)
diff --git
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstance.scala
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstance.scala
index 313772b0fc..3e2eabceea 100644
---
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstance.scala
+++
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstance.scala
@@ -19,11 +19,23 @@
package org.apache.texera.amber.core.storage
+import com.google.common.base.Ticker
+import com.google.common.cache.{
+ Cache,
+ CacheBuilder,
+ RemovalCause,
+ RemovalListener,
+ RemovalNotification
+}
+import com.google.common.util.concurrent.{ExecutionError,
UncheckedExecutionException}
+import com.typesafe.scalalogging.LazyLogging
import org.apache.texera.common.config.StorageConfig
import org.apache.texera.amber.util.IcebergUtil
import org.apache.iceberg.catalog.Catalog
-import scala.collection.mutable
+import java.time.Duration
+import java.util.concurrent.{Callable, ExecutionException}
+import scala.util.Try
/**
* IcebergCatalogInstance manages the Iceberg catalog clients used across the
Texera application.
@@ -36,11 +48,65 @@ import scala.collection.mutable
* Only the REST catalog varies by warehouse; the hadoop and postgres
catalogs are warehouse-agnostic
* and ignore the warehouse argument.
*
- * Access is synchronized because the same JVM serves multiple warehouses
concurrently.
+ * The cache is bounded (#7290): per-user warehouses (#6870) make the set of
catalogs a
+ * long-lived JVM touches unbounded, and each REST catalog holds an HTTP
client. An entry
+ * idle for the expiry window is closed -- nothing can be using it, and idle
entries are
+ * exactly what a long-lived JVM accumulates. An entry evicted by *size* is
only dropped,
+ * never closed: size pressure means more simultaneously hot warehouses than
the bound,
+ * and closing a hot catalog would fail the operations still using it. Load
degrades into
+ * rebuild churn (the dropped catalog decays once its in-flight users
finish), not errors.
+ *
+ * Callers must therefore resolve their catalog per logical operation instead
of holding
+ * one across an execution -- that is also what keeps a dropped catalog's
lifetime bounded
+ * by the operation using it (see IcebergDocument / IcebergTableWriter).
+ *
+ * Only *evicted* entries are closed. A catalog displaced by
[[replaceInstance]] is the
+ * caller's to manage: whoever replaces an entry may still hold (and restore)
the old
+ * reference -- tests wrap-and-restore the shared catalog, and endpoint
reconfiguration
+ * (#7358) will swap catalogs the same way.
*/
-object IcebergCatalogInstance {
+object IcebergCatalogInstance extends LazyLogging {
+
+ // Sizing mirrors HuggingFaceModelResource's bounded-cache precedent:
generous enough
+ // that eviction never hits a warehouse in active use, small enough to bound
the JVM.
+ private val CatalogCacheMaxSize = 64L
+ private val CatalogCacheExpireAfterAccess = Duration.ofMinutes(60)
+
+ /**
+ * Builds a catalog cache with the eviction wiring `getInstance` relies on.
+ * Package-private so the spec can exercise size and idle eviction on
isolated
+ * instances with a manual ticker, instead of flooding the JVM-wide cache
below.
+ */
+ private[storage] def buildCatalogCache(
+ maximumSize: Long,
+ expireAfterAccess: Duration,
+ ticker: Ticker
+ ): Cache[String, Catalog] =
+ CacheBuilder
+ .newBuilder()
+ .maximumSize(maximumSize)
+ .expireAfterAccess(expireAfterAccess)
+ .ticker(ticker)
+ .removalListener(new RemovalListener[String, Catalog] {
+ override def onRemoval(notification: RemovalNotification[String,
Catalog]): Unit =
+ // Close ONLY idle-expired entries. A size-evicted catalog may be
mid-operation
+ // (overload = more hot warehouses than the bound) and a replaced
one is still
+ // the replacing caller's (wrap-and-restore in tests,
reconfiguration later);
+ // both are dropped un-closed and decay once their last user
finishes.
+ if (notification.getCause == RemovalCause.EXPIRED) {
+ notification.getValue match {
+ case closeable: AutoCloseable =>
+ Try(closeable.close()).failed.foreach(error =>
+ logger.warn(s"failed to close expired catalog
'${notification.getKey}'", error)
+ )
+ case _ =>
+ }
+ }
+ })
+ .build[String, Catalog]()
- private val catalogs = mutable.Map.empty[String, Catalog]
+ private val catalogs: Cache[String, Catalog] =
+ buildCatalogCache(CatalogCacheMaxSize, CatalogCacheExpireAfterAccess,
Ticker.systemTicker())
// Cache key for the warehouse-agnostic catalog types. Not a legal warehouse
name,
// so it cannot collide with a REST warehouse.
@@ -70,11 +136,34 @@ object IcebergCatalogInstance {
*/
def getInstance(warehouse: Option[String] = None): Catalog = {
val name = warehouse.getOrElse(defaultWarehouse)
- synchronized {
- catalogs.getOrElseUpdate(cacheKey(name), createCatalog(name))
- }
+ getOrLoad(catalogs, cacheKey(name), () => createCatalog(name))
}
+ /**
+ * `Cache.get` wraps loader failures (`UncheckedExecutionException`,
`ExecutionException`,
+ * `ExecutionError`); unwrap them so `createCatalog` failures keep the
types they had
+ * before the cache existed. Package-private so the spec can pin the
unwrapping against
+ * an isolated cache with a throwing loader.
+ */
+ private[storage] def getOrLoad(
+ cache: Cache[String, Catalog],
+ key: String,
+ loader: () => Catalog
+ ): Catalog =
+ try {
+ // get(key, loader) locks per key, not globally: a cache miss's REST
config
+ // round trip no longer blocks lookups of other warehouses.
+ cache.get(
+ key,
+ new Callable[Catalog] {
+ override def call(): Catalog = loader()
+ }
+ )
+ } catch {
+ case e @ (_: UncheckedExecutionException | _: ExecutionException | _:
ExecutionError) =>
+ throw e.getCause
+ }
+
private def createCatalog(warehouse: String): Catalog =
StorageConfig.icebergCatalogType match {
case "hadoop" =>
@@ -103,7 +192,5 @@ object IcebergCatalogInstance {
* @param warehouse the warehouse to cache it under; `None` uses the
configured default.
*/
def replaceInstance(catalog: Catalog, warehouse: Option[String] = None):
Unit =
- synchronized {
- catalogs(cacheKey(warehouse.getOrElse(defaultWarehouse))) = catalog
- }
+ catalogs.put(cacheKey(warehouse.getOrElse(defaultWarehouse)), catalog)
}
diff --git
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
index cc414825d9..3f3f131ada 100644
---
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
+++
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocument.scala
@@ -72,7 +72,10 @@ private[storage] class IcebergDocument[T >: Null <: AnyRef](
private val lock = new ReentrantReadWriteLock()
- @transient lazy val catalog: Catalog =
IcebergCatalogInstance.getInstance(warehouse)
+ // Resolved per use, never held: the catalog cache is bounded and closes
evicted
+ // entries (#7290), so a pinned reference could outlive its catalog. A
public def
+ // (not a lazy val) also means a replaced/rebuilt catalog is picked up
immediately.
+ def catalog: Catalog = IcebergCatalogInstance.getInstance(warehouse)
/**
* Returns the URI of the table location.
@@ -94,8 +97,11 @@ private[storage] class IcebergDocument[T >: Null <: AnyRef](
override def clear(): Unit =
withWriteLock(lock) {
val identifier = TableIdentifier.of(tableNamespace, tableName)
- if (catalog.tableExists(identifier)) {
- catalog.dropTable(identifier)
+ // One resolve for the whole check-then-drop: both steps must address
the same
+ // catalog even if the cache entry is replaced between them (#7290).
+ val currentCatalog = catalog
+ if (currentCatalog.tableExists(identifier)) {
+ currentCatalog.dropTable(identifier)
}
}
@@ -141,7 +147,7 @@ private[storage] class IcebergDocument[T >: Null <: AnyRef](
override def writer(writerIdentifier: String): BufferedItemWriter[T] = {
new IcebergTableWriter[T](
writerIdentifier,
- catalog,
+ warehouse,
tableNamespace,
tableName,
tableSchema,
@@ -164,8 +170,9 @@ private[storage] class IcebergDocument[T >: Null <: AnyRef](
withReadLock(lock) {
new Iterator[T] {
private val iteLock = new ReentrantLock()
- // Load the table instance, initially the table instance may not exist
- private var table: Option[Table] = loadTableMetadata()
+ // No eager load: the constructor-time seekToUsableFile() below
resolves the
+ // table, so loading here would be an immediately-overwritten REST
round trip.
+ private var table: Option[Table] = None
// Last seen snapshot id(logically it's like a version number). While
reading, new snapshots may be created
private var lastSnapshotId: Option[Long] = None
@@ -203,11 +210,12 @@ private[storage] class IcebergDocument[T >: Null <:
AnyRef](
throw new RuntimeException("seek operation should not be called")
}
- // refresh the table's snapshots
- if (table.isEmpty) {
- table = loadTableMetadata()
- }
- table.foreach(_.refresh())
+ // Re-resolve the table from the current catalog instead of
refreshing a
+ // pinned one (#7290): a Table held across polls keeps its REST
operations
+ // bound to a catalog the bounded cache may have closed, and
re-resolving
+ // also keeps this warehouse's cache entry live for as long as the
reader
+ // polls. Snapshot continuity lives in lastSnapshotId, not in the
Table.
+ table = loadTableMetadata()
// Retrieve and sort the file scan tasks by file sequence number.
// Materialize inside `Using.resource` so the `planFiles()`
diff --git
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriter.scala
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriter.scala
index 81b27d1139..b0f34f0b56 100644
---
a/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriter.scala
+++
b/common/workflow-core/src/main/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriter.scala
@@ -20,6 +20,7 @@
package org.apache.texera.amber.core.storage.result.iceberg
import org.apache.texera.common.config.StorageConfig
+import org.apache.texera.amber.core.storage.IcebergCatalogInstance
import org.apache.texera.amber.core.storage.model.BufferedItemWriter
import org.apache.texera.amber.util.IcebergUtil
import org.apache.iceberg.catalog.Catalog
@@ -41,7 +42,8 @@ import scala.collection.mutable.ArrayBuffer
* **Thread Safety**: This writer is **NOT thread-safe**, so only one thread
should call this writer.
*
* @param writerIdentifier a unique identifier used to prefix the created
files.
- * @param catalog the Iceberg catalog to manage table metadata.
+ * @param warehouse the warehouse whose catalog manages the table metadata;
`None` uses the
+ * configured default.
* @param tableNamespace the namespace of the Iceberg table.
* @param tableName the name of the Iceberg table.
* @param tableSchema the schema of the Iceberg table.
@@ -50,13 +52,17 @@ import scala.collection.mutable.ArrayBuffer
*/
private[storage] class IcebergTableWriter[T](
val writerIdentifier: String,
- val catalog: Catalog,
+ val warehouse: Option[String],
val tableNamespace: String,
val tableName: String,
val tableSchema: Schema,
val serde: (org.apache.iceberg.Schema, T) => Record
) extends BufferedItemWriter[T] {
+ // Resolved per use (#7290): the catalog cache is bounded and closes evicted
entries,
+ // so the writer must not pin one across its lifetime.
+ private def catalog: Catalog = IcebergCatalogInstance.getInstance(warehouse)
+
// Buffer to hold items before flushing to the table
private val buffer = new ArrayBuffer[T]()
// Incremental filename index, incremented each time a new buffer is flushed
@@ -66,12 +72,6 @@ private[storage] class IcebergTableWriter[T](
override val bufferSize: Int = StorageConfig.icebergTableCommitBatchSize
- // Load the Iceberg table
- private val table: Table =
- IcebergUtil
- .loadTableMetadata(catalog, tableNamespace, tableName)
- .get
-
/**
* Open the writer and clear the buffer.
*/
@@ -106,6 +106,12 @@ private[storage] class IcebergTableWriter[T](
*/
private def flushBuffer(): Unit = {
if (buffer.nonEmpty) {
+ // Resolve the table per flush (#7290): an eagerly-held Table would pin
REST
+ // operations backed by a catalog the bounded cache may close, and
resolving
+ // here also keeps this warehouse's cache entry live for the whole
execution.
+ val table: Table = IcebergUtil
+ .loadTableMetadata(catalog, tableNamespace, tableName)
+ .get
// Create a unique file path using the writer's identifier and the
filename index
val location = table.location().stripSuffix("/")
val filepathString = s"$location/${writerIdentifier}_$filenameIdx"
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstanceSpec.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstanceSpec.scala
new file mode 100644
index 0000000000..e14cac58f2
--- /dev/null
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/IcebergCatalogInstanceSpec.scala
@@ -0,0 +1,262 @@
+/*
+ * 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 "tolerate a catalog whose close fails, and still drop the entry"
in {
+ val ticker = new ManualTicker
+ val cache = IcebergCatalogInstance.buildCatalogCache(64,
Duration.ofMinutes(60), ticker)
+ val faulty = new FakeCatalog("faulty") {
+ override def close(): Unit = throw new IllegalStateException("close
failed")
+ }
+ cache.put("faulty", faulty)
+
+ ticker.advance(Duration.ofMinutes(61))
+ noException should be thrownBy cache.cleanUp()
+
+ cache.getIfPresent("faulty") shouldBe null
+ }
+
+ it should "leave a catalog that is not closable alone when it expires" in {
+ // Hadoop/postgres catalogs need not implement AutoCloseable; expiry must
not fail.
+ val ticker = new ManualTicker
+ val cache = IcebergCatalogInstance.buildCatalogCache(64,
Duration.ofMinutes(60), ticker)
+ val notClosable = new Catalog {
+ override def name(): String = "not-closable"
+ 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
+ }
+ cache.put("not-closable", notClosable)
+
+ ticker.advance(Duration.ofMinutes(61))
+ noException should be thrownBy cache.cleanUp()
+
+ cache.getIfPresent("not-closable") shouldBe null
+ }
+
+ 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")
+ )
+
+ // Errors ride the third wrapper, ExecutionError.
+ an[StackOverflowError] should be thrownBy
+ IcebergCatalogInstance.getOrLoad(cache, "fatal", () => throw new
StackOverflowError("boom"))
+ }
+
+ "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)),
+ Some("catalog-cache-spec-swap")
+ )
+ def catalogSeenAfterInstalling(catalog: Catalog): Catalog = {
+ IcebergCatalogInstance.replaceInstance(catalog,
Some("catalog-cache-spec-swap"))
+ document.catalog
+ }
+ val before = new FakeCatalog("before")
+ val after = new FakeCatalog("after")
+
+ catalogSeenAfterInstalling(before) should be theSameInstanceAs before
+ catalogSeenAfterInstalling(after) should be theSameInstanceAs after
+ }
+}
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/CountingCatalog.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/CountingCatalog.scala
new file mode 100644
index 0000000000..e62045cfc8
--- /dev/null
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/CountingCatalog.scala
@@ -0,0 +1,48 @@
+/*
+ * 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.result.iceberg
+
+import org.apache.iceberg.Table
+import org.apache.iceberg.catalog.{Catalog, Namespace, TableIdentifier}
+
+import java.util.concurrent.atomic.AtomicInteger
+
+/**
+ * Test helper: delegates to `delegate` while counting `loadTable` calls --
the
+ * discriminator for per-operation table resolution (#7290). A holder that
pins a
+ * `Table` (or refreshes a pinned one) touches the catalog once, at
construction;
+ * per-operation resolution touches it again on every flush/seek.
+ */
+class CountingCatalog(delegate: Catalog) extends Catalog {
+ val loadTableCalls = new AtomicInteger()
+ override def name(): String = "counting"
+ override def loadTable(identifier: TableIdentifier): Table = {
+ loadTableCalls.incrementAndGet()
+ delegate.loadTable(identifier)
+ }
+ override def tableExists(identifier: TableIdentifier): Boolean =
+ delegate.tableExists(identifier)
+ override def listTables(namespace: Namespace):
java.util.List[TableIdentifier] =
+ delegate.listTables(namespace)
+ override def dropTable(identifier: TableIdentifier, purge: Boolean): Boolean
=
+ delegate.dropTable(identifier, purge)
+ override def renameTable(from: TableIdentifier, to: TableIdentifier): Unit =
+ delegate.renameTable(from, to)
+}
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
index a1f4ea3023..fb2f5b4a15 100644
---
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergDocumentSpec.scala
@@ -328,4 +328,37 @@ class IcebergDocumentSpec extends AnyFlatSpec with
Matchers with BeforeAndAfterA
doc.asInstanceOf[IcebergDocument[Tuple]].tableNamespace shouldBe
tableNamespace
doc.asInstanceOf[IcebergDocument[Tuple]].tableName shouldBe name
}
+ it should "re-resolve its catalog on every seek, keeping a polling reader's
entry live" in {
+ // #7290 review: a reader that refreshed a pinned Table never touched the
catalog
+ // cache again, so a long poll looked idle -- expiry could close the
catalog under
+ // it. Per-seek re-resolution touches the cache at iterator construction
AND every
+ // seek; the pinned-refresh implementation stopped at the construction
touch.
+ val tableName = freshTableName()
+ IcebergUtil.createTable(
+ IcebergCatalogInstance.getInstance(),
+ tableNamespace,
+ tableName,
+ icebergSchema,
+ overrideIfExists = true
+ )
+ val counting = new CountingCatalog(IcebergCatalogInstance.getInstance())
+ IcebergCatalogInstance.replaceInstance(counting,
Some("iceberg-doc-spec-counting"))
+ val doc = new IcebergDocument[Tuple](
+ tableNamespace,
+ tableName,
+ icebergSchema,
+ serde,
+ deserde,
+ Some("iceberg-doc-spec-counting")
+ )
+ write(doc, (1 to 3).map(tuple))
+
+ val before = counting.loadTableCalls.get()
+ doc.get().toList should have size 3
+
+ // Exactly the construction-time seek and the final exhausted-files seek
resolve
+ // through the catalog: the pinned-refresh implementation touched it only
once,
+ // and an eager constructor-time load would add a wasted third round trip.
+ (counting.loadTableCalls.get() - before) shouldBe 2
+ }
}
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriterSpec.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriterSpec.scala
index 5c82392417..8b8013f6e0 100644
---
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriterSpec.scala
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/IcebergTableWriterSpec.scala
@@ -20,6 +20,7 @@
package org.apache.texera.amber.core.storage.result.iceberg
import org.apache.texera.amber.core.tuple.{AttributeType, Schema, Tuple}
+import org.apache.texera.amber.core.storage.IcebergCatalogInstance
import org.apache.texera.amber.util.IcebergUtil
import org.apache.iceberg.catalog.Catalog
import org.apache.iceberg.data.IcebergGenerics
@@ -43,9 +44,14 @@ class IcebergTableWriterSpec extends AnyFlatSpec with
BeforeAndAfterAll {
private val icebergSchema: IcebergSchema =
IcebergUtil.toIcebergSchema(amberSchema)
+ // The writer resolves its catalog from the shared cache per use (#7290), so
the
+ // spec's local catalog is registered under a spec-unique warehouse name.
+ private val specWarehouse = "iceberg-table-writer-spec"
+
override def beforeAll(): Unit = {
warehouseDir = Files.createTempDirectory("iceberg-table-writer-spec")
catalog = IcebergUtil.createHadoopCatalog("writer-spec", warehouseDir)
+ IcebergCatalogInstance.replaceInstance(catalog, Some(specWarehouse))
}
override def afterAll(): Unit = {
@@ -71,7 +77,7 @@ class IcebergTableWriterSpec extends AnyFlatSpec with
BeforeAndAfterAll {
)
new IcebergTableWriter[Tuple](
writerIdentifier,
- catalog,
+ Some(specWarehouse),
tableNamespace,
tableName,
icebergSchema,
@@ -157,4 +163,45 @@ class IcebergTableWriterSpec extends AnyFlatSpec with
BeforeAndAfterAll {
assert(dataFiles.nonEmpty)
assert(dataFiles.forall(_.contains("worker_42_")))
}
+ it should "resolve its table from the catalog installed at flush time, not
at construction" in {
+ // #7290: an eagerly-held Table would keep the writer bound to the catalog
that was
+ // cached when it was constructed -- a reference the bounded cache may
evict and
+ // close. Swapping the warehouse entry after construction must redirect
the flush.
+ val swapWarehouse = "iceberg-table-writer-spec-swap"
+ val tableName = s"tbl_${UUID.randomUUID().toString.replace("-", "")}"
+ IcebergUtil.createTable(
+ catalog,
+ tableNamespace,
+ tableName,
+ icebergSchema,
+ overrideIfExists = true
+ )
+
+ val constructionCatalog = new CountingCatalog(catalog)
+ IcebergCatalogInstance.replaceInstance(constructionCatalog,
Some(swapWarehouse))
+ val writer = new IcebergTableWriter[Tuple](
+ "writer_swap",
+ Some(swapWarehouse),
+ tableNamespace,
+ tableName,
+ icebergSchema,
+ IcebergUtil.toGenericRecord
+ )
+ writer.open()
+ writer.putOne(tuple(1))
+
+ // Install a different catalog before the flush; the writer must go
through it.
+ val flushCatalog = new CountingCatalog(catalog)
+ IcebergCatalogInstance.replaceInstance(flushCatalog, Some(swapWarehouse))
+ val loadsBeforeFlush = flushCatalog.loadTableCalls.get()
+
+ writer.close() // flushes the buffer
+
+ assert(
+ flushCatalog.loadTableCalls.get() > loadsBeforeFlush,
+ "the flush must load the table through the catalog installed at flush
time"
+ )
+ assert(readTuples(tableName) == List(tuple(1)), "the tuple must reach the
table")
+ }
+
}
diff --git
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/OnIcebergSpec.scala
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/OnIcebergSpec.scala
index 24b1e08158..b230aad94e 100644
---
a/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/OnIcebergSpec.scala
+++
b/common/workflow-core/src/test/scala/org/apache/texera/amber/core/storage/result/iceberg/OnIcebergSpec.scala
@@ -20,6 +20,7 @@
package org.apache.texera.amber.core.storage.result.iceberg
import org.apache.texera.amber.core.tuple.{AttributeType, Schema, Tuple}
+import org.apache.texera.amber.core.storage.IcebergCatalogInstance
import org.apache.texera.amber.util.IcebergUtil
import org.apache.iceberg.catalog.Catalog
import org.apache.iceberg.data.IcebergGenerics
@@ -56,9 +57,14 @@ class OnIcebergSpec extends AnyFlatSpec with
BeforeAndAfterAll {
private val icebergSchema: IcebergSchema =
IcebergUtil.toIcebergSchema(amberSchema)
+ // The writer resolves its catalog from the shared cache per use (#7290), so
the
+ // spec's local catalog is registered under a spec-unique warehouse name.
+ private val specWarehouse = "on-iceberg-spec"
+
override def beforeAll(): Unit = {
warehouseDir = Files.createTempDirectory("on-iceberg-spec")
catalog = IcebergUtil.createHadoopCatalog("on-iceberg-spec", warehouseDir)
+ IcebergCatalogInstance.replaceInstance(catalog, Some(specWarehouse))
}
override def afterAll(): Unit = {
@@ -88,7 +94,7 @@ class OnIcebergSpec extends AnyFlatSpec with
BeforeAndAfterAll {
private def appendSnapshot(tableName: String, ids: Seq[Int]): Unit = {
val writer = new IcebergTableWriter[Tuple](
s"writer_${UUID.randomUUID().toString.replace("-", "")}",
- catalog,
+ Some(specWarehouse),
tableNamespace,
tableName,
icebergSchema,