This is an automated email from the ASF dual-hosted git repository.

JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git


The following commit(s) were added to refs/heads/main by this push:
     new 5ed5f77d feat(table): support scan.watermark batch time travel (#677)
5ed5f77d is described below

commit 5ed5f77d7b4c050896bddfaa7038f4b3b05f9ccb
Author: kid <[email protected]>
AuthorDate: Sun Aug 9 21:08:21 2026 +0800

    feat(table): support scan.watermark batch time travel (#677)
---
 bindings/c/src/table.rs                            |  54 ++++++-
 bindings/python/src/read.rs                        |   3 +-
 bindings/python/tests/test_read.py                 |  20 +++
 .../datafusion/src/relation_planner.rs             |   4 +-
 .../datafusion/tests/time_travel_schema_tests.rs   |  65 ++++++++
 crates/paimon/src/spec/core_options.rs             |  55 ++++++-
 crates/paimon/src/table/mod.rs                     |   9 +-
 crates/paimon/src/table/snapshot_manager.rs        | 173 +++++++++++++++++++++
 crates/paimon/src/table/table_scan.rs              |   8 +-
 crates/paimon/src/table/time_travel.rs             | 170 ++++++++++++++++++--
 docs/src/sql.md                                    |  24 +++
 11 files changed, 555 insertions(+), 30 deletions(-)

diff --git a/bindings/c/src/table.rs b/bindings/c/src/table.rs
index 0ccf6d34..298a3c56 100644
--- a/bindings/c/src/table.rs
+++ b/bindings/c/src/table.rs
@@ -227,8 +227,9 @@ pub unsafe extern "C" fn paimon_table_free(table: *mut 
paimon_table) {
 }
 
 /// Time-travel selector option names, in the core's resolution priority order.
-const TIME_TRAVEL_SELECTORS: [&str; 4] = [
+const TIME_TRAVEL_SELECTORS: [&str; 5] = [
     "scan.timestamp-millis",
+    "scan.watermark",
     "scan.version",
     "scan.snapshot-id",
     "scan.tag-name",
@@ -300,8 +301,9 @@ pub unsafe extern "C" fn paimon_table_new_read_builder(
 
 /// Create a ReadBuilder from a Table with scan options (e.g. time-travel
 /// selectors `scan.snapshot-id` / `scan.tag-name` / `scan.timestamp-millis` /
-/// `scan.version`). At most one time-travel selector may be set. A selector 
that
-/// does not resolve to a snapshot is an error (never a silent read-of-latest).
+/// `scan.watermark` / `scan.version`). At most one time-travel selector may be
+/// set. A selector that does not resolve to a snapshot is an error (never a
+/// silent read-of-latest).
 ///
 /// # Safety
 /// `table` must be a valid pointer. `options` must be a valid pointer to
@@ -2298,11 +2300,53 @@ mod tests {
     }
 
     #[test]
-    fn unsupported_scan_option_is_rejected() {
+    fn watermark_conflicting_with_other_selector_is_rejected() {
+        unsafe {
+            let table = boxed_test_table();
+            let k1 = CString::new("scan.watermark").unwrap();
+            let v1 = CString::new("1").unwrap();
+            let k2 = CString::new("scan.snapshot-id").unwrap();
+            let v2 = CString::new("1").unwrap();
+            let opts = [opt(&k1, &v1), opt(&k2, &v2)];
+            let (code, message) = assert_rb_err_code_message(
+                paimon_table_new_read_builder_with_options(table, 
opts.as_ptr(), 2),
+            );
+            assert_eq!(code, PaimonErrorCode::InvalidInput as i32);
+            assert!(
+                message.contains("scan.watermark") && 
message.contains("scan.snapshot-id"),
+                "message should name both selectors, got: {message}"
+            );
+            paimon_table_free(table);
+        }
+    }
+
+    #[test]
+    fn unresolved_watermark_does_not_silently_read_latest() {
         unsafe {
+            // The test table commits no watermarks, so any watermark selector
+            // is unresolvable; the binding must error instead of falling back.
             let table = boxed_test_table();
             let k = CString::new("scan.watermark").unwrap();
-            let v = CString::new("0").unwrap();
+            let v = CString::new("1").unwrap();
+            let opts = [opt(&k, &v)];
+            let (code, message) = assert_rb_err_code_message(
+                paimon_table_new_read_builder_with_options(table, 
opts.as_ptr(), 1),
+            );
+            assert_eq!(code, PaimonErrorCode::InvalidInput as i32);
+            assert!(
+                message.contains("did not resolve"),
+                "message should report the selector did not resolve, got: 
{message}"
+            );
+            paimon_table_free(table);
+        }
+    }
+
+    #[test]
+    fn unsupported_scan_option_is_rejected() {
+        unsafe {
+            let table = boxed_test_table();
+            let k = CString::new("incremental-between").unwrap();
+            let v = CString::new("1,2").unwrap();
             let opts = [opt(&k, &v)];
             // Core's validate_scan_options rejects this before resolution; the
             // binding surfaces core's Unsupported code.
diff --git a/bindings/python/src/read.rs b/bindings/python/src/read.rs
index d5d63ece..db9f7208 100644
--- a/bindings/python/src/read.rs
+++ b/bindings/python/src/read.rs
@@ -31,8 +31,9 @@ use crate::error::to_py_err;
 use crate::predicate::dict_to_predicate;
 
 /// Time-travel selector option names, in the core's resolution priority order.
-const TIME_TRAVEL_SELECTORS: [&str; 4] = [
+const TIME_TRAVEL_SELECTORS: [&str; 5] = [
     "scan.timestamp-millis",
+    "scan.watermark",
     "scan.version",
     "scan.snapshot-id",
     "scan.tag-name",
diff --git a/bindings/python/tests/test_read.py 
b/bindings/python/tests/test_read.py
index 0f4600a4..49ed926d 100644
--- a/bindings/python/tests/test_read.py
+++ b/bindings/python/tests/test_read.py
@@ -610,6 +610,16 @@ def test_time_travel_unresolved_snapshot_raises():
             table.new_read_builder({"scan.snapshot-id": "999"})
 
 
+def test_time_travel_unresolved_watermark_raises():
+    with tempfile.TemporaryDirectory() as warehouse:
+        _make_two_snapshot_table(warehouse)
+        table = PaimonCatalog({"warehouse": warehouse}).get_table("tdb.t")
+        # The Rust commit path never writes watermarks, so no snapshot matches;
+        # the binding must raise instead of silently reading latest.
+        with pytest.raises(ValueError, match="did not resolve"):
+            table.new_read_builder({"scan.watermark": "1"})
+
+
 def test_unsupported_scan_option_raises_not_implemented():
     with tempfile.TemporaryDirectory() as warehouse:
         _make_two_snapshot_table(warehouse)
@@ -702,6 +712,16 @@ def test_time_travel_conflicting_selectors_raises():
         assert "scan.tag-name" in str(exc.value)
 
 
+def test_time_travel_watermark_conflicting_selector_raises():
+    with tempfile.TemporaryDirectory() as warehouse:
+        _make_two_snapshot_table(warehouse)
+        table = PaimonCatalog({"warehouse": warehouse}).get_table("tdb.t")
+        with pytest.raises(ValueError, match="Only one time-travel selector") 
as exc:
+            table.new_read_builder({"scan.watermark": "1", "scan.snapshot-id": 
"1"})
+        assert "scan.watermark" in str(exc.value)
+        assert "scan.snapshot-id" in str(exc.value)
+
+
 def test_split_serialize_produces_split_v1_binary():
     import struct
 
diff --git a/crates/integrations/datafusion/src/relation_planner.rs 
b/crates/integrations/datafusion/src/relation_planner.rs
index acbdfc2d..e51bdcd1 100644
--- a/crates/integrations/datafusion/src/relation_planner.rs
+++ b/crates/integrations/datafusion/src/relation_planner.rs
@@ -149,7 +149,7 @@ fn object_name_to_table_reference(
 /// Resolve `VERSION AS OF <expr>` into `scan.version` option.
 ///
 /// The raw value (integer or string) is passed through as-is.
-/// Resolution (tag vs snapshot id) happens at scan time in `TableScan`.
+/// Resolution (tag vs watermark vs snapshot id) happens at scan time in 
`TableScan`.
 fn resolve_version_as_of(expr: &ast::Expr) -> DFResult<HashMap<String, 
String>> {
     let version = match expr {
         ast::Expr::Value(v) => match &v.value {
@@ -163,7 +163,7 @@ fn resolve_version_as_of(expr: &ast::Expr) -> 
DFResult<HashMap<String, String>>
         },
         _ => {
             return Err(datafusion::error::DataFusionError::Plan(format!(
-                "Unsupported VERSION AS OF expression: {expr}. Expected an 
integer snapshot id or a tag name."
+                "Unsupported VERSION AS OF expression: {expr}. Expected an 
integer snapshot id, a tag name, or a quoted 'watermark-<value>'."
             )))
         }
     };
diff --git a/crates/integrations/datafusion/tests/time_travel_schema_tests.rs 
b/crates/integrations/datafusion/tests/time_travel_schema_tests.rs
index e7986653..747cd579 100644
--- a/crates/integrations/datafusion/tests/time_travel_schema_tests.rs
+++ b/crates/integrations/datafusion/tests/time_travel_schema_tests.rs
@@ -101,6 +101,19 @@ fn total_rows(batches: 
&[datafusion::arrow::record_batch::RecordBatch]) -> usize
     batches.iter().map(|b| b.num_rows()).sum()
 }
 
+fn set_snapshot_watermark(temp_dir: &TempDir, snapshot_id: i64, watermark: 
i64) {
+    let path = temp_dir
+        .path()
+        .join("default.db")
+        .join("t")
+        .join("snapshot")
+        .join(format!("snapshot-{snapshot_id}"));
+    let mut snapshot: serde_json::Value =
+        
serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap();
+    snapshot["watermark"] = serde_json::json!(watermark);
+    std::fs::write(path, serde_json::to_string(&snapshot).unwrap()).unwrap();
+}
+
 #[tokio::test]
 async fn test_version_as_of_uses_snapshot_schema() {
     let (_tmp, sql_context) = setup_evolved_table().await;
@@ -138,6 +151,58 @@ async fn test_version_as_of_uses_snapshot_schema() {
     assert_eq!(total_rows(&batches), 5);
 }
 
+#[tokio::test]
+async fn test_version_as_of_java_watermark_prefix() {
+    let temp_dir = TempDir::new().expect("Failed to create temp dir");
+    let warehouse = format!("file://{}", temp_dir.path().display());
+    let mut options = Options::new();
+    options.set(CatalogOptions::WAREHOUSE, warehouse);
+    let catalog = Arc::new(FileSystemCatalog::new(options).unwrap());
+    let sql_context = create_sql_context(catalog).await;
+
+    sql_context
+        .sql("CREATE TABLE paimon.default.t (id INT)")
+        .await
+        .unwrap()
+        .collect()
+        .await
+        .unwrap();
+    for id in 1..=3 {
+        sql_context
+            .sql(&format!("INSERT INTO paimon.default.t VALUES ({id})"))
+            .await
+            .unwrap()
+            .collect()
+            .await
+            .unwrap();
+    }
+    set_snapshot_watermark(&temp_dir, 1, 1);
+    set_snapshot_watermark(&temp_dir, 3, 10);
+
+    for (watermark, expected_rows) in [(1, 1), (9, 3), (10, 3)] {
+        let batches = sql_context
+            .sql(&format!(
+                "SELECT * FROM paimon.default.t VERSION AS OF 
'watermark-{watermark}'"
+            ))
+            .await
+            .unwrap()
+            .collect()
+            .await
+            .unwrap();
+        assert_eq!(total_rows(&batches), expected_rows);
+    }
+
+    let df = sql_context
+        .sql("SELECT * FROM paimon.default.t VERSION AS OF 'watermark-11'")
+        .await
+        .unwrap();
+    let err = df.collect().await.expect_err("watermark 11 must not match");
+    assert!(
+        err.to_string().contains("watermark[11]"),
+        "error should name the unmatched watermark: {err}"
+    );
+}
+
 #[tokio::test]
 async fn test_session_scan_version_uses_snapshot_schema() {
     let (_tmp, sql_context) = setup_evolved_table().await;
diff --git a/crates/paimon/src/spec/core_options.rs 
b/crates/paimon/src/spec/core_options.rs
index 0d2de319..f09a1e46 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -114,7 +114,7 @@ pub const SCAN_TAG_NAME_OPTION: &str = "scan.tag-name";
 const INCREMENTAL_BETWEEN_OPTION: &str = "incremental-between";
 const INCREMENTAL_BETWEEN_TIMESTAMP_OPTION: &str = 
"incremental-between-timestamp";
 const INCREMENTAL_BETWEEN_SCAN_MODE_OPTION: &str = 
"incremental-between-scan-mode";
-const SCAN_WATERMARK_OPTION: &str = "scan.watermark";
+pub const SCAN_WATERMARK_OPTION: &str = "scan.watermark";
 const SCAN_MODE_OPTION: &str = "scan.mode";
 const DEFAULT_SOURCE_SPLIT_TARGET_SIZE: i64 = 128 * 1024 * 1024;
 const DEFAULT_SOURCE_SPLIT_OPEN_FILE_COST: i64 = 4 * 1024 * 1024;
@@ -309,9 +309,12 @@ pub struct CoreOptions<'a> {
 #[derive(Debug, Clone, PartialEq, Eq)]
 pub(crate) enum TimeTravelSelector<'a> {
     TimestampMillis(i64),
+    /// `scan.watermark`: batch time travel to the earliest snapshot whose
+    /// watermark is greater than or equal to the value (millis).
+    Watermark(i64),
     /// `scan.version` (SQL `VERSION AS OF`): ambiguous by design. Resolved at
-    /// scan time as tag name (if a tag exists) → snapshot id (if parseable) →
-    /// error. `option_name` is kept for error attribution.
+    /// scan time as tag name (if a tag exists) → `watermark-<value>` → 
snapshot
+    /// id (if parseable) → error. `option_name` is kept for error attribution.
     Version {
         value: &'a str,
         option_name: &'static str,
@@ -426,7 +429,6 @@ impl<'a> CoreOptions<'a> {
             INCREMENTAL_BETWEEN_OPTION,
             INCREMENTAL_BETWEEN_TIMESTAMP_OPTION,
             INCREMENTAL_BETWEEN_SCAN_MODE_OPTION,
-            SCAN_WATERMARK_OPTION,
         ] {
             if self.options.contains_key(key) {
                 return Err(crate::Error::Unsupported {
@@ -442,6 +444,7 @@ impl<'a> CoreOptions<'a> {
                     SCAN_SNAPSHOT_ID_OPTION,
                     SCAN_TAG_NAME_OPTION,
                     SCAN_VERSION_OPTION,
+                    SCAN_WATERMARK_OPTION,
                 ]
             } else if mode.eq_ignore_ascii_case("from-timestamp") {
                 &[SCAN_TIMESTAMP_MILLIS_OPTION]
@@ -849,10 +852,13 @@ impl<'a> CoreOptions<'a> {
     }
 
     fn configured_time_travel_selectors(&self) -> Vec<&'static str> {
-        let mut selectors = Vec::with_capacity(4);
+        let mut selectors = Vec::with_capacity(5);
         if self.options.contains_key(SCAN_TIMESTAMP_MILLIS_OPTION) {
             selectors.push(SCAN_TIMESTAMP_MILLIS_OPTION);
         }
+        if self.options.contains_key(SCAN_WATERMARK_OPTION) {
+            selectors.push(SCAN_WATERMARK_OPTION);
+        }
         if self.options.contains_key(SCAN_VERSION_OPTION) {
             selectors.push(SCAN_VERSION_OPTION);
         }
@@ -883,6 +889,8 @@ impl<'a> CoreOptions<'a> {
 
         if let Some(ts) = self.parse_i64_option(SCAN_TIMESTAMP_MILLIS_OPTION)? 
{
             Ok(Some(TimeTravelSelector::TimestampMillis(ts)))
+        } else if let Some(watermark) = 
self.parse_i64_option(SCAN_WATERMARK_OPTION)? {
+            Ok(Some(TimeTravelSelector::Watermark(watermark)))
         } else if let Some(value) = 
self.options.get(SCAN_VERSION_OPTION).map(String::as_str) {
             Ok(Some(TimeTravelSelector::Version {
                 value,
@@ -2345,6 +2353,41 @@ mod tests {
         );
     }
 
+    #[test]
+    fn test_watermark_maps_to_watermark_selector() {
+        let options = HashMap::from([(SCAN_WATERMARK_OPTION.to_string(), 
"1234".to_string())]);
+        assert_eq!(
+            CoreOptions::new(&options)
+                .try_time_travel_selector()
+                .unwrap(),
+            Some(TimeTravelSelector::Watermark(1234))
+        );
+
+        // Strict numeric parsing, like scan.timestamp-millis.
+        let options = HashMap::from([(SCAN_WATERMARK_OPTION.to_string(), 
"abc".to_string())]);
+        assert!(CoreOptions::new(&options)
+            .try_time_travel_selector()
+            .is_err());
+    }
+
+    #[test]
+    fn test_watermark_conflicts_with_other_selectors() {
+        let options = HashMap::from([
+            (SCAN_WATERMARK_OPTION.to_string(), "1".to_string()),
+            (SCAN_TIMESTAMP_MILLIS_OPTION.to_string(), "2".to_string()),
+        ]);
+        let err = CoreOptions::new(&options)
+            .try_time_travel_selector()
+            .unwrap_err();
+        match err {
+            crate::Error::DataInvalid { message, .. } => {
+                assert!(message.contains(SCAN_WATERMARK_OPTION));
+                assert!(message.contains(SCAN_TIMESTAMP_MILLIS_OPTION));
+            }
+            other => panic!("unexpected: {other:?}"),
+        }
+    }
+
     #[test]
     fn test_snapshot_id_conflicts_with_version_lists_original_keys() {
         let options = HashMap::from([
@@ -2437,7 +2480,6 @@ mod tests {
             "incremental-between",
             "incremental-between-timestamp",
             "incremental-between-scan-mode",
-            "scan.watermark",
         ] {
             let options = HashMap::from([(key.to_string(), "x".to_string())]);
             let err = CoreOptions::new(&options)
@@ -2475,6 +2517,7 @@ mod tests {
             SCAN_SNAPSHOT_ID_OPTION,
             SCAN_TAG_NAME_OPTION,
             SCAN_VERSION_OPTION,
+            SCAN_WATERMARK_OPTION,
         ] {
             let options = HashMap::from([
                 ("scan.mode".to_string(), "from-snapshot".to_string()),
diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs
index 4078e3a2..b43b2b11 100644
--- a/crates/paimon/src/table/mod.rs
+++ b/crates/paimon/src/table/mod.rs
@@ -373,6 +373,7 @@ impl Table {
         let selector_changed = extra.keys().any(|k| {
             k == crate::spec::SCAN_VERSION_OPTION
                 || k == crate::spec::SCAN_TIMESTAMP_MILLIS_OPTION
+                || k == crate::spec::SCAN_WATERMARK_OPTION
                 || k == crate::spec::SCAN_SNAPSHOT_ID_OPTION
                 || k == crate::spec::SCAN_TAG_NAME_OPTION
         });
@@ -400,10 +401,10 @@ impl Table {
     ///
     /// Mirrors Java `AbstractFileStoreTable.copy(dynamicOptions)` →
     /// `tryTimeTravel`: if the merged options contain a time-travel selector
-    /// (`scan.version` / `scan.timestamp-millis` / `scan.snapshot-id` /
-    /// `scan.tag-name`) that resolves to a snapshot, the table's fields and
-    /// keys come from that snapshot's schema while the options stay the merged
-    /// ones (Java `TableSchema.copy(newOptions)`).
+    /// (`scan.version` / `scan.timestamp-millis` / `scan.watermark` /
+    /// `scan.snapshot-id` / `scan.tag-name`) that resolves to a snapshot, the
+    /// table's fields and keys come from that snapshot's schema while the
+    /// options stay the merged ones (Java `TableSchema.copy(newOptions)`).
     /// Like Java, resolution failures fall back silently to the current
     /// schema (the `if let Ok` below swallows them); an invalid selector
     /// still fails later at scan planning.
diff --git a/crates/paimon/src/table/snapshot_manager.rs 
b/crates/paimon/src/table/snapshot_manager.rs
index d7d8fbbd..cb71b31c 100644
--- a/crates/paimon/src/table/snapshot_manager.rs
+++ b/crates/paimon/src/table/snapshot_manager.rs
@@ -361,6 +361,76 @@ impl SnapshotManager {
         Ok(result)
     }
 
+    /// Returns the first snapshot whose watermark is later than or equal to 
the given
+    /// `watermark`. Snapshots without a watermark — `None`, or 
`Some(i64::MIN)`,
+    /// Flink's no-watermark sentinel — are skipped. If no such snapshot 
exists,
+    /// returns None.
+    ///
+    /// Uses binary search over the actual snapshot ID list to handle gaps from
+    /// deleted snapshots; watermarks are non-decreasing in snapshot order.
+    ///
+    /// Reference: 
[SnapshotManager.laterOrEqualWatermark](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/utils/SnapshotManager.java).
+    /// The Java binary search can retain the raw mid snapshot after walking
+    /// backwards over missing watermark metadata. This implementation only
+    /// returns a snapshot whose own effective watermark satisfies the 
predicate,
+    /// preserving the method's contract when watermark metadata is sparse.
+    pub async fn later_or_equal_watermark(
+        &self,
+        watermark: i64,
+    ) -> crate::Result<Option<Snapshot>> {
+        fn effective_watermark(snapshot: &Snapshot) -> Option<i64> {
+            snapshot.watermark().filter(|w| *w != i64::MIN)
+        }
+
+        let ids = self.list_all_ids().await?;
+        if ids.is_empty() {
+            return Ok(None);
+        }
+
+        // Find the first snapshot that carries a watermark.
+        let mut lo: usize = 0;
+        let (first, first_watermark) = loop {
+            if lo >= ids.len() {
+                return Ok(None);
+            }
+            let snapshot = self.get_snapshot(ids[lo]).await?;
+            if let Some(w) = effective_watermark(&snapshot) {
+                break (snapshot, w);
+            }
+            lo += 1;
+        };
+        if first_watermark >= watermark {
+            return Ok(Some(first));
+        }
+
+        let mut hi: usize = ids.len() - 1;
+        let mut result: Option<Snapshot> = None;
+        while lo <= hi {
+            let mid = lo + (hi - lo) / 2;
+            // A snapshot without a watermark takes the ordering position of 
the
+            // nearest earlier snapshot that carries one.
+            let mut pos = mid;
+            let mut snapshot = self.get_snapshot(ids[pos]).await?;
+            while effective_watermark(&snapshot).is_none() && pos > lo {
+                pos -= 1;
+                snapshot = self.get_snapshot(ids[pos]).await?;
+            }
+            match effective_watermark(&snapshot) {
+                // No watermark-bearing snapshot in [lo, mid]: skip the range.
+                None => lo = mid + 1,
+                Some(w) if w >= watermark => {
+                    result = Some(snapshot);
+                    if pos == 0 {
+                        break;
+                    }
+                    hi = pos - 1;
+                }
+                Some(_) => lo = mid + 1,
+            }
+        }
+        Ok(result)
+    }
+
     /// Returns the snapshot whose commit time is earlier than or equal to the 
given
     /// `timestamp_millis`. If no such snapshot exists, returns None.
     ///
@@ -447,6 +517,109 @@ mod tests {
             .build()
     }
 
+    fn test_snapshot_with_watermark(id: i64, watermark: Option<i64>) -> 
Snapshot {
+        Snapshot::builder()
+            .version(3)
+            .id(id)
+            .schema_id(0)
+            .base_manifest_list("base-list".to_string())
+            .delta_manifest_list("delta-list".to_string())
+            .commit_user("test-user".to_string())
+            .commit_identifier(0)
+            .commit_kind(CommitKind::APPEND)
+            .time_millis(1000 * id as u64)
+            .watermark(watermark)
+            .build()
+    }
+
+    async fn pick_watermark(sm: &SnapshotManager, w: i64) -> Option<i64> {
+        sm.later_or_equal_watermark(w)
+            .await
+            .unwrap()
+            .map(|s| s.id())
+    }
+
+    #[tokio::test]
+    async fn test_later_or_equal_watermark_empty() {
+        let (_, sm) = setup("memory:/test_watermark_empty").await;
+        assert!(sm.later_or_equal_watermark(100).await.unwrap().is_none());
+    }
+
+    #[tokio::test]
+    async fn test_later_or_equal_watermark_all_sentinel() {
+        // Mirrors Java SnapshotManagerTest.testLaterOrEqualWatermark: 
snapshots
+        // whose watermark is all the no-watermark sentinel never match.
+        let (_, sm) = setup("memory:/test_watermark_sentinel").await;
+        for id in 1..=3 {
+            sm.commit_snapshot(&test_snapshot_with_watermark(id, 
Some(i64::MIN)))
+                .await
+                .unwrap();
+        }
+        assert!(sm.later_or_equal_watermark(100).await.unwrap().is_none());
+    }
+
+    #[tokio::test]
+    async fn test_later_or_equal_watermark_picks_earliest_match() {
+        let (_, sm) = setup("memory:/test_watermark_earliest").await;
+        for (id, w) in [(1, 100), (2, 200), (3, 200), (4, 300)] {
+            sm.commit_snapshot(&test_snapshot_with_watermark(id, Some(w)))
+                .await
+                .unwrap();
+        }
+
+        assert_eq!(pick_watermark(&sm, 50).await, Some(1));
+        assert_eq!(pick_watermark(&sm, 100).await, Some(1));
+        assert_eq!(pick_watermark(&sm, 150).await, Some(2));
+        // Equal watermarks still select the earliest matching snapshot.
+        assert_eq!(pick_watermark(&sm, 200).await, Some(2));
+        assert_eq!(pick_watermark(&sm, 201).await, Some(4));
+        assert_eq!(pick_watermark(&sm, 300).await, Some(4));
+        // Later than every watermark: no match.
+        assert_eq!(pick_watermark(&sm, 301).await, None);
+    }
+
+    #[tokio::test]
+    async fn test_later_or_equal_watermark_skips_missing_watermarks() {
+        let (_, sm) = setup("memory:/test_watermark_skip_none").await;
+        sm.commit_snapshot(&test_snapshot_with_watermark(1, None))
+            .await
+            .unwrap();
+        sm.commit_snapshot(&test_snapshot_with_watermark(2, Some(200)))
+            .await
+            .unwrap();
+        sm.commit_snapshot(&test_snapshot_with_watermark(3, None))
+            .await
+            .unwrap();
+        sm.commit_snapshot(&test_snapshot_with_watermark(4, Some(300)))
+            .await
+            .unwrap();
+
+        assert_eq!(pick_watermark(&sm, 50).await, Some(2));
+        assert_eq!(pick_watermark(&sm, 200).await, Some(2));
+        assert_eq!(pick_watermark(&sm, 250).await, Some(4));
+        assert_eq!(pick_watermark(&sm, 301).await, None);
+    }
+
+    #[tokio::test]
+    async fn test_later_or_equal_watermark_with_id_gaps() {
+        // Deleted snapshots leave holes in the id list; selection must still 
work.
+        let (_, sm) = setup("memory:/test_watermark_gaps").await;
+        sm.commit_snapshot(&test_snapshot_with_watermark(2, Some(100)))
+            .await
+            .unwrap();
+        sm.commit_snapshot(&test_snapshot_with_watermark(5, None))
+            .await
+            .unwrap();
+        sm.commit_snapshot(&test_snapshot_with_watermark(9, Some(300)))
+            .await
+            .unwrap();
+
+        assert_eq!(pick_watermark(&sm, 100).await, Some(2));
+        assert_eq!(pick_watermark(&sm, 150).await, Some(9));
+        assert_eq!(pick_watermark(&sm, 300).await, Some(9));
+        assert_eq!(pick_watermark(&sm, 301).await, None);
+    }
+
     #[tokio::test]
     async fn test_commit_snapshot_first() {
         let (_, sm) = setup("memory:/test_commit_first").await;
diff --git a/crates/paimon/src/table/table_scan.rs 
b/crates/paimon/src/table/table_scan.rs
index 269dc415..43922255 100644
--- a/crates/paimon/src/table/table_scan.rs
+++ b/crates/paimon/src/table/table_scan.rs
@@ -1013,13 +1013,15 @@ impl<'a> PaimonTableScan<'a> {
     /// Plan the full scan: resolve snapshot (via options or latest), then 
read manifests and build DataSplits.
     ///
     /// Time travel is resolved from table options:
-    /// - only one of `scan.version`, `scan.timestamp-millis`,
+    /// - only one of `scan.version`, `scan.timestamp-millis`, 
`scan.watermark`,
     ///   `scan.snapshot-id`, `scan.tag-name` may be set
-    /// - `scan.version` → tag name (if exists) → snapshot id (if parseable) →
-    ///   error (ambiguous by design, like SQL `VERSION AS OF`)
+    /// - `scan.version` → tag name (if exists) → `watermark-<value>` → 
snapshot
+    ///   id (if parseable) → error (ambiguous by design, like SQL `VERSION AS 
OF`)
     /// - `scan.snapshot-id` → snapshot id only (never a tag lookup)
     /// - `scan.tag-name` → tag name only (never parsed as a snapshot id)
     /// - `scan.timestamp-millis` → find the latest snapshot <= that timestamp
+    /// - `scan.watermark` → find the earliest snapshot with watermark >= that
+    ///   value (snapshots without a watermark are skipped)
     /// - otherwise → read the latest snapshot
     ///
     /// Reference: 
[TimeTravelUtil.tryTravelToSnapshot](https://github.com/apache/paimon/blob/master/paimon-core/src/main/java/org/apache/paimon/table/source/snapshot/TimeTravelUtil.java)
diff --git a/crates/paimon/src/table/time_travel.rs 
b/crates/paimon/src/table/time_travel.rs
index 3a389e33..cfb80a3e 100644
--- a/crates/paimon/src/table/time_travel.rs
+++ b/crates/paimon/src/table/time_travel.rs
@@ -22,6 +22,8 @@ use crate::table::{SnapshotManager, Table, TagManager};
 use crate::Error;
 use std::collections::HashMap;
 
+const WATERMARK_PREFIX: &str = "watermark-";
+
 /// Resolve the snapshot selected by the time-travel options, if any.
 ///
 /// Returns `Ok(None)` when no time-travel selector is configured. Returns an
@@ -45,13 +47,27 @@ pub(crate) async fn travel_to_snapshot(
                 }),
             }
         }
+        Some(TimeTravelSelector::Watermark(w)) => {
+            resolve_watermark(snapshot_manager, w).await.map(Some)
+        }
         Some(TimeTravelSelector::Version {
             value: v,
             option_name,
         }) => {
-            // `scan.version` is ambiguous by design: tag first, then snapshot 
id.
+            // Match Java TimeTravelUtil.adaptScanVersion: tag first, then the
+            // `watermark-<value>` prefix, then snapshot id.
             if tag_manager.tag_exists(v).await? {
                 resolve_tag(tag_manager, v).await.map(Some)
+            } else if let Some(raw_watermark) = 
v.strip_prefix(WATERMARK_PREFIX) {
+                let watermark = raw_watermark
+                    .parse::<i64>()
+                    .map_err(|e| Error::DataInvalid {
+                        message: format!("{option_name} '{v}' has an invalid 
watermark value."),
+                        source: Some(Box::new(e)),
+                    })?;
+                resolve_watermark(snapshot_manager, watermark)
+                    .await
+                    .map(Some)
             } else if let Ok(id) = v.parse::<i64>() {
                 snapshot_manager.get_snapshot(id).await.map(Some)
             } else {
@@ -92,6 +108,22 @@ pub(crate) async fn travel_to_snapshot(
     }
 }
 
+async fn resolve_watermark(
+    snapshot_manager: &SnapshotManager,
+    watermark: i64,
+) -> crate::Result<Snapshot> {
+    match snapshot_manager.later_or_equal_watermark(watermark).await? {
+        Some(snapshot) => Ok(snapshot),
+        // Mirrors Java StaticFromWatermarkStartingScanner's error.
+        None => Err(Error::DataInvalid {
+            message: format!(
+                "There is currently no snapshot later than or equal to 
watermark[{watermark}]"
+            ),
+            source: None,
+        }),
+    }
+}
+
 /// Resolve the snapshot a read should use, including the latest-snapshot 
fallback.
 ///
 /// Reuses a snapshot cached by [`Table::copy_with_time_travel`] so every read 
path
@@ -138,7 +170,7 @@ async fn resolve_tag(tag_manager: &TagManager, name: &str) 
-> crate::Result<Snap
 mod tests {
     use crate::catalog::Identifier;
     use crate::io::{FileIO, FileIOBuilder};
-    use crate::spec::{DataType, IntType, Schema, TableSchema};
+    use crate::spec::{CommitKind, DataType, IntType, Schema, Snapshot, 
TableSchema};
     use crate::table::{SnapshotManager, Table, TableCommit, TableWrite, 
TagManager};
     use arrow_array::{Int32Array, RecordBatch};
     use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema 
as ArrowSchema};
@@ -255,6 +287,39 @@ mod tests {
         make_table(file_io, table_path, schema_v1())
     }
 
+    /// Table whose snapshots carry watermarks, committed directly through
+    /// `SnapshotManager` (the Rust commit path never writes watermarks):
+    /// snapshot 1 (watermark 100), snapshot 2 (no watermark), snapshot 3
+    /// (watermark 300), all on schema 0.
+    async fn setup_watermark_table() -> (FileIO, String) {
+        let file_io = FileIOBuilder::new("memory").build().unwrap();
+        let table_path = "memory:/watermark_table";
+        for dir in ["snapshot", "manifest"] {
+            file_io
+                .mkdirs(&format!("{table_path}/{dir}/"))
+                .await
+                .unwrap();
+        }
+        write_schema_file(&file_io, table_path, &schema_v0()).await;
+        let sm = SnapshotManager::new(file_io.clone(), table_path.to_string());
+        for (id, watermark) in [(1, Some(100)), (2, None), (3, Some(300))] {
+            let snapshot = Snapshot::builder()
+                .version(3)
+                .id(id)
+                .schema_id(0)
+                .base_manifest_list(format!("base-list-{id}"))
+                .delta_manifest_list(format!("delta-list-{id}"))
+                .commit_user("test-user".to_string())
+                .commit_identifier(0)
+                .commit_kind(CommitKind::APPEND)
+                .time_millis(1000 * id as u64)
+                .watermark(watermark)
+                .build();
+            sm.commit_snapshot(&snapshot).await.unwrap();
+        }
+        (file_io, table_path.to_string())
+    }
+
     fn options(pairs: &[(&str, &str)]) -> HashMap<String, String> {
         pairs
             .iter()
@@ -503,18 +568,105 @@ mod tests {
     }
 
     #[tokio::test]
-    async fn test_copy_with_time_travel_rejects_unsupported_scan_option() {
-        let (file_io, table_path) = setup_evolved_table().await;
-        let table = latest_table(&file_io, &table_path);
-        let err = table
-            .copy_with_time_travel(options(&[("scan.watermark", "5")]))
+    async fn test_copy_with_time_travel_resolves_watermark() {
+        let (file_io, table_path) = setup_watermark_table().await;
+        let table = make_table(&file_io, &table_path, schema_v0());
+
+        // Exact match on snapshot 1; snapshot 2 carries no watermark and is 
skipped.
+        let traveled = table
+            .copy_with_time_travel(options(&[("scan.watermark", "100")]))
             .await
-            .expect_err("unsupported scan option must fail");
+            .unwrap();
+        assert_eq!(traveled.travel_snapshot().map(|s| s.id()), Some(1));
+
+        // Between watermarks: the earliest snapshot with watermark >= the 
value.
+        let traveled = table
+            .copy_with_time_travel(options(&[("scan.watermark", "150")]))
+            .await
+            .unwrap();
+        assert_eq!(traveled.travel_snapshot().map(|s| s.id()), Some(3));
+        assert!(traveled.has_resolved_travel_snapshot());
+    }
+
+    #[tokio::test]
+    async fn test_scan_version_resolves_java_watermark_prefix_after_tag() {
+        let (file_io, table_path) = setup_watermark_table().await;
+        let table = make_table(&file_io, &table_path, schema_v0());
+
+        let traveled = table
+            .copy_with_time_travel(options(&[("scan.version", 
"watermark-150")]))
+            .await
+            .unwrap();
+        assert_eq!(traveled.travel_snapshot().map(|s| s.id()), Some(3));
+
+        // Java resolves an existing tag before interpreting the watermark 
prefix.
+        let sm = SnapshotManager::new(file_io.clone(), table_path.clone());
+        let snapshot1 = sm.get_snapshot(1).await.unwrap();
+        let tm = TagManager::new(file_io.clone(), table_path.clone());
+        tm.create("watermark-150", &snapshot1).await.unwrap();
+        let tagged = table
+            .copy_with_time_travel(options(&[("scan.version", 
"watermark-150")]))
+            .await
+            .unwrap();
+        assert_eq!(tagged.travel_snapshot().map(|s| s.id()), Some(1));
+
+        let err =
+            super::travel_to_snapshot(&sm, &tm, &options(&[("scan.version", 
"watermark-invalid")]))
+                .await
+                .expect_err("invalid watermark version must fail");
         assert!(
-            matches!(err, crate::Error::Unsupported { message } if 
message.contains("scan.watermark"))
+            matches!(err, crate::Error::DataInvalid { ref message, .. }
+                if message.contains("invalid watermark value")),
+            "expected watermark parse error, got {err:?}"
         );
     }
 
+    #[tokio::test]
+    async fn test_watermark_without_matching_snapshot_fails_at_scan() {
+        let (file_io, table_path) = setup_watermark_table().await;
+        let table = make_table(&file_io, &table_path, schema_v0());
+
+        // Like Java tryTravelToSnapshot, resolution failure falls back 
silently...
+        let unresolved = table
+            .copy_with_time_travel(options(&[("scan.watermark", "301")]))
+            .await
+            .unwrap();
+        assert!(!unresolved.has_resolved_travel_snapshot());
+
+        // ...and the error surfaces at scan planning, naming the watermark.
+        let err = unresolved
+            .new_read_builder()
+            .new_scan()
+            .plan()
+            .await
+            .expect_err("scan with unresolvable watermark must fail");
+        assert!(
+            matches!(err, crate::Error::DataInvalid { ref message, .. }
+                if message.contains("watermark[301]")),
+            "expected watermark error, got {err:?}"
+        );
+    }
+
+    #[tokio::test]
+    async fn test_watermark_selector_change_invalidates_resolved_snapshot() {
+        let (file_io, table_path) = setup_watermark_table().await;
+        let table = make_table(&file_io, &table_path, schema_v0());
+
+        let traveled = table
+            .copy_with_time_travel(options(&[("scan.watermark", "100")]))
+            .await
+            .unwrap();
+        assert_eq!(traveled.travel_snapshot().map(|s| s.id()), Some(1));
+
+        // Merging unrelated options keeps the resolved snapshot.
+        let recopied = traveled.copy_with_options(options(&[("k", "v")]));
+        assert_eq!(recopied.travel_snapshot().map(|s| s.id()), Some(1));
+
+        // Changing the watermark invalidates the cached resolution.
+        let changed = traveled.copy_with_options(options(&[("scan.watermark", 
"150")]));
+        assert!(changed.travel_snapshot().is_none());
+    }
+
     #[tokio::test]
     async fn test_has_resolved_travel_snapshot_reflects_resolution() {
         let (file_io, table_path) = setup_evolved_table().await;
diff --git a/docs/src/sql.md b/docs/src/sql.md
index be0ea7e8..d02217da 100644
--- a/docs/src/sql.md
+++ b/docs/src/sql.md
@@ -1539,6 +1539,30 @@ SELECT * FROM paimon.default.my_table TIMESTAMP AS OF 
'2024-01-01 00:00:00';
 
 This finds the latest snapshot whose commit time is less than or equal to the 
given timestamp. The timestamp is interpreted in the local timezone.
 
+### By Watermark
+
+Use `VERSION AS OF 'watermark-<value>'` syntax:
+
+```sql
+SELECT * FROM paimon.default.my_table
+VERSION AS OF 'watermark-1704067200000';
+```
+
+This resolves the tag first if a tag with that exact name exists. Otherwise,
+the suffix is parsed as a watermark in milliseconds. The session-scoped dynamic
+option `scan.watermark` is also available:
+
+```sql
+SET 'paimon.scan.watermark' = '1704067200000';
+SELECT * FROM paimon.default.my_table;
+RESET 'paimon.scan.watermark';
+```
+
+This reads the earliest snapshot whose watermark is greater than or equal to 
the
+given value (snapshots without a watermark are skipped). It is mutually
+exclusive with the other time-travel selectors. If no matching snapshot exists,
+scan planning fails.
+
 ## Dynamic Options (SET / RESET)
 
 Use `SET` to configure session-scoped Paimon dynamic options that apply to 
subsequent table loads:

Reply via email to