diff --git a/fluss-rust/bindings/cpp/include/fluss.hpp b/fluss-rust/bindings/cpp/include/fluss.hpp index 261c1d972da..f8b37efdf09 100644 --- a/fluss-rust/bindings/cpp/include/fluss.hpp +++ b/fluss-rust/bindings/cpp/include/fluss.hpp @@ -1548,6 +1548,8 @@ struct Configuration { int32_t scanner_log_fetch_wait_max_time_ms{500}; // Maximum bytes per fetch response per bucket for LogScanner (1 MB) int32_t scanner_log_fetch_max_bytes_for_bucket{1024 * 1024}; + // Maximum bytes per fetch response per bucket for a full KV scan (4 MB) + int32_t scanner_kv_fetch_max_bytes{4 * 1024 * 1024}; int64_t writer_batch_timeout_ms{100}; // Whether to enable idempotent writes bool writer_enable_idempotence{true}; diff --git a/fluss-rust/bindings/cpp/src/ffi_converter.hpp b/fluss-rust/bindings/cpp/src/ffi_converter.hpp index ae23003e738..ac5e5a6f6cd 100644 --- a/fluss-rust/bindings/cpp/src/ffi_converter.hpp +++ b/fluss-rust/bindings/cpp/src/ffi_converter.hpp @@ -228,6 +228,7 @@ inline ffi::FfiConfig to_ffi_config(const Configuration& config) { ffi_config.scanner_log_fetch_min_bytes = config.scanner_log_fetch_min_bytes; ffi_config.scanner_log_fetch_wait_max_time_ms = config.scanner_log_fetch_wait_max_time_ms; ffi_config.scanner_log_fetch_max_bytes_for_bucket = config.scanner_log_fetch_max_bytes_for_bucket; + ffi_config.scanner_kv_fetch_max_bytes = config.scanner_kv_fetch_max_bytes; ffi_config.writer_batch_timeout_ms = config.writer_batch_timeout_ms; ffi_config.writer_enable_idempotence = config.writer_enable_idempotence; ffi_config.writer_max_inflight_requests_per_bucket = diff --git a/fluss-rust/bindings/cpp/src/lib.rs b/fluss-rust/bindings/cpp/src/lib.rs index 7d1224b28c7..75d680a0b51 100644 --- a/fluss-rust/bindings/cpp/src/lib.rs +++ b/fluss-rust/bindings/cpp/src/lib.rs @@ -64,6 +64,7 @@ mod ffi { scanner_log_fetch_min_bytes: i32, scanner_log_fetch_wait_max_time_ms: i32, scanner_log_fetch_max_bytes_for_bucket: i32, + scanner_kv_fetch_max_bytes: i32, writer_batch_timeout_ms: i64, writer_enable_idempotence: bool, writer_max_inflight_requests_per_bucket: usize, @@ -1143,6 +1144,7 @@ fn new_connection(config: &ffi::FfiConfig) -> ffi::FfiPtrResult { scanner_log_fetch_min_bytes: config.scanner_log_fetch_min_bytes, scanner_log_fetch_wait_max_time_ms: config.scanner_log_fetch_wait_max_time_ms, scanner_log_fetch_max_bytes_for_bucket: config.scanner_log_fetch_max_bytes_for_bucket, + scanner_kv_fetch_max_bytes: config.scanner_kv_fetch_max_bytes, writer_enable_idempotence: config.writer_enable_idempotence, writer_max_inflight_requests_per_bucket: config.writer_max_inflight_requests_per_bucket, writer_buffer_memory_size: config.writer_buffer_memory_size, diff --git a/fluss-rust/bindings/elixir/lib/fluss/error.ex b/fluss-rust/bindings/elixir/lib/fluss/error.ex index b2f1ddc3b69..11096429416 100644 --- a/fluss-rust/bindings/elixir/lib/fluss/error.ex +++ b/fluss-rust/bindings/elixir/lib/fluss/error.ex @@ -97,6 +97,10 @@ defmodule Fluss.Error do | :invalid_alter_table_exception | :deletion_disabled_exception | :storage_backpressure_exception + | :scanner_expired_exception + | :unknown_scanner_id_exception + | :invalid_scan_request_exception + | :too_many_scanners | :client_error @type t :: %__MODULE__{code: code(), error_code: integer(), message: String.t()} diff --git a/fluss-rust/bindings/elixir/native/fluss_nif/src/atoms.rs b/fluss-rust/bindings/elixir/native/fluss_nif/src/atoms.rs index 7a7116422b2..f2ba90e75c1 100644 --- a/fluss-rust/bindings/elixir/native/fluss_nif/src/atoms.rs +++ b/fluss-rust/bindings/elixir/native/fluss_nif/src/atoms.rs @@ -101,6 +101,10 @@ rustler::atoms! { invalid_alter_table_exception, deletion_disabled_exception, storage_backpressure_exception, + scanner_expired_exception, + unknown_scanner_id_exception, + invalid_scan_request_exception, + too_many_scanners, client_error, } @@ -214,6 +218,10 @@ fn api_error_atom(code: i32) -> Atom { FlussError::InvalidAlterTableException => invalid_alter_table_exception(), FlussError::DeletionDisabledException => deletion_disabled_exception(), FlussError::StorageBackpressureException => storage_backpressure_exception(), + FlussError::ScannerExpired => scanner_expired_exception(), + FlussError::UnknownScannerId => unknown_scanner_id_exception(), + FlussError::InvalidScanRequest => invalid_scan_request_exception(), + FlussError::TooManyScanners => too_many_scanners(), } } diff --git a/fluss-rust/crates/fluss/src/client/metadata.rs b/fluss-rust/crates/fluss/src/client/metadata.rs index 6041d5f12ce..1c8265d326d 100644 --- a/fluss-rust/crates/fluss/src/client/metadata.rs +++ b/fluss-rust/crates/fluss/src/client/metadata.rs @@ -502,10 +502,17 @@ impl Metadata { #[cfg(test)] impl Metadata { pub(crate) fn new_for_test(cluster: Arc) -> Self { + Self::new_for_test_with_connections(cluster, Arc::new(RpcClient::new())) + } + + pub(crate) fn new_for_test_with_connections( + cluster: Arc, + connections: Arc, + ) -> Self { let (cluster_version_tx, _) = watch::channel(0); Metadata { cluster: RwLock::new(cluster), - connections: Arc::new(RpcClient::new()), + connections, bootstrap: Arc::from(""), cluster_version_tx, } diff --git a/fluss-rust/crates/fluss/src/client/table/batch_scanner.rs b/fluss-rust/crates/fluss/src/client/table/batch_scanner.rs index 1519ccc3e8a..6aee7be821a 100644 --- a/fluss-rust/crates/fluss/src/client/table/batch_scanner.rs +++ b/fluss-rust/crates/fluss/src/client/table/batch_scanner.rs @@ -287,9 +287,10 @@ async fn decode_kv_batch( table_info: &TableInfo, schema_getter: &ClientSchemaGetter, projected_fields: Option<&[usize]>, - raw: Vec, + raw: impl Into, limit: usize, ) -> Result { + let raw: Bytes = raw.into(); // No records: return an empty (projected) batch. if raw.is_empty() { return empty_record_batch(table_info.get_row_type(), projected_fields); @@ -306,7 +307,7 @@ async fn decode_kv_batch( source: None, })?; - let batch = ValueRecordBatch::new(Bytes::from(raw)); + let batch = ValueRecordBatch::new(raw); let ranges = batch.value_ranges()?; // Collect the distinct schema ids present, then build one decoder per id @@ -455,6 +456,9 @@ fn project_batch( } } +mod kv_batch_scanner; +pub use kv_batch_scanner::{KvBatchReadOutcome, KvBatchScanner, KvSnapshotScanner}; + #[cfg(test)] mod tests { use super::*; diff --git a/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner.rs b/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner.rs new file mode 100644 index 00000000000..6bf4894654f --- /dev/null +++ b/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner.rs @@ -0,0 +1,726 @@ +// 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. + +use super::decode_kv_batch; +use crate::client::ClientSchemaGetter; +use crate::client::metadata::Metadata; +use crate::error::{ApiError, Error, FlussError, Result}; +use crate::metadata::{TableBucket, TableInfo}; +use crate::proto::{PbScanReqForBucket, ScanKvResponse}; +use crate::record::ScanBatch; +use crate::rpc::message::ScanKvRequest; +use crate::rpc::{RpcClient, ServerConnection}; +use bytes::Bytes; +use log::debug; +use std::collections::{HashSet, VecDeque}; +use std::sync::Arc; +use std::time::Duration; +use tokio::task::JoinHandle; +use tokio::time::Instant; + +const DEFAULT_KV_POLL_TIMEOUT: Duration = Duration::from_millis(500); + +/// Streams every live row of a single primary-key bucket from the tablet +/// server's KV state via a sequence of `ScanKv` RPCs. The scan has snapshot +/// isolation: rows reflect the KV state at the moment the server opened its +/// snapshot; concurrent writes after that point are invisible. +/// +/// At most one RPC is in flight. The RPC runs in an owned Tokio task so a caller +/// timeout or cancellation does not discard a response after the server may +/// have advanced its RocksDB iterator. Once a response indicates more data, the +/// next continuation is started immediately so its network round-trip overlaps +/// decoding and consuming the current batch. +/// +/// Empty intermediate responses are skipped internally, so a returned batch +/// always carries at least one row. +pub struct KvBatchScanner { + bucket: TableBucket, + rpc_client: Arc, + metadata: Arc, + table_info: TableInfo, + schema_getter: Arc, + projected_fields: Option>, + batch_size_bytes: i32, + /// Leader connection, resolved lazily on the first `next_batch`. + connection: Option, + /// Server-assigned session id; `None` until the first response. + scanner_id: Option>, + /// Monotonic sequence number: 0 for the open request, incremented per + /// continuation (matches the server's in-order delivery validation). + call_seq_id: i32, + /// Number of `TooManyScanners` open retries already attempted. + open_retries: u32, + /// Earliest instant at which a `TooManyScanners` open retry may be sent. + /// Stored in scanner state so caller timeouts or cancellation do not shorten + /// the exponential backoff. + open_retry_at: Option, + /// Log high-watermark captured when the server opened the snapshot. + log_offset: Option, + /// The single open or continuation RPC currently executing. + in_flight: Option>>, + /// Raw records from the last completed RPC. Retained until decoding succeeds + /// so cancellation during an asynchronous schema lookup cannot lose rows. + pending_records: Option, + drained: bool, + closed: bool, +} + +/// Result of an owned ScanKV task. The open task also returns the resolved +/// connection so subsequent continuations stay on the server that owns the +/// scanner session. +struct ScanKvRpcResult { + connection: ServerConnection, + response: ScanKvResponse, +} + +/// Outcome of a KV batch read with a caller-supplied timeout. +#[derive(Debug)] +pub enum KvBatchReadOutcome { + /// A batch is available. + Batch(ScanBatch), + /// The current RPC or record decoding did not complete before the timeout. + /// The scanner remains valid and the caller may retry. + TimedOut, + /// The bucket has been fully drained. + Finished, +} + +impl KvBatchScanner { + const MAX_OPEN_RETRIES: u32 = 3; + const BASE_RETRY_DELAY_MS: u64 = 100; + + pub(crate) fn new( + rpc_client: Arc, + metadata: Arc, + table_info: TableInfo, + schema_getter: Arc, + projected_fields: Option>, + bucket: TableBucket, + batch_size_bytes: i32, + ) -> Self { + Self { + bucket, + rpc_client, + metadata, + table_info, + schema_getter, + projected_fields, + batch_size_bytes: batch_size_bytes.max(1), + connection: None, + scanner_id: None, + call_seq_id: 0, + open_retries: 0, + open_retry_at: None, + log_offset: None, + in_flight: None, + pending_records: None, + drained: false, + closed: false, + } + } + + /// The bucket scanned by this `KvBatchScanner`. + pub fn bucket(&self) -> &TableBucket { + &self.bucket + } + + /// The log high-watermark captured when the server opened the KV snapshot, + /// available after the first `next_batch`. It marks the log offset from + /// which a changelog tail can resume to read rows written after the + /// snapshot (snapshot + changelog handoff). + pub fn snapshot_log_offset(&self) -> Option { + self.log_offset + } + + /// Returns the next decoded batch, waiting until data arrives or the scan is + /// drained. + /// + /// Transient continuation failures are not retried because the server may + /// already have advanced the scanner cursor. An error leaves this scanner + /// spent; create a new scanner to restart the bucket from the beginning. + pub async fn next_batch(&mut self) -> Result> { + loop { + match self + .next_batch_with_timeout(DEFAULT_KV_POLL_TIMEOUT) + .await? + { + KvBatchReadOutcome::Batch(batch) => return Ok(Some(batch)), + KvBatchReadOutcome::TimedOut => continue, + KvBatchReadOutcome::Finished => return Ok(None), + } + } + } + + /// Returns the next batch while waiting for at most `timeout`. + /// + /// A timeout never cancels or resends the current `ScanKv` request. The + /// owned RPC task remains stored in the scanner, and a later call continues + /// waiting for that exact request and `call_seq_id`. + pub async fn next_batch_with_timeout( + &mut self, + timeout: Duration, + ) -> Result { + let start = Instant::now(); + + loop { + if self.is_fully_drained() { + return Ok(KvBatchReadOutcome::Finished); + } + if self.closed { + return Err(Error::UnexpectedError { + message: format!( + "KvBatchScanner for bucket {} was closed before it was fully drained", + self.bucket + ), + source: None, + }); + } + + if let Some(raw) = self.pending_records.clone() { + let Some(remaining) = remaining_timeout(start, timeout) else { + return Ok(KvBatchReadOutcome::TimedOut); + }; + let decoded = tokio::time::timeout( + remaining, + decode_kv_batch( + &self.table_info, + &self.schema_getter, + self.projected_fields.as_deref(), + raw, + usize::MAX, + ), + ) + .await; + let batch = match decoded { + Ok(Ok(batch)) => batch, + Ok(Err(error)) => { + self.terminate(); + return Err(error); + } + Err(_) => return Ok(KvBatchReadOutcome::TimedOut), + }; + + // Clear only after decoding succeeds. If the decode future is + // cancelled while fetching an older schema, the original bytes + // remain available for the next call. + self.pending_records = None; + return Ok(KvBatchReadOutcome::Batch(ScanBatch::new( + self.bucket.clone(), + batch, + 0, + ))); + } + + if let Some(retry_at) = self.open_retry_at { + let now = Instant::now(); + if now < retry_at { + let Some(remaining) = remaining_timeout(start, timeout) else { + return Ok(KvBatchReadOutcome::TimedOut); + }; + let retry_delay = retry_at - now; + if retry_delay >= remaining { + tokio::time::sleep(remaining).await; + return Ok(KvBatchReadOutcome::TimedOut); + } + tokio::time::sleep(retry_delay).await; + } + self.open_retry_at = None; + } + + if self.in_flight.is_none() { + if self.scanner_id.is_none() { + self.start_open_request(); + } else { + self.start_continuation_or_terminate()?; + } + } + + let Some(remaining) = remaining_timeout(start, timeout) else { + return Ok(KvBatchReadOutcome::TimedOut); + }; + let task_result = { + let task = self + .in_flight + .as_mut() + .expect("ScanKV request must be in flight"); + wait_for_task(task, remaining).await + }; + + let joined = match task_result { + Some(joined) => joined, + None => return Ok(KvBatchReadOutcome::TimedOut), + }; + self.in_flight = None; + + let rpc_result = match joined { + Ok(Ok(result)) => result, + Ok(Err(error)) => { + self.terminate(); + return Err(error); + } + Err(error) => { + self.terminate(); + return Err(Error::UnexpectedError { + message: format!("ScanKV RPC task failed: {error}"), + source: Some(Box::new(error)), + }); + } + }; + self.connection = Some(rpc_result.connection); + let mut response = rpc_result.response; + + if let Some(code) = response.error_code + && code != FlussError::None.code() + { + let retry_delay = self + .handle_error_response(code, response.error_message.take()) + .await?; + if let Some(delay) = retry_delay { + self.open_retry_at = Some(Instant::now() + delay); + continue; + } + } + + let response_has_scanner_id = response.scanner_id.is_some(); + if let Some(id) = response.scanner_id.take() { + self.scanner_id = Some(id); + } + if self.log_offset.is_none() { + self.log_offset = response.log_offset; + } + + let Some(has_more_results) = response.has_more_results else { + self.terminate(); + return Err(Error::UnexpectedError { + message: "ScanKV response did not include has_more_results".to_string(), + source: None, + }); + }; + if has_more_results && !response_has_scanner_id { + self.terminate(); + return Err(Error::UnexpectedError { + message: "ScanKV response reported more results without a scanner id" + .to_string(), + source: None, + }); + } + + self.pending_records = response.records.take().filter(|raw| !raw.is_empty()); + + if has_more_results { + // Pipeline the next continuation before decoding or returning + // the current records, matching the Java scanner. + self.start_continuation_or_terminate()?; + } else { + // A terminal response means the server has already closed the + // session; no explicit close request is needed. + self.drained = true; + } + } + } + + /// Drains the scanner into all of its batches. + pub async fn collect_all_batches(&mut self) -> Result> { + let mut batches = Vec::new(); + while let Some(batch) = self.next_batch().await? { + batches.push(batch); + } + Ok(batches) + } + + /// Best-effort close of the server-side scanner session. Only meaningful for + /// a scanner abandoned mid-scan; a drained scanner is already closed by the + /// server, and any failure here is reclaimed by the server-side session TTL. + /// A scanner closed before it is drained remains incomplete; subsequent + /// reads return an error rather than reporting normal end-of-scan. + pub async fn close(&mut self) -> Result<()> { + if self.closed || self.is_fully_drained() { + return Ok(()); + } + self.closed = true; + // A terminal response may already have arrived while its records are + // still waiting to be decoded. Discarding those records is an + // incomplete close, not a successfully drained scan. + self.drained = false; + if let Some(task) = self.in_flight.take() { + task.abort(); + } + self.pending_records = None; + // Dispatch the close in an owned task. The RPC is best effort, and this + // keeps close cancellation-safe if the caller drops this future. + self.send_best_effort_close(); + Ok(()) + } + + fn start_continuation_or_terminate(&mut self) -> Result<()> { + match self.start_continuation() { + Ok(()) => Ok(()), + Err(error) => { + self.terminate(); + Err(error) + } + } + } + + fn start_open_request(&mut self) { + let bucket_req = PbScanReqForBucket { + table_id: self.bucket.table_id(), + partition_id: self.bucket.partition_id(), + bucket_id: self.bucket.bucket_id(), + limit: None, + }; + // call_seq_id stays 0 for the open request (including open retries). + self.call_seq_id = 0; + let request = ScanKvRequest::new( + None, + Some(bucket_req), + Some(self.call_seq_id), + Some(self.batch_size_bytes), + None, + ); + + let metadata = Arc::clone(&self.metadata); + let rpc_client = Arc::clone(&self.rpc_client); + let table_path = self.table_info.table_path.clone(); + let bucket = self.bucket.clone(); + self.in_flight = Some(tokio::spawn(async move { + let leader = metadata + .leader_for(&table_path, &bucket) + .await? + .ok_or_else(|| { + Error::leader_not_available(format!( + "No leader found for table bucket: {bucket}" + )) + })?; + let connection = rpc_client.get_connection(&leader).await?; + let response = connection.request(request).await?; + Ok(ScanKvRpcResult { + connection, + response, + }) + })); + } + + fn start_continuation(&mut self) -> Result<()> { + if self.in_flight.is_some() { + return Err(Error::UnexpectedError { + message: "KvBatchScanner attempted to start concurrent ScanKV requests".to_string(), + source: None, + }); + } + let connection = self + .connection + .clone() + .ok_or_else(|| Error::UnexpectedError { + message: "KvBatchScanner continuation issued without an open connection" + .to_string(), + source: None, + })?; + let scanner_id = self + .scanner_id + .clone() + .ok_or_else(|| Error::UnexpectedError { + message: "KvBatchScanner continuation issued without a scanner id".to_string(), + source: None, + })?; + self.call_seq_id += 1; + let request = ScanKvRequest::new( + Some(scanner_id), + None, + Some(self.call_seq_id), + Some(self.batch_size_bytes), + None, + ); + self.in_flight = Some(tokio::spawn(async move { + let response = connection.request(request).await?; + Ok(ScanKvRpcResult { + connection, + response, + }) + })); + Ok(()) + } + + /// Handles an error response. `Ok(Some(delay))` means the caller should + /// retry opening after the delay; all other errors are terminal. + async fn handle_error_response( + &mut self, + code: i32, + message: Option, + ) -> Result> { + let error = FlussError::for_code(code); + let api_error = ApiError { + code, + message: message.unwrap_or_else(|| error.message().to_string()), + }; + + match error { + // Retry the open with exponential backoff — only before a session + // was established, so no rows can be skipped or repeated. + FlussError::TooManyScanners + if self.scanner_id.is_none() && self.open_retries < Self::MAX_OPEN_RETRIES => + { + let delay_ms = Self::BASE_RETRY_DELAY_MS * (1u64 << self.open_retries); + self.open_retries += 1; + Ok(Some(Duration::from_millis(delay_ms))) + } + // Stale leader: refresh metadata so a fresh scan resolves the new + // leader. Auto-restarting here would silently swap the RocksDB + // snapshot and break snapshot isolation. + FlussError::NotLeaderOrFollower => { + self.terminate(); + self.refresh_bucket_metadata().await; + Err(Error::FlussAPIError { api_error }) + } + // The server-side session is already gone; skip the close request. + FlussError::ScannerExpired | FlussError::UnknownScannerId => { + self.scanner_id = None; + self.closed = true; + Err(Error::FlussAPIError { api_error }) + } + _ => { + self.terminate(); + Err(Error::FlussAPIError { api_error }) + } + } + } + + async fn refresh_bucket_metadata(&self) { + let partition_ids = metadata_refresh_partition_ids(&self.bucket); + let result = if partition_ids.is_empty() { + self.metadata + .update_table_metadata(&self.table_info.table_path) + .await + } else { + self.metadata + .update_tables_metadata( + &HashSet::from([&self.table_info.table_path]), + &HashSet::new(), + partition_ids, + ) + .await + }; + if let Err(error) = result { + debug!( + "Failed to refresh metadata after NotLeaderOrFollower for bucket {}: {}", + self.bucket, error + ); + } + } + + /// Marks the scanner closed, aborts the local RPC task and sends a + /// best-effort close request. The server-side TTL is the final fallback. + fn terminate(&mut self) { + if self.closed { + return; + } + let fully_drained = self.is_fully_drained(); + self.closed = true; + if let Some(task) = self.in_flight.take() { + task.abort(); + } + self.open_retry_at = None; + self.pending_records = None; + if !fully_drained { + self.drained = false; + } + self.send_best_effort_close(); + } + + fn send_best_effort_close(&self) { + if self.drained { + return; + } + let (Some(connection), Some(scanner_id)) = + (self.connection.clone(), self.scanner_id.clone()) + else { + return; + }; + let batch_size_bytes = self.batch_size_bytes; + if let Ok(runtime) = tokio::runtime::Handle::try_current() { + runtime.spawn(async move { + let request = ScanKvRequest::new( + Some(scanner_id), + None, + None, + Some(batch_size_bytes), + Some(true), + ); + let _ = connection.request(request).await; + }); + } + } + + fn is_fully_drained(&self) -> bool { + self.drained && self.pending_records.is_none() + } +} + +impl Drop for KvBatchScanner { + fn drop(&mut self) { + self.terminate(); + } +} + +/// Full primary-key table scan: scans a fixed set of buckets sequentially, each +/// with its own [`KvBatchScanner`]. Buckets are scanned lazily and one at a +/// time, so only a single `ScanKv` session is open at any moment. +pub struct KvSnapshotScanner { + pending: VecDeque, + current: Option, + terminal_failure: Option, +} + +impl KvSnapshotScanner { + pub(crate) fn new(scanners: Vec) -> Self { + Self { + pending: scanners.into(), + current: None, + terminal_failure: None, + } + } + + /// Returns the next non-empty [`ScanBatch`] across all buckets, advancing to + /// the next bucket as each one drains, or `None` once every bucket is done. + pub async fn next_batch(&mut self) -> Result> { + loop { + match self + .next_batch_with_timeout(DEFAULT_KV_POLL_TIMEOUT) + .await? + { + KvBatchReadOutcome::Batch(batch) => return Ok(Some(batch)), + KvBatchReadOutcome::TimedOut => continue, + KvBatchReadOutcome::Finished => return Ok(None), + } + } + } + + /// Returns the next batch across all buckets while waiting for at most + /// `timeout`. A timeout leaves the current bucket scanner resumable. + pub async fn next_batch_with_timeout( + &mut self, + timeout: Duration, + ) -> Result { + let start = Instant::now(); + loop { + if let Some(failure) = &self.terminal_failure { + return Err(Error::UnexpectedError { + message: format!( + "KvSnapshotScanner cannot be resumed: {failure}. Create a new scanner to \ + restart the whole-table scan" + ), + source: None, + }); + } + if self.current.is_none() { + self.current = self.pending.pop_front(); + if self.current.is_none() { + return Ok(KvBatchReadOutcome::Finished); + } + } + let Some(remaining) = remaining_timeout(start, timeout) else { + return Ok(KvBatchReadOutcome::TimedOut); + }; + let result = self + .current + .as_mut() + .expect("current scanner present") + .next_batch_with_timeout(remaining) + .await; + match result { + Err(error) => { + let bucket = self + .current + .as_ref() + .expect("current scanner present") + .bucket() + .clone(); + self.fail(format!("bucket {bucket}: {error}")); + return Err(error); + } + Ok(KvBatchReadOutcome::Batch(batch)) => { + return Ok(KvBatchReadOutcome::Batch(batch)); + } + Ok(KvBatchReadOutcome::TimedOut) => { + return Ok(KvBatchReadOutcome::TimedOut); + } + Ok(KvBatchReadOutcome::Finished) => self.current = None, + } + } + } + + fn fail(&mut self, failure: String) { + self.current = None; + self.pending.clear(); + self.terminal_failure = Some(failure); + } + + /// Drains the whole-table scan into all of its batches. + pub async fn collect_all_batches(&mut self) -> Result> { + let mut batches = Vec::new(); + while let Some(batch) = self.next_batch().await? { + batches.push(batch); + } + Ok(batches) + } + + /// Closes the active bucket scanner. Pending scanners are unopened and + /// therefore hold no server-side resources. Closing before every bucket is + /// fully consumed leaves this whole-table scanner incomplete, so subsequent + /// reads return an error rather than reporting normal end-of-scan. + pub async fn close(&mut self) -> Result<()> { + let incomplete = !self.pending.is_empty() + || self + .current + .as_ref() + .map(|scanner| !scanner.is_fully_drained()) + .unwrap_or(false); + if let Some(scanner) = self.current.as_mut() { + scanner.close().await?; + } + self.current = None; + self.pending.clear(); + if incomplete && self.terminal_failure.is_none() { + self.terminal_failure = + Some("scanner was closed before every bucket was fully drained".to_string()); + } + Ok(()) + } +} + +fn metadata_refresh_partition_ids(bucket: &TableBucket) -> Vec { + bucket.partition_id().into_iter().collect() +} + +fn remaining_timeout(start: Instant, timeout: Duration) -> Option { + let elapsed = start.elapsed(); + if elapsed >= timeout { + None + } else { + Some(timeout - elapsed) + } +} + +async fn wait_for_task( + task: &mut JoinHandle, + timeout: Duration, +) -> Option> { + tokio::time::timeout(timeout, task).await.ok() +} + +#[cfg(test)] +mod tests; diff --git a/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner/tests.rs b/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner/tests.rs new file mode 100644 index 00000000000..6bdfb9fdf6a --- /dev/null +++ b/fluss-rust/crates/fluss/src/client/table/batch_scanner/kv_batch_scanner/tests.rs @@ -0,0 +1,561 @@ +// 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. + +use super::*; +use crate::client::admin::FlussAdmin; +use crate::client::metadata::Metadata; +use crate::cluster::{BucketLocation, Cluster, ServerNode, ServerType}; +use crate::metadata::{DataField, DataTypes, PhysicalTablePath, SchemaInfo, TableInfo, TablePath}; +use crate::record::kv::{SCHEMA_ID_LENGTH, ValueRecordBatch}; +use crate::row::binary::BinaryWriter; +use crate::row::compacted::CompactedRowWriter; +use crate::rpc::test_utils::{ + FramedRequest, install_duplex_connection, read_framed_request, write_error_response, + write_success_response, +}; +use crate::rpc::{ApiKey, RpcClient}; +use crate::test_utils::build_table_info_with_columns; +use bytes::Bytes; +use prost::Message; +use std::collections::HashMap; +use std::sync::Arc; +use std::time::Duration; +use tokio::io::{AsyncReadExt, DuplexStream}; +use tokio::sync::oneshot; + +fn build_two_col_table_info() -> TableInfo { + build_table_info_with_columns( + TablePath::new("db".to_string(), "tbl".to_string()), + 42, + 1, + vec![ + DataField::new("id", DataTypes::int(), None), + DataField::new("name", DataTypes::string(), None), + ], + ) +} + +fn protocol_test_cluster(table_info: &TableInfo, bucket: &TableBucket) -> Arc { + let server = ServerNode::new( + 1, + "scan-kv-test".to_string(), + 9124, + ServerType::TabletServer, + ); + let table_path = Arc::new(table_info.table_path.clone()); + let physical_table_path = Arc::new(match bucket.partition_id() { + Some(partition_id) => PhysicalTablePath::of_partitioned( + Arc::clone(&table_path), + Some(format!("partition-{partition_id}")), + ), + None => PhysicalTablePath::of(Arc::clone(&table_path)), + }); + let location = BucketLocation::new( + bucket.clone(), + Some(server.clone()), + Arc::clone(&physical_table_path), + ); + let partition_ids = bucket + .partition_id() + .map(|partition_id| HashMap::from([(Arc::clone(&physical_table_path), partition_id)])) + .unwrap_or_default(); + Arc::new(Cluster::new( + None, + HashMap::from([(server.id(), server)]), + HashMap::from([(Arc::clone(&physical_table_path), vec![location.clone()])]), + HashMap::from([(bucket.clone(), location)]), + HashMap::from([(table_info.table_path.clone(), table_info.table_id)]), + HashMap::from([(table_info.table_path.clone(), table_info.clone())]), + partition_ids, + )) +} + +fn protocol_test_scanner( + table_info: TableInfo, + bucket: TableBucket, +) -> (KvBatchScanner, DuplexStream) { + let cluster = protocol_test_cluster(&table_info, &bucket); + let leader = cluster + .leader_for(&bucket) + .expect("test bucket leader") + .clone(); + let rpc_client = Arc::new(RpcClient::new()); + let metadata = Arc::new(Metadata::new_for_test_with_connections( + cluster, + Arc::clone(&rpc_client), + )); + let server_stream = install_duplex_connection(&rpc_client, &leader); + let schema_getter = Arc::new(ClientSchemaGetter::new( + table_info.table_path.clone(), + Arc::new(FlussAdmin::new( + Arc::clone(&rpc_client), + Arc::clone(&metadata), + )), + SchemaInfo::new(table_info.get_schema().clone(), table_info.get_schema_id()), + )); + ( + KvBatchScanner::new( + rpc_client, + metadata, + table_info, + schema_getter, + None, + bucket, + 1024, + ), + server_stream, + ) +} + +fn decode_scan_request(request: &FramedRequest) -> crate::proto::ScanKvRequest { + assert_eq!(request.api_key, i16::from(ApiKey::ScanKv)); + assert_eq!(request.api_version, 0); + crate::proto::ScanKvRequest::decode(request.body.as_slice()).expect("decode ScanKV request") +} + +fn kv_record_bytes(schema_id: i16, id: i32, name: &str) -> Bytes { + let row = compacted(2, |writer| { + writer.write_int(id); + writer.write_string(name); + }); + Bytes::copy_from_slice(value_batch(&[(schema_id, row)]).data()) +} + +fn value_batch(records: &[(i16, Vec)]) -> ValueRecordBatch { + let mut body = Vec::new(); + for (schema_id, row) in records { + let record_len = (SCHEMA_ID_LENGTH + row.len()) as i32; + body.extend_from_slice(&record_len.to_le_bytes()); + body.extend_from_slice(&schema_id.to_le_bytes()); + body.extend_from_slice(row); + } + let mut bytes = Vec::new(); + bytes.extend_from_slice(&((1 + 4 + body.len()) as i32).to_le_bytes()); + bytes.push(0); + bytes.extend_from_slice(&(records.len() as i32).to_le_bytes()); + bytes.extend_from_slice(&body); + ValueRecordBatch::new(Bytes::from(bytes)) +} + +fn compacted(field_count: usize, write: impl FnOnce(&mut CompactedRowWriter)) -> Vec { + let mut writer = CompactedRowWriter::new(field_count); + write(&mut writer); + writer.to_bytes().as_ref().to_vec() +} + +#[tokio::test] +async fn kv_scanner_sequences_and_pipelines_protocol_requests() { + let table_info = build_two_col_table_info(); + let schema_id = table_info.get_schema_id() as i16; + let bucket = TableBucket::new(table_info.table_id, 0); + let (mut scanner, mut server_stream) = protocol_test_scanner(table_info, bucket.clone()); + let scanner_id = vec![1, 2, 3, 4]; + let expected_scanner_id = scanner_id.clone(); + let (open_seen_tx, open_seen_rx) = oneshot::channel(); + let (release_open_tx, release_open_rx) = oneshot::channel(); + let (continuation_seen_tx, continuation_seen_rx) = oneshot::channel(); + let (release_continuation_tx, release_continuation_rx) = oneshot::channel(); + + let server = tokio::spawn(async move { + let open_frame = read_framed_request(&mut server_stream).await; + let open = decode_scan_request(&open_frame); + assert_eq!(open.scanner_id, None); + assert_eq!(open.call_seq_id, Some(0)); + assert_eq!(open.batch_size_bytes, Some(1024)); + assert_eq!(open.close_scanner, None); + let open_bucket = open.bucket_scan_req.expect("open bucket"); + assert_eq!(open_bucket.table_id, bucket.table_id()); + assert_eq!(open_bucket.partition_id, bucket.partition_id()); + assert_eq!(open_bucket.bucket_id, bucket.bucket_id()); + open_seen_tx.send(()).expect("report open request"); + release_open_rx.await.expect("release open response"); + + write_success_response( + &mut server_stream, + open_frame.request_id, + &ScanKvResponse { + scanner_id: Some(scanner_id.clone()), + has_more_results: Some(true), + records: Some(kv_record_bytes(schema_id, 1, "one")), + log_offset: Some(99), + ..Default::default() + }, + ) + .await; + + let continuation_1_frame = read_framed_request(&mut server_stream).await; + let continuation_1 = decode_scan_request(&continuation_1_frame); + assert_eq!(continuation_1.scanner_id, Some(scanner_id.clone())); + assert_eq!(continuation_1.bucket_scan_req, None); + assert_eq!(continuation_1.call_seq_id, Some(1)); + continuation_seen_tx + .send(()) + .expect("report pipelined continuation"); + release_continuation_rx + .await + .expect("release first continuation"); + + write_success_response( + &mut server_stream, + continuation_1_frame.request_id, + &ScanKvResponse { + scanner_id: Some(scanner_id.clone()), + has_more_results: Some(true), + records: Some(kv_record_bytes(schema_id, 2, "two")), + ..Default::default() + }, + ) + .await; + + let continuation_2_frame = read_framed_request(&mut server_stream).await; + let continuation_2 = decode_scan_request(&continuation_2_frame); + assert_eq!(continuation_2.scanner_id, Some(scanner_id)); + assert_eq!(continuation_2.bucket_scan_req, None); + assert_eq!(continuation_2.call_seq_id, Some(2)); + write_success_response( + &mut server_stream, + continuation_2_frame.request_id, + &ScanKvResponse { + has_more_results: Some(false), + ..Default::default() + }, + ) + .await; + }); + + assert!(matches!( + scanner + .next_batch_with_timeout(Duration::from_millis(20)) + .await + .expect("open timeout"), + KvBatchReadOutcome::TimedOut + )); + open_seen_rx.await.expect("server observed open"); + release_open_tx.send(()).expect("release open response"); + + let first = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect("first batch"); + let KvBatchReadOutcome::Batch(first) = first else { + panic!("expected first batch"); + }; + assert_eq!(first.batch().num_rows(), 1); + assert_eq!(scanner.snapshot_log_offset(), Some(99)); + continuation_seen_rx + .await + .expect("continuation must be sent before returning the first batch"); + release_continuation_tx + .send(()) + .expect("release continuation response"); + + let second = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect("second batch"); + let KvBatchReadOutcome::Batch(second) = second else { + panic!("expected second batch"); + }; + assert_eq!(second.batch().num_rows(), 1); + assert_eq!( + scanner.scanner_id.as_deref(), + Some(expected_scanner_id.as_slice()) + ); + assert!(matches!( + scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect("terminal response"), + KvBatchReadOutcome::Finished + )); + server.await.expect("protocol server"); +} + +#[tokio::test(start_paused = true)] +async fn kv_scanner_retries_too_many_scanners_as_a_fresh_open() { + let table_info = build_two_col_table_info(); + let bucket = TableBucket::new(table_info.table_id, 0); + let (mut scanner, mut server_stream) = protocol_test_scanner(table_info, bucket); + + let server = tokio::spawn(async move { + let first_frame = read_framed_request(&mut server_stream).await; + let first = decode_scan_request(&first_frame); + assert_eq!(first.scanner_id, None); + assert_eq!(first.call_seq_id, Some(0)); + write_success_response( + &mut server_stream, + first_frame.request_id, + &ScanKvResponse { + error_code: Some(FlussError::TooManyScanners.code()), + error_message: Some("busy".to_string()), + ..Default::default() + }, + ) + .await; + + let retry_frame = read_framed_request(&mut server_stream).await; + let retry = decode_scan_request(&retry_frame); + assert_eq!(retry.scanner_id, None); + assert_eq!(retry.call_seq_id, Some(0)); + assert_eq!(retry.bucket_scan_req, first.bucket_scan_req); + write_success_response( + &mut server_stream, + retry_frame.request_id, + &ScanKvResponse { + has_more_results: Some(false), + ..Default::default() + }, + ) + .await; + }); + + assert!(matches!( + scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect("retry open scanner"), + KvBatchReadOutcome::Finished + )); + assert_eq!(scanner.open_retries, 1); + server.await.expect("retry protocol server"); +} + +#[tokio::test] +async fn kv_scanner_close_sends_close_for_the_open_session() { + let table_info = build_two_col_table_info(); + let schema_id = table_info.get_schema_id() as i16; + let bucket = TableBucket::new(table_info.table_id, 0); + let (mut scanner, mut server_stream) = protocol_test_scanner(table_info, bucket); + let scanner_id = vec![9, 8, 7]; + let expected_scanner_id = scanner_id.clone(); + + let server = tokio::spawn(async move { + let open_frame = read_framed_request(&mut server_stream).await; + write_success_response( + &mut server_stream, + open_frame.request_id, + &ScanKvResponse { + scanner_id: Some(scanner_id.clone()), + has_more_results: Some(true), + records: Some(kv_record_bytes(schema_id, 1, "one")), + ..Default::default() + }, + ) + .await; + + let next_frame = read_framed_request(&mut server_stream).await; + let next = decode_scan_request(&next_frame); + let (close_frame, close) = if next.close_scanner == Some(true) { + (next_frame, next) + } else { + assert_eq!(next.call_seq_id, Some(1)); + let close_frame = read_framed_request(&mut server_stream).await; + let close = decode_scan_request(&close_frame); + (close_frame, close) + }; + assert_eq!(close.scanner_id, Some(scanner_id)); + assert_eq!(close.bucket_scan_req, None); + assert_eq!(close.call_seq_id, None); + assert_eq!(close.close_scanner, Some(true)); + write_success_response( + &mut server_stream, + close_frame.request_id, + &ScanKvResponse::default(), + ) + .await; + }); + + assert!(matches!( + scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect("first batch"), + KvBatchReadOutcome::Batch(_) + )); + assert_eq!( + scanner.scanner_id.as_deref(), + Some(expected_scanner_id.as_slice()) + ); + scanner.close().await.expect("close scanner"); + let error = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect_err("an explicitly closed scanner must not report completion"); + assert!( + error + .to_string() + .contains("closed before it was fully drained") + ); + server.await.expect("close protocol server"); + + let table_info = build_two_col_table_info(); + let bucket = TableBucket::new(table_info.table_id, 0); + let (bucket_scanner, _server_stream) = protocol_test_scanner(table_info, bucket); + let mut snapshot_scanner = KvSnapshotScanner::new(vec![bucket_scanner]); + snapshot_scanner.close().await.expect("close snapshot"); + let error = snapshot_scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect_err("an explicitly closed snapshot must not report completion"); + assert!(error.to_string().contains("closed before every bucket")); +} + +#[tokio::test] +async fn not_leader_refresh_is_partition_aware_and_synchronous() { + let table_info = build_two_col_table_info(); + let partition_id = 77; + let bucket = TableBucket::new_with_partition(table_info.table_id, Some(partition_id), 0); + let (mut scanner, mut server_stream) = protocol_test_scanner(table_info, bucket); + let (refresh_seen_tx, refresh_seen_rx) = oneshot::channel(); + let (release_refresh_tx, release_refresh_rx) = oneshot::channel(); + + let server = tokio::spawn(async move { + let scan_frame = read_framed_request(&mut server_stream).await; + write_success_response( + &mut server_stream, + scan_frame.request_id, + &ScanKvResponse { + error_code: Some(FlussError::NotLeaderOrFollower.code()), + error_message: Some("moved".to_string()), + ..Default::default() + }, + ) + .await; + + let refresh_frame = read_framed_request(&mut server_stream).await; + assert_eq!(refresh_frame.api_key, i16::from(ApiKey::MetaData)); + assert_eq!(refresh_frame.api_version, 0); + let refresh = crate::proto::MetadataRequest::decode(refresh_frame.body.as_slice()) + .expect("decode metadata refresh"); + assert_eq!(refresh.partitions_id, vec![partition_id]); + assert_eq!(refresh.table_path.len(), 1); + assert_eq!(refresh.table_path[0].database_name, "db"); + assert_eq!(refresh.table_path[0].table_name, "tbl"); + refresh_seen_tx.send(()).expect("report metadata refresh"); + release_refresh_rx.await.expect("release metadata refresh"); + write_error_response( + &mut server_stream, + refresh_frame.request_id, + 1, + "refresh failed", + ) + .await; + }); + + let mut scan = tokio::spawn(async move { + let result = scanner + .next_batch_with_timeout(Duration::from_secs(5)) + .await; + (scanner, result) + }); + tokio::select! { + refresh = refresh_seen_rx => { + refresh.expect("partition metadata refresh request"); + } + completed = &mut scan => { + let (_, result) = completed.expect("scanner task"); + panic!("scanner returned before sending metadata refresh: {result:?}"); + } + _ = tokio::time::sleep(Duration::from_secs(1)) => { + panic!("partition metadata refresh must be sent"); + } + } + assert!( + !scan.is_finished(), + "NotLeader must not return before metadata refresh completes" + ); + release_refresh_tx + .send(()) + .expect("complete metadata refresh"); + let (scanner, result) = tokio::time::timeout(Duration::from_secs(1), &mut scan) + .await + .expect("scanner must return after refresh completes") + .expect("scanner task"); + let error = result.expect_err("NotLeader remains the authoritative error"); + assert_eq!(error.api_error(), Some(FlussError::NotLeaderOrFollower)); + assert!(scanner.closed); + tokio::time::timeout(Duration::from_secs(1), server) + .await + .expect("NotLeader protocol server must complete") + .expect("NotLeader protocol server"); +} + +#[tokio::test] +async fn kv_snapshot_scanner_stays_terminal_after_bucket_failure() { + let table_info = build_two_col_table_info(); + let (first, mut first_server_stream) = + protocol_test_scanner(table_info.clone(), TableBucket::new(table_info.table_id, 0)); + let (second, mut second_server_stream) = + protocol_test_scanner(table_info.clone(), TableBucket::new(table_info.table_id, 1)); + let first_server = tokio::spawn(async move { + let open_frame = read_framed_request(&mut first_server_stream).await; + write_success_response( + &mut first_server_stream, + open_frame.request_id, + &ScanKvResponse { + scanner_id: Some(vec![1, 2, 3]), + ..Default::default() + }, + ) + .await; + }); + let second_server = tokio::spawn(async move { + let mut byte = [0]; + if let Ok(Ok(_)) = tokio::time::timeout( + Duration::from_millis(100), + second_server_stream.read_exact(&mut byte), + ) + .await + { + panic!("whole-table scanner advanced to the next bucket"); + } + }); + + let mut scanner = KvSnapshotScanner::new(vec![first, second]); + let original_error = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect_err("first bucket must fail"); + assert!( + original_error + .to_string() + .contains("did not include has_more_results") + ); + + let error = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect_err("a failed whole-table scanner must remain terminal"); + assert!( + error + .to_string() + .contains("KvSnapshotScanner cannot be resumed") + ); + + scanner.close().await.expect("close terminal scanner"); + let error = scanner + .next_batch_with_timeout(Duration::from_secs(1)) + .await + .expect_err("close must not turn a failed scan into normal completion"); + assert!( + error + .to_string() + .contains("KvSnapshotScanner cannot be resumed") + ); + first_server.await.expect("first bucket protocol server"); + second_server.await.expect("second bucket protocol server"); +} diff --git a/fluss-rust/crates/fluss/src/client/table/mod.rs b/fluss-rust/crates/fluss/src/client/table/mod.rs index 1069ba08234..c28dfa1515c 100644 --- a/fluss-rust/crates/fluss/src/client/table/mod.rs +++ b/fluss-rust/crates/fluss/src/client/table/mod.rs @@ -37,7 +37,7 @@ mod scanner; mod upsert; pub use append::{AppendWriter, TableAppend}; -pub use batch_scanner::LimitBatchScanner; +pub use batch_scanner::{KvBatchReadOutcome, KvBatchScanner, KvSnapshotScanner, LimitBatchScanner}; pub use lookup::{LookupResult, Lookuper, PrefixKeyLookuper, TableLookup, TablePrefixLookup}; pub use reader::{ BoundedCollectOutcome, BoundedLogReadRange, RecordBatchLogReader, RecordBatchReadOutcome, diff --git a/fluss-rust/crates/fluss/src/client/table/scanner.rs b/fluss-rust/crates/fluss/src/client/table/scanner.rs index 6aba4c306fa..a1c3fa3e2a6 100644 --- a/fluss-rust/crates/fluss/src/client/table/scanner.rs +++ b/fluss-rust/crates/fluss/src/client/table/scanner.rs @@ -19,7 +19,7 @@ use crate::client::ClientSchemaGetter; use crate::client::connection::FlussConnection; use crate::client::credentials::SecurityTokenManager; use crate::client::metadata::Metadata; -use crate::client::table::batch_scanner::LimitBatchScanner; +use crate::client::table::batch_scanner::{KvBatchScanner, KvSnapshotScanner, LimitBatchScanner}; use crate::client::table::log_fetch_buffer::{ CompletedFetch, DefaultCompletedFetch, FetchErrorAction, FetchErrorContext, FetchErrorLogLevel, FetchResult, LogFetchBuffer, NO_FILTERED_END_OFFSET, RemotePendingFetch, @@ -222,6 +222,123 @@ impl<'a> TableScan<'a> { )) } + /// Pre-seeds a schema getter with the current schema; older versions are + /// fetched lazily during KV decode. Mirrors `Table::new_lookup`. + fn build_kv_schema_getter(&self) -> Result> { + let latest = SchemaInfo::new( + self.table_info.get_schema().clone(), + self.table_info.get_schema_id(), + ); + Ok(Arc::new(ClientSchemaGetter::new( + self.table_info.table_path.clone(), + self.conn.get_admin()?, + latest, + ))) + } + + /// Rejects a full KV scan on a log table or when a limit is configured. + fn ensure_kv_scan_supported(&self) -> Result<()> { + self.reject_filter("KV scanner")?; + if !self.table_info.has_primary_key() { + return Err(Error::UnsupportedOperation { + message: format!( + "Full KV scan is only supported for primary key tables. Table: {}", + self.table_info.table_path + ), + }); + } + if let Some(limit) = self.limit { + return Err(Error::UnsupportedOperation { + message: format!( + "Full KV scan doesn't support limit pushdown; use create_bucket_batch_scanner for a bounded scan. Table: {}, requested limit: {limit}", + self.table_info.table_path + ), + }); + } + Ok(()) + } + + /// Creates a full scan of a single primary-key bucket via the `ScanKv` RPC. + /// + /// Reads the current committed state of the bucket (one row per primary + /// key, already merged — not the changelog) with snapshot isolation. + /// Requires a primary-key table and no configured limit. Creation is cheap; + /// the first `ScanKv` request runs on the first + /// [`KvBatchScanner::next_batch`]. + pub fn create_bucket_kv_scanner(self, table_bucket: TableBucket) -> Result { + self.ensure_kv_scan_supported()?; + if table_bucket.table_id() != self.table_info.table_id { + return Err(Error::IllegalArgument { + message: format!( + "Bucket table_id {} does not match scan table_id {}", + table_bucket.table_id(), + self.table_info.table_id + ), + }); + } + let num_buckets = self.table_info.get_num_buckets(); + if table_bucket.bucket_id() < 0 || table_bucket.bucket_id() >= num_buckets { + return Err(Error::IllegalArgument { + message: format!( + "Bucket id {} out of range for table with {num_buckets} buckets", + table_bucket.bucket_id() + ), + }); + } + let schema_getter = self.build_kv_schema_getter()?; + let batch_size_bytes = self.conn.config().scanner_kv_fetch_max_bytes; + Ok(KvBatchScanner::new( + self.conn.get_connections(), + self.metadata.clone(), + self.table_info, + schema_getter, + self.projected_fields, + table_bucket, + batch_size_bytes, + )) + } + + /// Creates a full scan of an entire primary-key table, scanning every + /// `(partition, bucket)` sequentially via the `ScanKv` RPC. + /// + /// Requires a primary-key table and no configured limit or filter. For a + /// partitioned table, partition metadata is resolved once when the scanner + /// is created. + pub async fn create_kv_scanner(self) -> Result { + self.ensure_kv_scan_supported()?; + let schema_getter = self.build_kv_schema_getter()?; + let batch_size_bytes = self.conn.config().scanner_kv_fetch_max_bytes; + let rpc_client = self.conn.get_connections(); + let table_id = self.table_info.table_id; + let num_buckets = self.table_info.get_num_buckets(); + let partition_ids = if self.table_info.is_partitioned() { + self.conn + .get_admin()? + .list_partition_infos(&self.table_info.table_path) + .await? + .into_iter() + .map(|partition| Some(partition.get_partition_id())) + .collect() + } else { + vec![None] + }; + let mut scanners = Vec::with_capacity(partition_ids.len() * num_buckets as usize); + for partition_id in partition_ids { + for bucket_id in 0..num_buckets { + scanners.push(KvBatchScanner::new( + rpc_client.clone(), + self.metadata.clone(), + self.table_info.clone(), + schema_getter.clone(), + self.projected_fields.clone(), + TableBucket::new_with_partition(table_id, partition_id, bucket_id), + batch_size_bytes, + )); + } + } + Ok(KvSnapshotScanner::new(scanners)) + } + /// Projects the scan to only include specified columns by their indices. /// /// # Arguments diff --git a/fluss-rust/crates/fluss/src/config.rs b/fluss-rust/crates/fluss/src/config.rs index f2d932f7dd2..feb7d2f7a0b 100644 --- a/fluss-rust/crates/fluss/src/config.rs +++ b/fluss-rust/crates/fluss/src/config.rs @@ -35,6 +35,7 @@ const DEFAULT_SCANNER_LOG_FETCH_MIN_BYTES: i32 = 1; const DEFAULT_SCANNER_LOG_FETCH_WAIT_MAX_TIME_MS: i32 = 500; const DEFAULT_WRITER_BATCH_TIMEOUT_MS: i64 = 100; const DEFAULT_SCANNER_LOG_FETCH_MAX_BYTES_FOR_BUCKET: i32 = 1024 * 1024; +const DEFAULT_SCANNER_KV_FETCH_MAX_BYTES: i32 = 4 * 1024 * 1024; const DEFAULT_WRITER_MAX_INFLIGHT_REQUESTS_PER_BUCKET: usize = 5; const DEFAULT_WRITER_BUFFER_MEMORY_SIZE: usize = 64 * 1024 * 1024; // 64MB, matching Java const DEFAULT_WRITER_BUFFER_WAIT_TIMEOUT_MS: u64 = u64::MAX; @@ -139,6 +140,12 @@ pub struct Config { #[arg(long, default_value_t = DEFAULT_SCANNER_LOG_FETCH_MAX_BYTES_FOR_BUCKET)] pub scanner_log_fetch_max_bytes_for_bucket: i32, + /// Maximum bytes of record data the server returns per `ScanKv` request when + /// performing a full primary-key table scan (`BatchScanner` without a limit). + /// Default: 4194304 (4MB, matching Java `client.scanner.kv.fetch.max-bytes`) + #[arg(long, default_value_t = DEFAULT_SCANNER_KV_FETCH_MAX_BYTES)] + pub scanner_kv_fetch_max_bytes: i32, + /// Whether to enable idempotent writes. When enabled, each batch is tagged with /// a server-allocated writer ID and per-bucket sequence number so the server can /// detect and deduplicate retried batches. @@ -256,6 +263,10 @@ impl std::fmt::Debug for Config { "scanner_log_fetch_max_bytes_for_bucket", &self.scanner_log_fetch_max_bytes_for_bucket, ) + .field( + "scanner_kv_fetch_max_bytes", + &self.scanner_kv_fetch_max_bytes, + ) .field( "scanner_log_fetch_wait_max_time_ms", &self.scanner_log_fetch_wait_max_time_ms, @@ -311,6 +322,7 @@ impl Default for Config { scanner_log_fetch_min_bytes: DEFAULT_SCANNER_LOG_FETCH_MIN_BYTES, scanner_log_fetch_wait_max_time_ms: DEFAULT_SCANNER_LOG_FETCH_WAIT_MAX_TIME_MS, scanner_log_fetch_max_bytes_for_bucket: DEFAULT_SCANNER_LOG_FETCH_MAX_BYTES_FOR_BUCKET, + scanner_kv_fetch_max_bytes: DEFAULT_SCANNER_KV_FETCH_MAX_BYTES, writer_batch_timeout_ms: DEFAULT_WRITER_BATCH_TIMEOUT_MS, writer_enable_idempotence: true, writer_max_inflight_requests_per_bucket: @@ -382,6 +394,9 @@ impl Config { if self.scanner_log_fetch_max_bytes <= 0 { return Err("scanner_log_fetch_max_bytes must be > 0".to_string()); } + if self.scanner_kv_fetch_max_bytes <= 0 { + return Err("scanner_kv_fetch_max_bytes must be > 0".to_string()); + } if self.scanner_log_fetch_max_bytes < self.scanner_log_fetch_min_bytes { return Err( "scanner_log_fetch_max_bytes must be >= scanner_log_fetch_min_bytes".to_string(), @@ -571,6 +586,15 @@ mod tests { assert!(config.validate_scanner().is_err()); } + #[test] + fn test_scanner_kv_fetch_max_bytes_zero() { + let config = Config { + scanner_kv_fetch_max_bytes: 0, + ..Config::default() + }; + assert!(config.validate_scanner().is_err()); + } + #[test] fn test_scanner_fetch_negative_wait() { let config = Config { diff --git a/fluss-rust/crates/fluss/src/rpc/fluss_api_error.rs b/fluss-rust/crates/fluss/src/rpc/fluss_api_error.rs index 10d0dd37969..18d402d6244 100644 --- a/fluss-rust/crates/fluss/src/rpc/fluss_api_error.rs +++ b/fluss-rust/crates/fluss/src/rpc/fluss_api_error.rs @@ -171,6 +171,14 @@ pub enum FlussError { InvalidAlterTableException = 56, /// Deletion operations are disabled on this table. DeletionDisabledException = 57, + /// The scanner session has expired due to inactivity. + ScannerExpired = 66, + /// The scanner id is not recognized by the server. + UnknownScannerId = 67, + /// The scan request is invalid. + InvalidScanRequest = 68, + /// The per-bucket or per-server scanner session limit has been reached. + TooManyScanners = 69, /// The KV storage engine rejected a write due to backpressure. StorageBackpressureException = 72, } @@ -301,6 +309,12 @@ impl FlussError { FlussError::DeletionDisabledException => { "Deletion operations are disabled on this table." } + FlussError::ScannerExpired => "The scanner session has expired due to inactivity.", + FlussError::UnknownScannerId => "The scanner id is not recognized by the server.", + FlussError::InvalidScanRequest => "The scan request is invalid.", + FlussError::TooManyScanners => { + "The per-bucket or per-server scanner session limit has been reached." + } FlussError::StorageBackpressureException => { "The tablet server has rejected the write because the KV storage engine has reached its write-pressure threshold." } @@ -378,6 +392,10 @@ impl FlussError { 55 => FlussError::IneligibleReplicaException, 56 => FlussError::InvalidAlterTableException, 57 => FlussError::DeletionDisabledException, + 66 => FlussError::ScannerExpired, + 67 => FlussError::UnknownScannerId, + 68 => FlussError::InvalidScanRequest, + 69 => FlussError::TooManyScanners, 72 => FlussError::StorageBackpressureException, _ => FlussError::UnknownServerError, } @@ -421,6 +439,8 @@ mod tests { FlussError::for_code(72), FlussError::StorageBackpressureException ); + assert_eq!(FlussError::for_code(66), FlussError::ScannerExpired); + assert_eq!(FlussError::for_code(69), FlussError::TooManyScanners); assert_eq!(FlussError::for_code(9999), FlussError::UnknownServerError); } @@ -505,6 +525,10 @@ mod tests { FlussError::FencedLeaderEpochException, FlussError::FencedTieringEpochException, FlussError::RetriableAuthenticateException, + FlussError::ScannerExpired, + FlussError::UnknownScannerId, + FlussError::InvalidScanRequest, + FlussError::TooManyScanners, ]; for err in &non_retriable { assert!(!err.is_retriable(), "{err:?} should not be retriable"); diff --git a/fluss-rust/crates/fluss/src/rpc/message/scan_kv.rs b/fluss-rust/crates/fluss/src/rpc/message/scan_kv.rs index 081b49169ba..2c063d93057 100644 --- a/fluss-rust/crates/fluss/src/rpc/message/scan_kv.rs +++ b/fluss-rust/crates/fluss/src/rpc/message/scan_kv.rs @@ -28,7 +28,6 @@ pub struct ScanKvRequest { } impl ScanKvRequest { - #[allow(dead_code)] pub(crate) fn new( scanner_id: Option>, bucket_scan_req: Option, diff --git a/fluss-rust/crates/fluss/src/rpc/mod.rs b/fluss-rust/crates/fluss/src/rpc/mod.rs index a4b191d8a2b..c0b8fa52400 100644 --- a/fluss-rust/crates/fluss/src/rpc/mod.rs +++ b/fluss-rust/crates/fluss/src/rpc/mod.rs @@ -27,6 +27,8 @@ pub use error::*; mod server_connection; pub use server_connection::*; mod convert; +#[cfg(test)] +pub(crate) mod test_utils; mod transport; pub(crate) use convert::*; diff --git a/fluss-rust/crates/fluss/src/rpc/server_connection.rs b/fluss-rust/crates/fluss/src/rpc/server_connection.rs index a8b36cbecc2..469a1d469bb 100644 --- a/fluss-rust/crates/fluss/src/rpc/server_connection.rs +++ b/fluss-rust/crates/fluss/src/rpc/server_connection.rs @@ -238,6 +238,17 @@ impl RpcClient { Ok(new_server) } + #[cfg(test)] + pub(crate) fn insert_connection_for_test( + &self, + server_node: &ServerNode, + connection: ServerConnection, + ) { + self.connections + .write() + .insert(server_node.uid().to_owned(), connection); + } + async fn connect(&self, server_node: &ServerNode) -> Result { let url = server_node.url(); let transport = Transport::connect(&url, self.timeout) @@ -643,6 +654,15 @@ where } } +#[cfg(test)] +pub(crate) fn server_connection_from_duplex(stream: tokio::io::DuplexStream) -> ServerConnection { + Arc::new(ServerConnectionInner::new( + BufStream::new(Transport::from_duplex(stream)), + usize::MAX, + Arc::from("scan-kv-test"), + )) +} + impl Drop for ServerConnectionInner { fn drop(&mut self) { // todo: should remove from server_connections map? diff --git a/fluss-rust/crates/fluss/src/rpc/test_utils.rs b/fluss-rust/crates/fluss/src/rpc/test_utils.rs new file mode 100644 index 00000000000..d676f2c6d72 --- /dev/null +++ b/fluss-rust/crates/fluss/src/rpc/test_utils.rs @@ -0,0 +1,99 @@ +// 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. + +use super::{RpcClient, server_connection_from_duplex}; +use crate::cluster::ServerNode; +use crate::proto::ErrorResponse; +use prost::Message; +use tokio::io::{AsyncReadExt, AsyncWriteExt, DuplexStream}; + +pub(crate) struct FramedRequest { + pub(crate) api_key: i16, + pub(crate) api_version: i16, + pub(crate) request_id: i32, + pub(crate) body: Vec, +} + +pub(crate) fn install_duplex_connection( + rpc_client: &RpcClient, + server_node: &ServerNode, +) -> DuplexStream { + let (client_stream, server_stream) = tokio::io::duplex(64 * 1024); + rpc_client + .insert_connection_for_test(server_node, server_connection_from_duplex(client_stream)); + server_stream +} + +pub(crate) async fn read_framed_request(stream: &mut DuplexStream) -> FramedRequest { + let length = stream.read_i32().await.expect("request frame length"); + assert!(length >= 8, "request frame must contain its header"); + let mut payload = vec![0; length as usize]; + stream + .read_exact(&mut payload) + .await + .expect("request frame payload"); + FramedRequest { + api_key: i16::from_be_bytes([payload[0], payload[1]]), + api_version: i16::from_be_bytes([payload[2], payload[3]]), + request_id: i32::from_be_bytes([payload[4], payload[5], payload[6], payload[7]]), + body: payload[8..].to_vec(), + } +} + +pub(crate) async fn write_success_response( + stream: &mut DuplexStream, + request_id: i32, + response: &impl Message, +) { + let mut payload = Vec::new(); + payload.push(0); + payload.extend_from_slice(&request_id.to_be_bytes()); + response + .encode(&mut payload) + .expect("encode successful response"); + write_response_frame(stream, payload).await; +} + +pub(crate) async fn write_error_response( + stream: &mut DuplexStream, + request_id: i32, + code: i32, + message: &str, +) { + let mut payload = Vec::new(); + payload.push(1); + payload.extend_from_slice(&request_id.to_be_bytes()); + ErrorResponse { + error_code: code, + error_message: Some(message.to_string()), + } + .encode(&mut payload) + .expect("encode error response"); + write_response_frame(stream, payload).await; +} + +async fn write_response_frame(stream: &mut DuplexStream, payload: Vec) { + stream + .write_i32(payload.len() as i32) + .await + .expect("response frame length"); + stream + .write_all(&payload) + .await + .expect("response frame payload"); + stream.flush().await.expect("flush response"); +} diff --git a/fluss-rust/crates/fluss/src/rpc/transport.rs b/fluss-rust/crates/fluss/src/rpc/transport.rs index a6f721f6aaa..aa5a911774f 100644 --- a/fluss-rust/crates/fluss/src/rpc/transport.rs +++ b/fluss-rust/crates/fluss/src/rpc/transport.rs @@ -20,12 +20,20 @@ use std::ops::DerefMut; use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Duration; +#[cfg(test)] +use tokio::io::DuplexStream; use tokio::io::{AsyncRead, AsyncWrite, ReadBuf}; use tokio::net::TcpStream; #[derive(Debug)] pub enum Transport { - Plain { inner: TcpStream }, + Plain { + inner: TcpStream, + }, + #[cfg(test)] + Test { + inner: DuplexStream, + }, } impl AsyncRead for Transport { @@ -36,6 +44,8 @@ impl AsyncRead for Transport { ) -> Poll> { match self.deref_mut() { Self::Plain { inner } => Pin::new(inner).poll_read(cx, buf), + #[cfg(test)] + Self::Test { inner } => Pin::new(inner).poll_read(cx, buf), } } } @@ -48,23 +58,34 @@ impl AsyncWrite for Transport { ) -> Poll> { match self.deref_mut() { Self::Plain { inner } => Pin::new(inner).poll_write(cx, buf), + #[cfg(test)] + Self::Test { inner } => Pin::new(inner).poll_write(cx, buf), } } fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { match self.deref_mut() { Self::Plain { inner } => Pin::new(inner).poll_flush(cx), + #[cfg(test)] + Self::Test { inner } => Pin::new(inner).poll_flush(cx), } } fn poll_shutdown(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { match self.deref_mut() { Self::Plain { inner } => Pin::new(inner).poll_shutdown(cx), + #[cfg(test)] + Self::Test { inner } => Pin::new(inner).poll_shutdown(cx), } } } impl Transport { + #[cfg(test)] + pub(crate) fn from_duplex(inner: DuplexStream) -> Self { + Self::Test { inner } + } + pub async fn connect(server: &str, timeout: Option) -> Result { let tcp_stream = Self::connect_timeout(server, timeout).await?; Ok(Transport::Plain { inner: tcp_stream }) diff --git a/fluss-rust/crates/fluss/tests/integration/batch_scanner.rs b/fluss-rust/crates/fluss/tests/integration/batch_scanner.rs index 57002cde88d..e068ba362c4 100644 --- a/fluss-rust/crates/fluss/tests/integration/batch_scanner.rs +++ b/fluss-rust/crates/fluss/tests/integration/batch_scanner.rs @@ -18,12 +18,18 @@ #[cfg(test)] mod batch_scanner_test { - use crate::integration::utils::{create_table, get_shared_cluster, wait_for_table_ready}; + use crate::integration::utils::{ + create_partitions, create_table, get_shared_cluster, wait_for_partition_buckets_ready, + wait_for_table_buckets_ready, wait_for_table_ready, + }; use arrow::array::{Array, Int32Array, Int64Array, StringArray, record_batch}; + use fluss::client::FlussConnection; + use fluss::config::Config; use fluss::metadata::{ AddColumn, AlterTableChanges, ColumnPositionType, DataTypes, JsonSerde, LogFormat, Schema, TableBucket, TableDescriptor, TablePath, }; + use fluss::predicate::col; use fluss::row::GenericRow; use futures::TryStreamExt; use std::collections::HashMap; @@ -591,4 +597,373 @@ mod batch_scanner_test { "create_record_batch_log_scanner must reject a configured limit" ); } + + // ---- full KV scan (ScanKv) --------------------------------------------- + + fn id_name_pk_descriptor(num_buckets: i32) -> TableDescriptor { + TableDescriptor::builder() + .schema( + Schema::builder() + .column("id", DataTypes::int()) + .column("name", DataTypes::string()) + .primary_key(vec!["id"]) + .expect("primary key") + .build() + .expect("schema"), + ) + .distributed_by(Some(num_buckets), vec!["id".to_string()]) + .build() + .expect("descriptor") + } + + async fn create_id_name_pk_table<'a>( + connection: &'a FlussConnection, + table_name: &str, + num_buckets: i32, + ) -> fluss::client::FlussTable<'a> { + let admin = connection.get_admin().expect("admin"); + let table_path = TablePath::new("fluss", table_name); + create_table(&admin, &table_path, &id_name_pk_descriptor(num_buckets)).await; + if num_buckets == 1 { + wait_for_table_ready(&admin, &table_path).await; + } else { + let buckets: Vec = (0..num_buckets).collect(); + wait_for_table_buckets_ready(&admin, &table_path, &buckets).await; + } + connection.get_table(&table_path).await.expect("table") + } + + async fn upsert_id_name_rows( + table: &fluss::client::FlussTable<'_>, + rows: &HashMap, + ) { + let writer = table + .new_upsert() + .expect("upsert") + .create_writer() + .expect("writer"); + for (id, name) in rows { + let mut row = GenericRow::new(2); + row.set_field(0, *id); + row.set_field(1, name.as_str()); + writer.upsert(&row).expect("upsert row"); + } + writer.flush().await.expect("flush"); + } + + fn id_name_rows(ids: impl IntoIterator) -> HashMap { + ids.into_iter() + .map(|id| (id, format!("name-{id}"))) + .collect() + } + + /// Collect every (id, name) pair across all batches, asserting each key is + /// seen exactly once — a full KV scan returns merged state, not a changelog. + fn collect_id_name(batches: &[fluss::record::ScanBatch]) -> HashMap { + let mut seen: HashMap = HashMap::new(); + for scan_batch in batches { + let rows = scan_batch.batch(); + let ids = rows + .column(0) + .as_any() + .downcast_ref::() + .expect("id column Int32"); + let names = rows + .column(1) + .as_any() + .downcast_ref::() + .expect("name column Utf8"); + for i in 0..rows.num_rows() { + let prev = seen.insert(ids.value(i), names.value(i).to_string()); + assert!( + prev.is_none(), + "key {} returned more than once", + ids.value(i) + ); + } + } + seen + } + + /// Bucket and whole-table entry points both cover the complete merged state. + #[tokio::test] + async fn kv_scanners_cover_bucket_and_whole_table_reads() { + let connection = get_shared_cluster().get_fluss_connection().await; + let table = create_id_name_pk_table(&connection, "test_kv_scan_current_state", 3).await; + + let mut empty_scanner = table + .new_scan() + .create_kv_scanner() + .await + .expect("create empty whole-table scanner"); + assert!( + empty_scanner + .collect_all_batches() + .await + .expect("scan empty table") + .is_empty() + ); + + let expected = id_name_rows(1..=20); + upsert_id_name_rows(&table, &expected).await; + + let mut table_scanner = table + .new_scan() + .create_kv_scanner() + .await + .expect("create whole-table scanner"); + let whole_table = table_scanner + .collect_all_batches() + .await + .expect("scan whole table"); + assert_eq!(collect_id_name(&whole_table), expected); + + let mut across_buckets = HashMap::new(); + for bucket_id in 0..3 { + let bucket = TableBucket::new(table.get_table_info().table_id, bucket_id); + let mut scanner = table + .new_scan() + .create_bucket_kv_scanner(bucket.clone()) + .expect("create bucket scanner"); + let batches = scanner.collect_all_batches().await.expect("scan bucket"); + assert!(batches.iter().all(|batch| batch.bucket() == &bucket)); + for (id, name) in collect_id_name(&batches) { + assert!( + across_buckets.insert(id, name).is_none(), + "key {id} returned by multiple buckets" + ); + } + } + assert_eq!(across_buckets, expected); + } + + /// Whole-table scans enumerate every partition and bucket once. + #[tokio::test] + async fn kv_whole_table_scanner_covers_partitioned_table() { + let cluster = get_shared_cluster(); + let connection = cluster.get_fluss_connection().await; + let admin = connection.get_admin().expect("admin"); + + let table_path = TablePath::new("fluss", "test_kv_scan_partitioned_table"); + let descriptor = TableDescriptor::builder() + .schema( + Schema::builder() + .column("id", DataTypes::int()) + .column("name", DataTypes::string()) + .column("region", DataTypes::string()) + .primary_key(vec!["id", "region"]) + .expect("primary key") + .build() + .expect("schema"), + ) + .distributed_by(Some(2), vec!["id".to_string()]) + .partitioned_by(vec!["region"]) + .build() + .expect("descriptor"); + create_table(&admin, &table_path, &descriptor).await; + create_partitions(&admin, &table_path, "region", &["US", "EU"]).await; + wait_for_partition_buckets_ready(&admin, &table_path, "US", &[0, 1]).await; + wait_for_partition_buckets_ready(&admin, &table_path, "EU", &[0, 1]).await; + + let table = connection.get_table(&table_path).await.expect("table"); + let writer = table + .new_upsert() + .expect("upsert") + .create_writer() + .expect("writer"); + + let rows = [ + (1, "name-1", "US"), + (2, "name-2", "US"), + (3, "name-3", "EU"), + (4, "name-4", "EU"), + ]; + for &(id, name, region) in &rows { + let mut row = GenericRow::new(3); + row.set_field(0, id); + row.set_field(1, name); + row.set_field(2, region); + writer.upsert(&row).expect("upsert row"); + } + writer.flush().await.expect("flush"); + + let mut scanner = table + .new_scan() + .create_kv_scanner() + .await + .expect("create partitioned whole-table scanner"); + let batches = scanner + .collect_all_batches() + .await + .expect("scan partitioned primary-key table"); + + let seen = collect_id_name(&batches); + let expected: HashMap = rows + .iter() + .map(|(id, name, _)| (*id, (*name).to_string())) + .collect(); + assert_eq!(seen, expected); + + let regions: std::collections::HashSet = batches + .iter() + .flat_map(|batch| { + let rows = batch.batch(); + let regions = rows + .column(2) + .as_any() + .downcast_ref::() + .expect("region column"); + (0..rows.num_rows()) + .map(|row| regions.value(row).to_string()) + .collect::>() + }) + .collect(); + assert_eq!(regions, ["US".to_string(), "EU".to_string()].into()); + } + + /// A small fetch size forces continuation RPCs; rows written after the open + /// response must remain invisible to the bucket snapshot. + #[tokio::test] + async fn kv_bucket_scanner_preserves_snapshot_across_continuations() { + let cluster = get_shared_cluster(); + let setup_connection = cluster.get_fluss_connection().await; + let setup_table = + create_id_name_pk_table(&setup_connection, "test_kv_scan_snapshot", 1).await; + let table_path = setup_table.table_path().clone(); + + let scan_connection = FlussConnection::new(Config { + bootstrap_servers: cluster.plaintext_bootstrap_servers().to_string(), + writer_acks: "all".to_string(), + scanner_kv_fetch_max_bytes: 128, + ..Config::default() + }) + .await + .expect("small-fetch connection"); + let table = scan_connection.get_table(&table_path).await.expect("table"); + let initial_rows: HashMap = (0..20) + .map(|id| (id, format!("initial-value-{id:02}-with-padding"))) + .collect(); + upsert_id_name_rows(&table, &initial_rows).await; + + let bucket = TableBucket::new(table.get_table_info().table_id, 0); + let mut scanner = table + .new_scan() + .create_bucket_kv_scanner(bucket) + .expect("create snapshot scanner"); + let first_batch = scanner + .next_batch() + .await + .expect("open snapshot scan") + .expect("non-empty table must return KV rows"); + assert!( + first_batch.batch().num_rows() < initial_rows.len(), + "small fetch size must force at least one continuation" + ); + assert!( + scanner.snapshot_log_offset().is_some(), + "open response must expose the snapshot log offset" + ); + + let post_snapshot_id = 20; + let writer = table + .new_upsert() + .expect("upsert") + .create_writer() + .expect("writer"); + let mut row = GenericRow::new(2); + row.set_field(0, post_snapshot_id); + row.set_field(1, "post-snapshot"); + writer.upsert(&row).expect("upsert post-snapshot row"); + writer.flush().await.expect("flush post-snapshot row"); + + let mut batches = vec![first_batch]; + batches.extend( + scanner + .collect_all_batches() + .await + .expect("continue snapshot scan"), + ); + let seen = collect_id_name(&batches); + assert_eq!(seen, initial_rows); + assert!(!seen.contains_key(&post_snapshot_id)); + } + + /// Unsupported table shapes, pushdowns, and bucket coordinates fail before ScanKV. + #[tokio::test] + async fn kv_scanner_rejects_unsupported_requests() { + let connection = get_shared_cluster().get_fluss_connection().await; + let admin = connection.get_admin().expect("admin"); + + let log_path = TablePath::new("fluss", "test_kv_scan_reject_log"); + let log_descriptor = TableDescriptor::builder() + .schema( + Schema::builder() + .column("id", DataTypes::int()) + .build() + .expect("schema"), + ) + .distributed_by(Some(1), vec!["id".to_string()]) + .build() + .expect("descriptor"); + create_table(&admin, &log_path, &log_descriptor).await; + let log_table = connection.get_table(&log_path).await.expect("log table"); + let log_bucket = TableBucket::new(log_table.get_table_info().table_id, 0); + assert!(log_table.new_scan().create_kv_scanner().await.is_err()); + assert!( + log_table + .new_scan() + .create_bucket_kv_scanner(log_bucket) + .is_err() + ); + + let table = create_id_name_pk_table(&connection, "test_kv_scan_reject_options", 1).await; + let table_id = table.get_table_info().table_id; + let bucket = TableBucket::new(table_id, 0); + assert!( + table + .new_scan() + .limit(5) + .expect("limit") + .create_kv_scanner() + .await + .is_err() + ); + assert!( + table + .new_scan() + .limit(5) + .expect("limit") + .create_bucket_kv_scanner(bucket.clone()) + .is_err() + ); + assert!( + table + .new_scan() + .filter(col("id").gt(0)) + .expect("filter") + .create_kv_scanner() + .await + .is_err() + ); + assert!( + table + .new_scan() + .filter(col("id").gt(0)) + .expect("filter") + .create_bucket_kv_scanner(bucket) + .is_err() + ); + + for invalid_bucket in [ + TableBucket::new(table_id + 9999, 0), + TableBucket::new(table_id, 99), + ] { + assert!( + table + .new_scan() + .create_bucket_kv_scanner(invalid_bucket) + .is_err() + ); + } + } } diff --git a/fluss-rust/website/docs/user-guide/cpp/api-reference.md b/fluss-rust/website/docs/user-guide/cpp/api-reference.md index 91fe29a8df4..1b841029320 100644 --- a/fluss-rust/website/docs/user-guide/cpp/api-reference.md +++ b/fluss-rust/website/docs/user-guide/cpp/api-reference.md @@ -35,6 +35,7 @@ Complete API reference for the Fluss C++ client. | `scanner_log_fetch_min_bytes` | `int32_t` | `1` | Minimum bytes the server must accumulate before returning a fetch response | | `scanner_log_fetch_wait_max_time_ms` | `int32_t` | `500` | Maximum time (ms) the server may wait to satisfy min-bytes | | `scanner_log_fetch_max_bytes_for_bucket`| `int32_t` | `1048576` (1 MB) | Maximum bytes per fetch response per bucket for LogScanner | +| `scanner_kv_fetch_max_bytes` | `int32_t` | `4194304` (4 MB) | Maximum record bytes returned by each full KV scan RPC | | `connect_timeout_ms` | `uint64_t` | `120000` | TCP connect timeout in milliseconds | | `security_protocol` | `std::string` | `"PLAINTEXT"` | `"PLAINTEXT"` (default) or `"sasl"` for SASL auth | | `security_sasl_mechanism` | `std::string` | `"PLAIN"` | SASL mechanism (only `"PLAIN"` is supported) | diff --git a/fluss-rust/website/docs/user-guide/rust/api-reference.md b/fluss-rust/website/docs/user-guide/rust/api-reference.md index dfc87baebf3..0f1678f1765 100644 --- a/fluss-rust/website/docs/user-guide/rust/api-reference.md +++ b/fluss-rust/website/docs/user-guide/rust/api-reference.md @@ -27,6 +27,7 @@ Complete API reference for the Fluss Rust client. | `scanner_log_fetch_min_bytes` | `i32` | `1` | Minimum bytes the server must accumulate before returning a fetch response | | `scanner_log_fetch_wait_max_time_ms` | `i32` | `500` | Maximum time (ms) the server may wait to satisfy min-bytes | | `scanner_log_fetch_max_bytes_for_bucket`| `i32` | `1048576` (1 MB) | Maximum bytes per fetch response per bucket for LogScanner | +| `scanner_kv_fetch_max_bytes` | `i32` | `4194304` (4 MB) | Maximum record bytes returned by each full KV scan RPC | | `connect_timeout_ms` | `u64` | `120000` | TCP connect timeout in milliseconds | | `security_protocol` | `String` | `"PLAINTEXT"` | `PLAINTEXT` (default) or `sasl` for SASL auth | | `security_sasl_mechanism` | `String` | `"PLAIN"` | SASL mechanism (only `PLAIN` is supported) | @@ -146,10 +147,12 @@ series are shared by `AppendWriter` (log tables) and `UpsertWriter` (PK tables). | `fn project(self, indices: &[usize]) -> Result` | Project columns by index | | `fn project_by_name(self, names: &[&str]) -> Result` | Project columns by name | | `fn limit(self, n: i32) -> Result` | Set a row limit (enables `create_bucket_batch_scanner`; rejected by log scanners) | -| `fn filter(self, predicate: Predicate) -> Result` | Push a predicate down to log scanners; whole batches are pruned by statistics, and returned batches can still hold non-matching rows (see [Filter Pushdown](example/filter-pushdown.md); rejected by `create_bucket_batch_scanner`) | +| `fn filter(self, predicate: Predicate) -> Result` | Push a predicate down to log scanners; whole batches are pruned by statistics, and returned batches can still hold non-matching rows (see [Filter Pushdown](example/filter-pushdown.md); rejected by bounded and KV batch scanners) | | `fn create_log_scanner(self) -> Result` | Create a record-based log scanner; on a primary-key table, subscribes to its CDC changelog (per-record `ChangeType`) | | `fn create_record_batch_log_scanner(self) -> Result` | Create an Arrow batch-based log scanner (log tables only — no per-record change types) | | `fn create_bucket_batch_scanner(self, bucket: TableBucket) -> Result` | Bounded scan of one bucket (requires `limit`; runs on first `next_batch`) | +| `fn create_bucket_kv_scanner(self, bucket: TableBucket) -> Result` | Full current-state scan of one primary-key bucket | +| `async fn create_kv_scanner(self) -> Result` | Full current-state scan of every bucket, including all current partitions | ## `LogScanner` @@ -254,6 +257,46 @@ server-deduplicated state); yields a single batch of at most `n` rows. | `async fn collect_all_batches(&mut self) -> Result>` | Drain into all batches | | `fn bucket(&self) -> &TableBucket` | The scanned bucket | +## `KvBatchScanner` + +Full current-state scan of one primary-key bucket via the `ScanKv` RPC. It +returns one merged row per live primary key rather than changelog events. +Projection is applied client-side. `limit` and filter pushdown are rejected. + +The server opens a RocksDB snapshot for the bucket on the first read and keeps +that snapshot across continuation requests. At most one request is in flight. +A polling timeout does not cancel or resend that request, so the scanner can be +polled again without skipping or duplicating rows. + +| Method | Description | +|------------------------------------------------------------------------------------|-------------| +| `async fn next_batch(&mut self) -> Result>` | Wait for the next non-empty batch, or `None` at EOF | +| `async fn next_batch_with_timeout(&mut self, timeout: Duration) -> Result` | Return a batch, timeout, or completion while preserving the in-flight request | +| `async fn collect_all_batches(&mut self) -> Result>` | Drain the bucket snapshot | +| `async fn close(&mut self) -> Result<()>` | Best-effort close of an unfinished server scanner; later reads fail unless it was already drained | +| `fn bucket(&self) -> &TableBucket` | The scanned bucket | +| `fn snapshot_log_offset(&self) -> Option` | Log high-watermark captured when the snapshot opened | + +`KvBatchReadOutcome` is `Batch(ScanBatch)`, `TimedOut`, or `Finished`. +Continuation failures are terminal because blindly retrying after the server +may have advanced its cursor could skip or duplicate rows. Start a new scanner +to restart the bucket. + +## `KvSnapshotScanner` + +Whole-table primary-key scan. It captures the current partition list when +created, then scans every `(partition, bucket)` sequentially. Only one +server-side scanner is open at a time, avoiding one pinned RocksDB snapshot per +bucket. Each bucket receives its own snapshot when that bucket is opened; this +is not a single atomic snapshot across the whole table. + +| Method | Description | +|------------------------------------------------------------------------------------|-------------| +| `async fn next_batch(&mut self) -> Result>` | Read the next batch across all buckets | +| `async fn next_batch_with_timeout(&mut self, timeout: Duration) -> Result` | Poll the current bucket without losing its in-flight request | +| `async fn collect_all_batches(&mut self) -> Result>` | Drain all captured partitions and buckets | +| `async fn close(&mut self) -> Result<()>` | Close the active bucket and discard unopened buckets; later reads fail unless all buckets were drained | + ## `ScanRecord` | Method | Description | diff --git a/fluss-rust/website/docs/user-guide/rust/example/primary-key-tables.md b/fluss-rust/website/docs/user-guide/rust/example/primary-key-tables.md index 5645ff8f123..617ddcba87a 100644 --- a/fluss-rust/website/docs/user-guide/rust/example/primary-key-tables.md +++ b/fluss-rust/website/docs/user-guide/rust/example/primary-key-tables.md @@ -167,6 +167,34 @@ start, `subscribe` at a specific offset instead of `EARLIEST_OFFSET`. To fetch all rows sharing a common primary-key prefix (by choosing a bucket key that's a strict prefix of the primary key), see [Prefix Lookup](./prefix-lookup.md). +## Full KV Scan + +To scan the merged current view of a primary-key table without supplying keys, +create a KV scanner. The whole-table scanner captures the current partition +list, then visits every `(partition, bucket)` sequentially. + +```rust +let mut scanner = table.new_scan().create_kv_scanner().await?; + +while let Some(batch) = scanner.next_batch().await? { + println!( + "bucket={} rows={}", + batch.bucket(), + batch.batch().num_rows() + ); +} +``` + +Use `project` or `project_by_name` before `create_kv_scanner` to return only +selected columns. Projection is performed client-side. Filter pushdown and +`limit` are not supported by full KV scans. + +Each bucket is read from a server-side snapshot that remains stable across +continuation RPCs. Buckets are opened lazily, so a whole-table scan is not one +atomic point-in-time snapshot across every bucket. To scan one explicit bucket, +including a bucket in a partitioned table, use `create_bucket_kv_scanner` with +the corresponding `TableBucket`. + ## Limit Scan To read up to `n` rows of a bucket's current state without supplying keys, use a batch scanner. The server returns the deduplicated current rows as Arrow batches, which is convenient for previews or DataFusion sources.