krisnaru opened a new issue, #17830:
URL: https://github.com/apache/iceberg/issues/17830

   ### Apache Iceberg version
   
   1.11.0 (latest release)
   
   ### Query engine
   
   Spark
   
   ### Please describe the bug 🐞
   
     ## Summary
     
   When a table is loaded with `snapshot-loading-mode=refs`, 
`RESTSessionCatalog.loadTable()` wraps the returned `TableMetadata` with a 
`snapshotsSupplier` so that touching a snapshot outside the retained-refs 
window transparently re-fetches the full snapshot history (`snapshots=all`). 
This lazy fallback is **lost on every subsequent `table.refresh()` (and 
commit)**, because `RESTTableOperations` neither sends the snapshot mode on 
refresh nor re-installs the supplier. The refreshed metadata is therefore 
*partial with no fallback*, and any  code path that resolves a snapshot outside 
the refs window returns `null` and throws.
   
   Batch reads are unaffected (they resolve through the initially-loaded 
metadata). **Spark Structured Streaming and Flink streaming are  affected**, 
because both call `table.refresh()` every micro-batch/cycle and then walk 
snapshot ancestry incrementally.
   
     ## Root cause
   
     Two spots in `RESTTableOperations`:
     
     1. **`refresh()`** (`RESTTableOperations.java:150-154`) issues a plain GET 
with no `snapshots` param — it does not mirror the `SnapshotMode` handling in 
`RESTSessionCatalog.loadInternal()`:
        ```java
        return updateCurrentMetadata(
            client.get(path, LoadTableResponse.class, readHeaders, 
ErrorHandlers.tableErrorHandler()));
        ```
   
     2. **`updateCurrentMetadata()`** (`RESTTableOperations.java:288-298`) 
stores the response verbatim (via `checkUUID`) with no supplier re-install:
        ```java
        this.current = checkUUID(current, response.tableMetadata());
        ```
   
   So even if the server withheld history, nothing re-installs the lazy loader 
that `RESTSessionCatalog.loadTable()` installed on the initial load 
(`RESTSessionCatalog.java:543-554`). `RESTTableOperations` already holds 
everything the supplier needs (`client`, `path`,`readHeaders`) — the supplier 
is just a `GET path?snapshots=all`. No catalog plumbing, wire-format, or spec 
change is required.
     
     ## Reproduction
     1. Load a table via REST with `snapshot-loading-mode=refs`.
     2. Advance the table so the current metadata location changes and old 
snapshots fall outside the retained refs.
     3. Call `table.refresh()`.
     4. Resolve a snapshot id that is outside the refs window — what streaming 
does when the consumed offset lags retained refs (backlog, restart from an old 
checkpoint, or start-from-timestamp / oldest-ancestor).
   
     **Expected:** the snapshot is lazily re-fetched, as it is right after 
`loadTable`.
     **Actual:** `snapshotsLoaded == true`, supplier is `null`, resolution 
returns `null`, and streaming throws — Spark: *"Cannot load current offset … 
expired or removed"*; Flink: *"Cannot find snapshot"*.
   
     ## Proposed fix (client-only, ~1 class)
   
     1. Thread `RESTCatalogProperties.SnapshotMode` into `RESTTableOperations` 
as a field (constructor param), passed from the `newTableOps(...)` builders 
(`RESTSessionCatalog.java:1257` and `:1289`), which already have `snapshotMode` 
in scope. `SnapshotMode` is already a public enum  in `RESTCatalogProperties` 
(`RESTCatalogProperties.java:70-72`) — **no visibility change needed.**
     
     2. `refresh()` sends the mode, mirroring `loadInternal`:
        ```java
        client.get(path, snapshotModeToParam(mode), LoadTableResponse.class,
            readHeaders, ErrorHandlers.tableErrorHandler());
        ```
   
     3. `updateCurrentMetadata()` re-installs the supplier in `REFS` mode, 
