PaulLiang1 opened a new issue, #17485:
URL: https://github.com/apache/iceberg/issues/17485
### Apache Iceberg version
`1.10.1` (dropdown selection). The defective declaration is **byte-identical
in `1.11.0` and on
`main`** at the time of writing.
### Query engine
Spark (dropdown selection) β Structured Streaming on Spark `3.5.3`. The
defect is not Spark-specific;
it needs only more than 100 concurrent `S3FileIO` instances in one JVM with
AAL enabled.
### Catalog configuration used
```properties
spark.sql.catalog.demo =
org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.demo.catalog-impl =
org.apache.iceberg.aws.glue.GlueCatalog
spark.sql.catalog.demo.io-impl =
org.apache.iceberg.aws.s3.S3FileIO
spark.sql.catalog.demo.warehouse = s3://<bucket>/<prefix>
spark.sql.catalog.demo.s3.analytics-accelerator.enabled = true
spark.sql.catalog.demo.s3.crt.enabled = false # reproduces
with either client
```
Everything else, in Iceberg and in AAL, is at its default β in particular
AAL's
`physicalio.thread.pool.size` (96 threads **per factory**),
`max.memory.limit` (2 GB per factory) and
`small.objects.prefetching.enabled` (`true`). The reproduction bundle points
`s3.endpoint` at a local
MinIO instance; nothing depends on that choice.
### Full version set
| Component | Version |
|---|---|
| `org.apache.iceberg:iceberg-aws`, `iceberg-core` | `1.10.1` (also
`1.11.0`, `main`) |
| `software.amazon.s3.analyticsaccelerator:analyticsaccelerator-s3` |
`1.3.1` |
| `org.apache.spark:spark-sql_2.12` | `3.5.3` |
| `software.amazon.awssdk` (s3, kms, sts, glue, dynamodb) | `2.29.52`
(affected production path used BOM `2.42.13`) |
| `software.amazon.awssdk.crt:aws-crt` | `0.43.4` β **not required**;
reproduced with `s3.crt.enabled=false` |
| `com.github.ben-manes.caffeine:caffeine` | `3.1.8` in the reproduction |
| JDK | Corretto `17.0.15`; also `eclipse-temurin:17` |
| Kubernetes (containerised runs) | `kind` node image `kindest/node:v1.35.0`
|
| S3 endpoint for reproduction | MinIO `RELEASE.2025-04-22T22-12-26Z` |
---
## Please describe the bug π
**In one line:** with the S3 Analytics Accelerator enabled, an application
that holds more than 100
`S3FileIO` instances in one JVM will silently stop making progress β reads
hang forever instead of
failing, so nothing crashes and nothing alerts.
With `s3.analytics-accelerator.enabled=true`, a JVM holding **more than 100
live `S3FileIO`
instances** starts closing AAL reader thread pools that other threads are
actively reading through.
The reads do not fail loudly β they **hang forever** β so the application
appears healthy while doing
no work.
### What was observed
A Spark Structured Streaming driver running several dozen concurrent
queries, each with its own
Spark session and therefore its own catalog and `S3FileIO`:
- The Spark application stayed `RUNNING`; the driver was healthy.
- **About three-quarters of the driver's queries were absent from the Spark
UI.** They hung before
`query.start()` registered them with the `StreamingQueryManager`, so they
were not shown as
failed β they were not shown at all.
- No exception, no crash, no retry, no application-level alert.
- Source lag grew silently for several hours.
- The only evidence visible from inside the application was AAL telemetry:
**tens of thousands of**
`[failure]` lines, all on `*.metadata.json` reads, all
`RejectedExecutionException`. None before
AAL was enabled, none after it was disabled.
A representative line:
```
[failure] block.manager.make.range.available(generation=0, thread_id=<n>,
range=8000-15999,
etag="<etag>", uri=s3://<bucket>/<table>/metadata/<n>.metadata.json,
range.effective=8000-73535): <duration> ns
[java.util.concurrent.RejectedExecutionException: 'Task β¦ rejected from
ThreadPoolExecutor@β¦[Shutting down, pool size = 1, active threads = 1,
completed tasks = 0]']
```
Thread dump of a hung reader (from the attached reproduction):
```
java.util.concurrent.CountDownLatch.await(CountDownLatch.java:230)
β¦io.physical.data.Block.awaitData(Block.java:187)
β¦io.physical.data.Block.read(Block.java:125)
β¦io.physical.data.Blob.read(Blob.java:164)
β¦io.physical.impl.PhysicalIOImpl.read(PhysicalIOImpl.java:159)
β¦S3SeekableInputStream.read(S3SeekableInputStream.java:151)
```
### The sequence, in order
```text
Query A thread factory cache Factory F Query Z
thread
(holds factory F) (static, max 100) (96 readers) (another
query)
| | | |
1 |---- get(key A) ---> | |
<---- factory F ----| | |
| A begins reading table metadata through F |
| | | |
2 | <---------------------------------|
| get(key Z) -- the 101st distinct key |
| past maximumSize(100), so Caffeine must evict
something.
| It picks factory F. Query A is never consulted. |
| | | |
3 | |=================> |
| removalListener -> close() -> shutdown() |
| | | |
4 |------- read() its next range -------> |
5 <===== RejectedExecutionException ====| |
| AAL leaves its read buffer registered but unfilled,
then
| waits on it with no timeout -> parks forever |
| | | |
| query.start() is never reached, so the query is never
| registered: no UI entry, no metrics, no crash, no
retry
v
time
```
---
## Root cause in Iceberg
`aws/src/main/java/org/apache/iceberg/aws/s3/AnalyticsAcceleratorUtil.java:47-55`:
```java
private static final Cache<Pair<S3AsyncClient, S3FileIOProperties>,
S3SeekableInputStreamFactory>
STREAM_FACTORY_CACHE =
Caffeine.newBuilder()
.maximumSize(100)
.removalListener(
(RemovalListener<
Pair<S3AsyncClient, S3FileIOProperties>,
S3SeekableInputStreamFactory>)
(key, factory, cause) -> close(factory))
.build();
```
Three properties combine badly:
**1. The key is object identity, so an entry can never be shared across
`S3FileIO` instances.**
To be precise, because this cuts both ways: *within* one `S3FileIO` the
cache works as intended.
`S3InputFile.fromLocation` passes `client.s3Async()` and
`client.s3FileIOProperties()`
(`S3InputFile.java:127-135`), `PrefixedS3Client` holds the properties as a
`private final` field and
memoises the async client (`PrefixedS3Client.java:33,37,97-106`), and
`org.apache.iceberg.util.Pair`
delegates `equals`/`hashCode` to its components β so every read after the
first through the same
`S3FileIO` hits.
What cannot happen is sharing *between* instances. `PrefixedS3Client`
constructs a **fresh**
`S3FileIOProperties` per instance (`PrefixedS3Client.java:50`) and builds
its **own**
`S3AsyncClient`, and neither type overrides `equals`/`hashCode`. So two
`S3FileIO` instances with
byte-identical configuration are always two entries, and **the number of
entries equals the number of
live `S3FileIO` instances**. The bound therefore acts as a population limit
on `S3FileIO`, which is
not what a reader of this code would expect a read-path cache to be.
**2. The removal listener cannot distinguish who initiated the removal.** It
closes the factory for
*every* `RemovalCause`. That is correct for `EXPLICIT` β the `invalidate()`
that `cleanupCache()`
performs when the owning `S3FileIO` closes. It is unsafe for `SIZE`, which
Caffeine decides on its
own, with no knowledge of whether a caller obtained that factory from
`get()` moments earlier and is
mid-read.
**3. `close()` is destructive and immediate.**
`S3SeekableInputStreamFactory.close()` ends in
`threadPool.shutdown()`, so every subsequent read submitted by any holder is
rejected.
What AAL then does with that rejection is why this surfaces as a hang rather
than an error, and it is
worth one sentence of gloss since it is another project's internals: AAL
registers its read buffer
*before* submitting the work that fills it, the submit fails, nothing marks
the buffer as failed, and a
reader waiting on that buffer waits on a signal that will never come β with
no timeout
(`Block.awaitData()`). That half is filed separately as
awslabs/analytics-accelerator-s3#369, with a reproduction that uses
no Iceberg at all. **This issue is about Iceberg closing a shared resource
that callers still
hold.**
### Why `S3FileIO` instance count scales with concurrency
This is worth stating because "more than 100 `S3FileIO` in one JVM" sounds
unlikely until you look at
how engines create them:
- `FileIO` lifetime in Iceberg is effectively per-`TableOperations`;
`FileIOTracker` closes it only
when the `TableOperations` is collected.
- In Spark, **each session gets its own `CatalogManager`** and therefore its
own catalog plugin
instances. `SparkSession.newSession()` passes `parentSessionState = None`
(`SparkSession.scala:251-258`), and `CatalogManager` holds its own
`catalogs` map
(`CatalogManager.scala:48-54`).
- Independently of that, Spark's `StreamExecution` constructor does
`sparkSession.cloneSession()` (`StreamExecution.scala:197`, v3.5.3), so
**every streaming query
gets its own session β and its own catalogs and `S3FileIO` β whether or
not the application forks
sessions itself.**
So an application with N concurrent streaming queries and M AAL-enabled
catalogs holds on the order
of N Γ M live `S3FileIO` instances. A few dozen queries across two catalogs
already exceeds the
bound of 100; applications with hundreds of queries are far past it.
---
## Evidence
### The cache behaviour, driven directly
The reproduction inline at the bottom of this issue loads the real
`STREAM_FACTORY_CACHE` out of the shipped
`iceberg-aws-1.10.1` jar by reflection β nothing about the cache is
re-implemented. No AWS account, no
network, ~45 s:
```
130 S3FileIO-equivalent lookups -> factories created = 130 <-- no
cross-instance sharing
cache size after 130 inserts = 100
first factory still referenced by caller? yes
first factory still IN the cache? false
first factory's reader pool isShutdown() = true
```
Then a read through that still-referenced factory hangs indefinitely (stack
above).
### Controlled experiment β the bound is the variable
Against a real S3 endpoint, **readers held constant at 130** so thread
count, socket count and
endpoint load are identical across rows; the only variable is the number of
distinct `S3FileIO`
instances, i.e. distinct cache keys:
| AAL | Readers | Distinct `S3FileIO` | OK | Threw | **Hung** |
|---|---|---|---|---|---|
| disabled | 130 | 130 | 130 | 0 | 0 |
| enabled | 130 | 99 | 130 | 0 | 0 |
| enabled | 130 | **130** | 100 | 20 | **10** |
| enabled | 130 | 1 | 130 | 0 | 0 |
In the broken row **exactly 100 readers succeed** β the value of
`maximumSize`.
### Real Spark, separate executor JVMs
To check the executor side, a further harness runs Spark `3.5.3` in
`local-cluster[2,8,3072]` (two genuinely separate executor JVMs), 260 tasks,
each task carrying its
own `FileIO` serialized as task data β the way `SerializableTable` ships a
table's `FileIO` to the
executors that read it. The `FileIO` is built through
`CatalogUtil.loadFileIO(io-impl, props, conf)`.
| `io-impl` | Executor JVMs | Distinct `FileIO` per JVM | Reads OK |
|---|---|---|---|
| `S3FileIO` | 2 | **125 / 135** | 260 |
| a JVM-singleton `FileIO` (see workarounds) | 2 | **1 / 1** | 260 |
Two conclusions. The instance multiplication is real on executors too β
125β135 cache keys per JVM.
But note the "Reads OK" column: **the hang did not reproduce on the executor
side**, even well past
the bound, and raising executor cores from 2 to 8 did not change that. An
executor's read is
short-lived, so eviction tends to close a factory whose reader has already
returned. The driver is
exposed because it holds many *concurrent, long-lived* reads while the cache
churns underneath them.
So on executors this manifests as resource cost rather than the hang.
---
## Proposed fix
Options, worst to best:
1. **Give `S3FileIOProperties` value equality.** Collapses one half of the
key, but the
`S3AsyncClient` half is still per-instance, so cardinality barely
changes. *Insufficient.*
2. **Remove the size bound.** No eviction, so no premature close β but each
abandoned factory keeps
96 reader threads, a maintenance thread and a 2 GB blob-store budget,
released only by
`cleanupCache`. This trades a liveness bug for a leak, and the leak
already has an open report
(#15898, below). *Not recommended.*
3. **Reference-count the cached factory** and close only when the last
stream closes. Correct in
principle, but callers routinely abandon streams to the GC, so a naive
refcount may never reach
zero. Would need `Cleaner`/phantom-reference backup in a hot path.
4. **Bind the factory's lifetime to the owner's, not to a cache bound.** One
factory per
`S3AsyncClient`, created with it and closed with it. `PrefixedS3Client`
already owns exactly that
lifecycle and already calls `cleanupCache` on close. No size bound is
needed, because the
population is then bounded by the number of clients. **Recommended.**
### There is a precedent in this repository
Of the twelve files in `apache/iceberg` that attach a Caffeine
`removalListener`, I inspected the
nine distinct ones. **`AnalyticsAcceleratorUtil` is the only one that pairs
a hard `maximumSize`
bound with closing a shared resource in the listener.**
| Cache | Eviction driven by | Closes a resource on removal? |
|---|---|---|
| `AnalyticsAcceleratorUtil` | **hard `maximumSize(100)`** | **yes** |
| `io.FileIOTracker` | `weakKeys()` β reachability | yes, but eviction
implies nobody holds the key |
| `rest.auth.AuthSessionCache` | `expireAfterAccess` | yes β time-driven,
re-creatable resource |
| `hive.CachedClientPool` | `expireAfter` | yes β time-driven |
| `io.ContentCache` | `softValues` + expiry | no |
| `CachingCatalog` | soft/weak + expiry | no |
| `spark.SparkExecutorCache` | size + expiry | no |
| `ManifestFiles` | size + weak/soft | no |
Every other close-on-removal cache in the project is driven by
**reachability** or **time**, never by
a count β because a count carries no information about whether anyone is
still using the value.
`FileIOTracker`'s `weakKeys()` pattern is the in-repo model for option 4.
### A minimal interim change
If a full lifetime redesign is too large for a point release, the smallest
correct-direction change
is to discriminate on `RemovalCause` β close on `EXPLICIT`/`REPLACED`
(owner-initiated) and not on
`SIZE`. I have that patch and it does stop the hang (verified: **130 readers
/ 130 keys β 130 ok,
0 failed, 0 hung**, against 100/20/10 unpatched). But I want to be straight
about the cost: it
converts the hang into exactly the leak described in option 2, so it is a
stopgap, not the fix.
I have working patches for both the interim change and the AAL-side fixes
(the latter attached to
awslabs/analytics-accelerator-s3#369), and can share them or open a PR.
---
## Related existing issues
- **#15898** (open) β *CachingCatalog does not close FileIO on cache
eviction, causing S3FileIO /
SDK v2 thread leak in long-running applications.* This is the mirror image
of the same missing
concept: there, `FileIO` instances proliferate and are **never** closed;
here, they are closed **too
eagerly, by the wrong party**. Both point at the absence of a clear owner
for `FileIO` and factory
lifetime, and I would suggest they be considered together.
- **#12891** β *AWS: Close the S3SeekableInputStreamFactory before removing
from cache* (merged
2025-05-26) introduced the `removalListener`. It was fixing a genuine
leak; the PR body does not
discuss which removal causes are safe to close on, which is the gap.
- **#12299** β the original AAL integration, which introduced the cache with
`maximumSize(100)` and
the identity key. There is no stated rationale for the bound in the PR
body, commit message, or
review discussion. Worth noting that a reviewer asked for wider review at
the time β
*"I would like more eyes on this PR since AWS FileIO has a pretty big
blast radius. Have you posted
this on the iceberg devlist?"* β and it was not taken to the dev list.
- **#13133 / #13134** β an earlier bug in the same `S3FileIO.close()` β AAL
cleanup path.
- **#12799 / #12827** β relevant to the workaround below: per-storage-prefix
clients exist
deliberately so vended credentials stay scoped, and the
single-shared-client alternative (#12827)
was closed unmerged.
### Please gate the default-on work on this
Epic **#14350** (*Turn S3 Analytics Accelerator on by default*) was closed
`not_planned` by a stale
bot with the "Default On" item unticked β it stalled rather than being
decided against. AAL is
default-off in `1.10.1`, `1.11.0` and `main`, so today the blast radius is
opt-in users. If that epic
is revived before this is fixed, the failure ships to every `S3FileIO` user.
---
## Workarounds for anyone hitting this now
| Workaround | Effect |
|---|---|
| **Set `s3.analytics-accelerator.enabled=false`** on the catalog |
Prevents. Complete and immediate. |
| **Share one `S3FileIO` per JVM** via a delegating `io-impl` (see below) |
Prevents: one cache key, so eviction never fires. Measured clean at 130
concurrent readers, 1 `FileIO` identity per executor JVM, and *faster* than the
AAL-disabled baseline. |
| Enable AAL on only one catalog, or only on applications with few
concurrent queries | Reduces probability only. Establishes no invariant and
does not survive scaling up. |
| Tune `s3.analytics-accelerator.physicalio.thread.pool.size` /
`max.memory.limit` | No effect. These are per-factory resource knobs; neither
reads the cache bound, which is a `private static final` literal. |
| Set `cache-enabled=true` on the catalog | No effect. `CachingCatalog`
caches `Table` objects, not `FileIO` instances. |
| Disable AAL small-object prefetching | **Do not.** Measured worse:
converts thrown errors into silent hangs (0 thrown / 30 hung vs 20 / 10). |
### Notes on the shared-`FileIO` workaround
I have a working implementation of this. Four
things it must get right, each of which fails silently if missed:
1. **`readResolve()`** returning the JVM singleton. Without it, every task
that deserializes the
`FileIO` builds its own delegate and the multiplication returns per
executor with no error.
Verified: 1 identity and 1 delegate per executor JVM across 2 JVMs and
260 tasks.
2. **A no-op `close()`.** Iceberg closes a `FileIO` per table/broadcast
lifecycle; if that closed the
shared delegate, one table finishing would break every other reader.
3. **Every capability interface.** `S3FileIO` implements `DelegateFileIO`,
`SupportsRecoveryOperations` and `SupportsStorageCredentials`, and
Iceberg probes these with
`instanceof`. A wrapper missing one silently loses the capability β for
example bulk delete
quietly degrading to per-file.
4. **Vended credentials are a hard blocker.** A JVM-wide delegate can only
hold one credential set,
so this is **unsafe for REST catalogs that vend per-prefix credentials**.
The reference
implementation throws from `setCredentials` rather than silently applying
one catalog's credentials
to another. Use it only for statically-credentialed catalogs.
Also worth stating: this is mitigation by staying under an undocumented
library constant, not a fix.
It needs an invariant test on live instance count or it regresses the next
time an application scales
up.
**Detection**, since the failure is silent: alert on the AAL telemetry
failure line
(`block.manager.make.range.available` together with `failure`), which has a
zero baseline in a healthy
system. Do not rely on application health signals β they are emitted by the
component that hung.
---
## Willingness to contribute
- [x] I can contribute a fix for this bug independently
Patches for both the interim `RemovalCause` change and the AAL-side fixes
have been compiled and
tested against the released artifacts. I would welcome direction on which
of the four options above the maintainers prefer before opening a PR, since
option 4 touches
`PrefixedS3Client` lifetime and is a larger change than a point fix.
---
**AI Disclosure**
- Model: Claude Opus 4.6
- Platform/Tool: Claude Code
- Human Oversight: fully reviewed
- Prompt Summary: Investigate a production incident in which streaming
queries silently failed to
start with the S3 Analytics Accelerator enabled; identify root cause from
pinned sources, build
runnable reproductions, verify candidate patches, and draft an upstream
bug report.
---
<details>
<summary><b>Full reproduction source</b> β
<code>IcebergCacheMre.java</code>, drives the real static cache by
reflection</summary>
```java
package software.amazon.s3.analyticsaccelerator;
import com.github.benmanes.caffeine.cache.Cache;
import java.io.ByteArrayInputStream;
import java.io.IOException;
import java.lang.reflect.Field;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.iceberg.aws.s3.S3FileIOProperties;
import org.apache.iceberg.util.Pair;
import software.amazon.awssdk.services.s3.S3AsyncClient;
import software.amazon.s3.analyticsaccelerator.request.GetRequest;
import software.amazon.s3.analyticsaccelerator.request.HeadRequest;
import software.amazon.s3.analyticsaccelerator.request.ObjectClient;
import software.amazon.s3.analyticsaccelerator.request.ObjectContent;
import software.amazon.s3.analyticsaccelerator.request.ObjectMetadata;
import software.amazon.s3.analyticsaccelerator.util.OpenStreamInformation;
import software.amazon.s3.analyticsaccelerator.util.S3URI;
/**
* Minimal reproduction of the ICEBERG half of the defect: {@code
AnalyticsAcceleratorUtil}'s static,
* identity-keyed, {@code maximumSize(100)} factory cache closes an {@code
* S3SeekableInputStreamFactory} on size eviction, i.e. while a caller still
holds and uses it.
*
* <p>Drives the REAL static cache out of the shipped {@code iceberg-aws}
jar by reflection β nothing
* about the cache is re-implemented or simulated. No AWS account, no
network, no credentials.
*
* <p>Lives in AAL's package only so it can read the package-private {@code
getThreadPool()} accessor
* to show that an evicted factory's pool has been shut down.
*
* <p>The consequence of that shutdown β a read that hangs forever rather
than failing β is a separate
* AAL defect, reproduced independently by {@code AalMre} with no Iceberg on
the classpath at all.
*
* <p>Usage: {@code IcebergCacheMre [broken|fixed]}. Exits non-zero on unmet
expectations.
*/
public final class IcebergCacheMre {
private static final int OBJECT_LEN = 128 * 1024;
private static String mode = "broken";
private static final java.util.List<String> FAILURES = new
java.util.ArrayList<>();
private static boolean fixed() {
return "fixed".equals(mode);
}
public static void main(String[] args) throws Exception {
if (args.length > 0) {
mode = args[0];
}
System.out.println("iceberg factory-cache reproduction; expectation mode
= " + mode);
step1IdentityKeys();
step2CacheEvictsAndClosesLiveFactory();
banner("SUMMARY");
if (FAILURES.isEmpty()) {
System.out.println(" ALL EXPECTATIONS MET for mode=" + mode);
Runtime.getRuntime().halt(0);
}
System.out.println(" UNMET EXPECTATIONS for mode=" + mode + ":");
FAILURES.forEach(f -> System.out.println(" - " + f));
Runtime.getRuntime().halt(1);
}
private static void step1IdentityKeys() {
banner("STEP 1 cache key is identity, so distinct S3FileIO instances
never share an entry");
Map<String, String> props = new HashMap<>();
props.put("s3.analytics-accelerator.enabled", "true");
// PrefixedS3Client does exactly this, once per S3FileIO: `new
S3FileIOProperties(properties)`.
S3FileIOProperties a = new S3FileIOProperties(props);
S3FileIOProperties b = new S3FileIOProperties(props);
S3AsyncClient client = fakeAsyncClient();
Pair<S3AsyncClient, S3FileIOProperties> k1 = Pair.of(client, a);
Pair<S3AsyncClient, S3FileIOProperties> k2 = Pair.of(client, b);
System.out.println(" same properties map -> a.equals(b) = " +
a.equals(b));
System.out.println(" same async client -> k1.equals(k2) = " +
k1.equals(k2));
System.out.println(" => two S3FileIO instances with IDENTICAL config
are two distinct keys.");
System.out.println(" (within ONE S3FileIO the key components are
the same objects, so"
+ " repeat reads DO hit -- the cache is only unshareable ACROSS
instances.)");
check(!a.equals(b), "S3FileIOProperties has no value equality, so two
instances are two keys");
check(!k1.equals(k2), "cache keys must be distinct");
}
private static void step2CacheEvictsAndClosesLiveFactory() throws
Exception {
banner("STEP 2 Iceberg's real STREAM_FACTORY_CACHE closes a factory
that is still in use");
Class<?> util =
Class.forName("org.apache.iceberg.aws.s3.AnalyticsAcceleratorUtil");
Field f = util.getDeclaredField("STREAM_FACTORY_CACHE");
f.setAccessible(true);
Cache<Pair<S3AsyncClient, S3FileIOProperties>,
S3SeekableInputStreamFactory> cache =
(Cache<Pair<S3AsyncClient, S3FileIOProperties>,
S3SeekableInputStreamFactory>) f.get(null);
System.out.println(" loaded " + util.getName() +
"#STREAM_FACTORY_CACHE from iceberg-aws-1.10.1");
// A tiny thread pool per factory keeps the MRE cheap; production
default is 96
// (PhysicalIOConfiguration.DEFAULT_THREAD_POOL_SIZE).
AtomicInteger created = new AtomicInteger();
S3SeekableInputStreamFactory first = null;
Pair<S3AsyncClient, S3FileIOProperties> firstKey = null;
for (int i = 0; i < 130; i++) {
Pair<S3AsyncClient, S3FileIOProperties> key =
Pair.of(fakeAsyncClient(), new S3FileIOProperties(new
HashMap<>()));
S3SeekableInputStreamFactory factory =
cache.get(
key,
k -> {
created.incrementAndGet();
return newFactory();
});
if (i == 0) {
first = factory;
firstKey = key;
}
}
cache.cleanUp(); // force Caffeine's pending maintenance so eviction
is deterministic here
System.out.println(" 130 S3FileIO-equivalent lookups -> factories
created = " + created.get());
System.out.println(" cache size after 130 inserts = " +
cache.estimatedSize());
System.out.println(" first factory still referenced by caller? yes");
System.out.println(" first factory still IN the cache? " +
(cache.getIfPresent(firstKey) != null));
ExecutorService pool = first.getThreadPool();
System.out.println(" first factory's reader pool isShutdown() = " +
pool.isShutdown());
check(created.get() == 130,
"130 distinct S3FileIO-equivalents must produce 130 entries (no
cross-instance sharing)");
check(cache.estimatedSize() <= 100, "cache must be bounded at
maximumSize(100)");
if (fixed()) {
boolean ok = !pool.isShutdown();
System.out.println((ok ? " [PASS] " : " [FAIL] ")
+ "size eviction must NOT close a factory the caller still
holds");
if (!ok) FAILURES.add("evicted factory was closed despite the
Iceberg removal-cause patch");
System.out.println(" => the caller's factory survives eviction; its
pool is still usable.");
} else {
check(pool.isShutdown(),
"the removal listener must have closed a factory the caller
still holds");
System.out.println(" => a caller holding this factory now submits
into a shut-down pool.");
}
}
private static S3SeekableInputStreamFactory newFactory() {
return new S3SeekableInputStreamFactory(
new FakeObjectClient(),
S3SeekableInputStreamConfiguration.DEFAULT);
}
private static S3AsyncClient fakeAsyncClient() {
return (S3AsyncClient)
java.lang.reflect.Proxy.newProxyInstance(
AalMre.class.getClassLoader(),
new Class<?>[] {S3AsyncClient.class},
(proxy, method, methodArgs) -> {
if ("hashCode".equals(method.getName())) return
System.identityHashCode(proxy);
if ("equals".equals(method.getName())) return proxy ==
methodArgs[0];
if ("toString".equals(method.getName())) return
"fake-s3-async";
return null;
});
}
private static final class FakeObjectClient implements ObjectClient {
@Override
public ObjectMetadata headObject(HeadRequest r, OpenStreamInformation
i) {
return
ObjectMetadata.builder().contentLength(OBJECT_LEN).etag("etag-1").build();
}
@Override
public ObjectContent getObject(GetRequest r, OpenStreamInformation i) {
return ObjectContent.builder().stream(new ByteArrayInputStream(new
byte[OBJECT_LEN])).build();
}
@Override
public void close() throws IOException {}
}
private static void banner(String s) {
System.out.println();
System.out.println("================================================================");
System.out.println(s);
System.out.println("================================================================");
}
private static void check(boolean cond, String what) {
System.out.println((cond ? " [PASS] " : " [FAIL] ") + what);
if (!cond) throw new AssertionError(what);
}
}
```
</details>
--
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]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]