apache/iceberg · error · IllegalStateException
Cannot load current offset at snapshot
Error message
Cannot load current offset at snapshot %d, the snapshot was expired or removed
What it means
SyncSparkMicroBatchPlanner.validateCurrentSnapshotExists verifies that the snapshot referenced by the current streaming offset still exists. If the table's currentSnapshot lookup for currentOffset.snapshotId() returns null, the snapshot was expired or removed (e.g. by expireSnapshots or retention cleanup) and the micro-batch cannot be planned, so an IllegalStateException is thrown.
Solutions
- Restart the stream with a fresh checkpoint (or reset the streaming offset) so it resumes from the current snapshot.
- Increase history.expire.max-snapshot-age and history.expire.min-snapshots-to-keep so snapshots outlive maximum stream downtime.
- Coordinate expireSnapshots scheduling with streaming jobs; never expire snapshots referenced by active checkpoints.
- If historical replay is not required, recreate the streaming query from the table's latest snapshot.
Example fix
// before
ALTER TABLE t SET TBLPROPERTIES ('history.expire.max-snapshot-age'='3600'); // stream offline longer than 1h
// after
ALTER TABLE t SET TBLPROPERTIES ('history.expire.max-snapshot-age'='604800', 'history.expire.min-snapshots-to-keep'='100'); Defensive patterns
Strategy: validation
Validate before calling
Snapshot snap = table.snapshot(offset.snapshotId());
if (snap == null) {
// snapshot expired: reset checkpoint or fail fast with an operator-visible message
} Type guard
static boolean snapshotAvailable(Table table, StreamingOffset offset) {
return table.snapshot(offset.snapshotId()) != null;
} Try / catch
try {
planner.planFiles(offset);
} catch (IllegalStateException e) {
if (e.getMessage().contains("the snapshot was expired or removed")) {
// restart from the latest snapshot with a new checkpoint and alert operators
} else throw e;
} Prevention
- Set history.expire.max-snapshot-age longer than worst-case stream downtime.
- Pause expireSnapshots jobs while streaming queries are offline or lagging.
- Monitor snapshot count/age against streaming checkpoint lag.
- Test failover: expire snapshots while a stream is stopped, then restart and verify.
When it happens
Trigger: Spark structured streaming over an Iceberg table where the snapshotId recorded in the streaming offset no longer exists — typically after expireSnapshots() ran with retention shorter than the stream's lag, or the table was rewritten/recreated while the stream was paused. Called from planFiles and latestOffset.
Common situations: Long streaming downtime combined with aggressive snapshot expiration (small history.expire.max-snapshot-age); scheduled maintenance jobs expiring snapshots while a stream consumes; failure to coordinate table retention settings with streaming checkpoint lifetimes.
Understand the failure class
Background: "Not found" and "does not exist" errors: why "Task not found", "No such folder", and "Can't find" fire when a lookup comes back empty — this error's family across 14 libraries.
Related errors
- ALTER VIEW AS is not supported. Use CREATE OR REPLACE VIEW…
- apply(value) is deprecated, use bind(Type).apply(value)
- apply(value) is deprecated, use bind(Type).apply(value)
- AS OF is not supported for changelogs
- bind is not implemented
AI-assisted analysis of apache/iceberg@86d9c8fc54 (2026-09-12).
Data as JSON: /api/errors/bac2aba41c2a9ae0.
Report an issue: GitHub.
Appendix: source
Thrown at spark/v4.0/spark/src/main/java/org/apache/iceberg/spark/source/SyncSparkMicroBatchPlanner.java:242
startPosOfSnapOffset = -1;
// if anyhow we are moving to next snapshot we should only scan addedFiles
scanAllFiles = false;
}
}
StreamingOffset latestStreamingOffset =
new StreamingOffset(curSnapshot.snapshotId(), curPos, scanAllFiles);
// if no new data arrived, then return null.
return latestStreamingOffset.equals(startingOffset) ? null : latestStreamingOffset;
}
@Override
public void stop() {}
private void validateCurrentSnapshotExists(Snapshot snapshot, StreamingOffset currentOffset) {
if (snapshot == null) {
throw new IllegalStateException(
String.format(
Locale.ROOT,
"Cannot load current offset at snapshot %d, the snapshot was expired or removed",
currentOffset.snapshotId()));
}
}
}
View on GitHub (pinned to 86d9c8fc54)