reusing the exact pattern from `RESTSessionCatalog.loadTable` (lines
     543-554):
        ```java
        private TableMetadata updateCurrentMetadata(LoadTableResponse response) 
{
          if (current == null
              || !Objects.equals(current.metadataFileLocation(), 
response.metadataLocation())) {
            TableMetadata refreshed = checkUUID(current, 
response.tableMetadata());
            if (snapshotMode == SnapshotMode.REFS) {
              refreshed = TableMetadata.buildFrom(refreshed)
                  .withMetadataLocation(response.metadataLocation())
                  .setPreviousFileLocation(null)
                  .setSnapshotsSupplier(() ->
                      client.get(path, snapshotModeToParam(SnapshotMode.ALL), 
LoadTableResponse.class,
                          readHeaders, ErrorHandlers.tableErrorHandler())
                        .tableMetadata().snapshots())
                  .discardChanges().build();
            }
            this.current = refreshed;
          }
          return current;
     4. Resolve a snapshot id that is outside the refs window — what streaming 
does when the consumed offset lags retained refs (backlog, restart  from an old 
checkpoint, or start-from-timestamp / oldest-ancestor).
   
     **Expected:** the snapshot is lazily re-fetched, as it is right after 
`loadTable`.
     **Actual:** `snapshotsLoaded == true`, supplier is `null`, resolution 
returns `null`, and streaming throws — Spark: *"Cannot load current offset … 
expired or removed"*; Flink: *"Cannot find snapshot"*.
   
     ## Proposed fix (client-only, ~1 class)
   
     1. Thread `RESTCatalogProperties.SnapshotMode` into `RESTTableOperations` 
as a field (constructor param), passed from the `newTableOps(...)`
     builders (`RESTSessionCatalog.java:1257` and `:1289`), which already have 
`snapshotMode` in scope. `SnapshotMode` is already a public enum in 
`RESTCatalogProperties` (`RESTCatalogProperties.java:70-72`) — **no visibility 
change needed.**
   
     2. `refresh()` sends the mode, mirroring `loadInternal`:
        ```java
        client.get(path, snapshotModeToParam(mode), LoadTableResponse.class,
            readHeaders, ErrorHandlers.tableErrorHandler());
        ```
   
     3. `updateCurrentMetadata()` re-installs the supplier in `REFS` mode, 
reusing the exact pattern from `RESTSessionCatalog.loadTable` (lines 543-554):
        ```java
        private TableMetadata updateCurrentMetadata(LoadTableResponse response) 
{
          if (current == null
              || !Objects.equals(current.metadataFileLocation(), 
response.metadataLocation())) {
            TableMetadata refreshed = checkUUID(current, 
response.tableMetadata());
            if (snapshotMode == SnapshotMode.REFS) {
              refreshed = TableMetadata.buildFrom(refreshed)
                  .withMetadataLocation(response.metadataLocation())
                  .setPreviousFileLocation(null)
                  .setSnapshotsSupplier(() ->
                      client.get(path, snapshotModeToParam(SnapshotMode.ALL), 
LoadTableResponse.class,
                          readHeaders, ErrorHandlers.tableErrorHandler())
                        .tableMetadata().snapshots())
                  .discardChanges().build();
            }
            this.current = refreshed;
          }
          return current;
        }
        ```
        (`snapshotModeToParam` is currently a private static helper in 
`RESTSessionCatalog`, line 438 — either duplicate the one-liner or lift it to 
`RESTCatalogProperties` alongside the enum.)
   
    ## Behavior / cost
   
     - In the default (`ALL`) mode this is a **no-op** — identical to today 
(`SNAPSHOT_LOADING_MODE_DEFAULT = SnapshotMode.ALL`,  
`RESTCatalogProperties.java:32`).
     - In `refs` mode it's strictly a correctness improvement: a lagging 
streaming job triggers one full `snapshots=all` fetch on the refresh cycle that 
actually needs history. Only jobs that need history pay, and it's far cheaper 
than every `loadTable` returning full history.
   
    ## Tests
   
     - A `TestRESTCatalog`/mock case that loads a table in `refs` mode, 
refreshes to a new metadata location, then resolves an old (out-of-refs) 
snapshot id and asserts it lazily reloads — today that returns `null`.
     - A streaming-style ancestry-walk assertion (`SnapshotUtil.snapshotAfter` 
/ `ancestorsBetween` / `oldestAncestor`) after a refresh.
   
   ## Relationship to #14398
   
   PR #14398 (*Core: Freshness-aware table loading in REST catalog*) is 
adjacent but does **not** address this. It adds ETag/304 caching on the  
`loadTable` path, explicitly keeps `RESTTableOperations` out of freshness-aware 
loading, and explicitly **defers** the  refs/partial-snapshot interaction 
(caching a "partially loaded snapshot list" whose lazy-loaded remainder "won't 
be reflected in the cache"
     — "too complicated … for the initial version"). Verified on `main`: the 
ETag/`tableCache` path exists only in `loadTable`; 
`RESTTableOperations.refresh()` still issues an unconditioned GET. This issue 
tracks the `refresh()` supplier drop that #14398 left open.
   
     ## Why it matters
   
   This is the prerequisite for safely changing the REST `loadTable` default to 
`refs`: with the fix, `refs` is safe for both batch and streaming on Spark and 
Flink. Without it, streaming clients must stay on `all`.
   
   
   ### Willingness to contribute
   
   - [x] I can contribute a fix for this bug independently
   - [ ] I would be willing to contribute a fix for this bug with guidance from 
the Iceberg community
   - [ ] I cannot contribute a fix for this bug at this time


-- 
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]

Reply via email to