diff --git a/crates/integrations/datafusion/src/system_tables/consumers.rs b/crates/integrations/datafusion/src/system_tables/consumers.rs new file mode 100644 index 000000000..b8a33d4d5 --- /dev/null +++ b/crates/integrations/datafusion/src/system_tables/consumers.rs @@ -0,0 +1,190 @@ +// 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. + +//! Mirrors Java [ConsumersTable](https://github.com/apache/paimon/blob/release-1.3/paimon-core/src/main/java/org/apache/paimon/table/system/ConsumersTable.java). + +use std::collections::HashSet; +use std::sync::{Arc, OnceLock}; + +use async_trait::async_trait; +use datafusion::arrow::array::{Int64Array, RecordBatch, StringArray}; +use datafusion::arrow::datatypes::{DataType, Field, Schema, SchemaRef}; +use datafusion::catalog::Session; +use datafusion::common::ScalarValue; +use datafusion::datasource::memory::MemorySourceConfig; +use datafusion::datasource::{TableProvider, TableType}; +use datafusion::error::Result as DFResult; +use datafusion::logical_expr::{Expr, Operator, TableProviderFilterPushDown}; +use datafusion::physical_plan::ExecutionPlan; +use paimon::table::Table; + +use crate::error::to_datafusion_error; + +pub(super) fn build(table: Table) -> DFResult> { + Ok(Arc::new(ConsumersTable { table })) +} + +fn consumers_schema() -> SchemaRef { + static SCHEMA: OnceLock = OnceLock::new(); + SCHEMA + .get_or_init(|| { + Arc::new(Schema::new(vec![ + Field::new("consumer_id", DataType::Utf8, false), + Field::new("next_snapshot_id", DataType::Int64, false), + ])) + }) + .clone() +} + +#[derive(Debug)] +struct ConsumersTable { + table: Table, +} + +#[async_trait] +impl TableProvider for ConsumersTable { + fn schema(&self) -> SchemaRef { + consumers_schema() + } + + fn table_type(&self) -> TableType { + TableType::View + } + + async fn scan( + &self, + _state: &dyn Session, + projection: Option<&Vec>, + filters: &[Expr], + _limit: Option, + ) -> DFResult> { + let manager = self.table.consumer_manager(); + let requested_ids = requested_consumer_ids(filters); + let consumers = crate::runtime::await_with_runtime(async move { + match requested_ids { + Some(ids) => manager.list_by_ids(&ids).await, + None => manager.list_all().await, + } + }) + .await + .map_err(to_datafusion_error)?; + + let schema = consumers_schema(); + let batch = RecordBatch::try_new( + schema.clone(), + vec![ + Arc::new(StringArray::from_iter_values( + consumers.iter().map(|(id, _)| id.as_str()), + )), + Arc::new(Int64Array::from_iter_values( + consumers.iter().map(|(_, next_snapshot)| *next_snapshot), + )), + ], + )?; + + Ok(MemorySourceConfig::try_new_exec( + &[vec![batch]], + schema, + projection.cloned(), + )?) + } + + fn supports_filters_pushdown( + &self, + filters: &[&Expr], + ) -> DFResult> { + Ok(filters + .iter() + .map(|filter| { + if consumer_ids_from_filter(filter).is_some() { + TableProviderFilterPushDown::Inexact + } else { + TableProviderFilterPushDown::Unsupported + } + }) + .collect()) + } +} + +fn requested_consumer_ids(filters: &[Expr]) -> Option> { + let ids = filters + .iter() + .filter_map(consumer_ids_from_filter) + .reduce(|mut left, right| { + left.retain(|id| right.contains(id)); + left + })?; + let mut ids = ids.into_iter().collect::>(); + ids.sort_unstable(); + Some(ids) +} + +fn consumer_ids_from_filter(filter: &Expr) -> Option> { + match filter { + Expr::BinaryExpr(binary) if binary.op == Operator::Eq => { + consumer_id_literal(binary.left.as_ref(), binary.right.as_ref()) + .or_else(|| consumer_id_literal(binary.right.as_ref(), binary.left.as_ref())) + .map(|id| HashSet::from([id])) + } + Expr::BinaryExpr(binary) if binary.op == Operator::And => { + match ( + consumer_ids_from_filter(binary.left.as_ref()), + consumer_ids_from_filter(binary.right.as_ref()), + ) { + (Some(mut left), Some(right)) => { + left.retain(|id| right.contains(id)); + Some(left) + } + (Some(ids), None) | (None, Some(ids)) => Some(ids), + (None, None) => None, + } + } + Expr::BinaryExpr(binary) if binary.op == Operator::Or => { + let mut left = consumer_ids_from_filter(binary.left.as_ref())?; + left.extend(consumer_ids_from_filter(binary.right.as_ref())?); + Some(left) + } + Expr::InList(in_list) + if !in_list.negated && is_consumer_id_column(in_list.expr.as_ref()) => + { + in_list.list.iter().map(string_literal).collect() + } + _ => None, + } +} + +fn consumer_id_literal(column: &Expr, literal: &Expr) -> Option { + is_consumer_id_column(column) + .then_some(literal) + .and_then(string_literal) +} + +fn is_consumer_id_column(expr: &Expr) -> bool { + matches!(expr, Expr::Column(column) if column.name == "consumer_id") +} + +fn string_literal(expr: &Expr) -> Option { + match expr { + Expr::Literal( + ScalarValue::Utf8(Some(value)) + | ScalarValue::LargeUtf8(Some(value)) + | ScalarValue::Utf8View(Some(value)), + _, + ) => Some(value.clone()), + _ => None, + } +} diff --git a/crates/integrations/datafusion/src/system_tables/mod.rs b/crates/integrations/datafusion/src/system_tables/mod.rs index 392db7da5..3b3a0b098 100644 --- a/crates/integrations/datafusion/src/system_tables/mod.rs +++ b/crates/integrations/datafusion/src/system_tables/mod.rs @@ -30,6 +30,7 @@ use paimon::table::Table; use crate::error::to_datafusion_error; mod branches; +mod consumers; mod files; mod manifests; mod options; @@ -49,6 +50,7 @@ type Builder = fn(Table) -> DFResult>; // metadata via `Catalog::list_partitions`). const TABLES: &[(&str, Builder)] = &[ ("branches", branches::build), + ("consumers", consumers::build), ("files", files::build), ("manifests", manifests::build), ("options", options::build), @@ -62,6 +64,7 @@ const TABLES: &[(&str, Builder)] = &[ const SYSTEM_TABLE_NAMES: &[&str] = &[ "branches", + "consumers", "files", "manifests", "options", @@ -195,6 +198,9 @@ mod tests { assert!(is_registered("branches")); assert!(is_registered("Branches")); assert!(is_registered("BRANCHES")); + assert!(is_registered("consumers")); + assert!(is_registered("Consumers")); + assert!(is_registered("CONSUMERS")); assert!(is_registered("files")); assert!(is_registered("Files")); assert!(is_registered("FILES")); diff --git a/crates/integrations/datafusion/tests/system_tables.rs b/crates/integrations/datafusion/tests/system_tables.rs index 582d6ef10..1a64ccb77 100644 --- a/crates/integrations/datafusion/tests/system_tables.rs +++ b/crates/integrations/datafusion/tests/system_tables.rs @@ -530,6 +530,93 @@ async fn test_snapshots_system_table() { ); } +#[tokio::test] +async fn test_consumers_system_table() { + let (ctx, _catalog, tmp) = create_context().await; + let table = format!("paimon.default.{FIXTURE_TABLE}$consumers"); + + let batches = run_sql(&ctx, &format!("SELECT * FROM {table}")).await; + assert!(!batches.is_empty(), "$consumers should return a batch"); + let schema = batches[0].schema(); + assert_eq!(schema.field(0).name(), "consumer_id"); + assert_eq!(schema.field(0).data_type(), &DataType::Utf8); + assert_eq!(schema.field(1).name(), "next_snapshot_id"); + assert_eq!(schema.field(1).data_type(), &DataType::Int64); + assert_eq!( + batches.iter().map(|batch| batch.num_rows()).sum::(), + 0 + ); + + let consumer_dir = tmp + .path() + .join("default.db") + .join(FIXTURE_TABLE) + .join("consumer"); + std::fs::create_dir_all(&consumer_dir).unwrap(); + std::fs::write(consumer_dir.join("consumer-id2"), r#"{"nextSnapshot":6}"#).unwrap(); + std::fs::write( + consumer_dir.join("consumer-id1"), + r#"{"nextSnapshot":5,"ignored":"value"}"#, + ) + .unwrap(); + std::fs::write(consumer_dir.join("not-a-consumer"), "not json").unwrap(); + + let batches = run_sql(&ctx, &format!("SELECT * FROM {table} ORDER BY consumer_id")).await; + let batch = &batches[0]; + let ids = batch + .column(0) + .as_any() + .downcast_ref::() + .expect("consumer_id is Utf8"); + let snapshots = batch + .column(1) + .as_any() + .downcast_ref::() + .expect("next_snapshot_id is Int64"); + assert_eq!(ids.iter().flatten().collect::>(), vec!["id1", "id2"]); + assert_eq!(snapshots.values(), &[5, 6]); + + // A filtered point lookup must not read unrelated consumer files. + std::fs::write(consumer_dir.join("consumer-bad"), "{").unwrap(); + let cases: [(&str, &[(&str, i64)]); 5] = [ + ("consumer_id = 'id1'", &[("id1", 5)]), + ("consumer_id IN ('id2', 'missing')", &[("id2", 6)]), + ("consumer_id = 'missing'", &[]), + ("consumer_id = 'id1 '", &[]), + ("consumer_id IN ('id1', 'id1 ')", &[("id1", 5)]), + ]; + for (predicate, expected) in cases { + let batches = run_sql( + &ctx, + &format!("SELECT consumer_id, next_snapshot_id FROM {table} WHERE {predicate}"), + ) + .await; + let mut actual = Vec::new(); + for batch in batches { + let ids = batch + .column(0) + .as_any() + .downcast_ref::() + .unwrap(); + let snapshots = batch + .column(1) + .as_any() + .downcast_ref::() + .unwrap(); + for row in 0..batch.num_rows() { + actual.push((ids.value(row).to_string(), snapshots.value(row))); + } + } + assert_eq!( + actual, + expected + .iter() + .map(|(id, snapshot)| (id.to_string(), *snapshot)) + .collect::>() + ); + } +} + #[tokio::test] async fn test_branches_system_table_empty_when_no_branch_dir() { let (ctx, _catalog, _tmp) = create_context().await; diff --git a/crates/paimon/src/io/storage.rs b/crates/paimon/src/io/storage.rs index 5ba31e715..24885d290 100644 --- a/crates/paimon/src/io/storage.rs +++ b/crates/paimon/src/io/storage.rs @@ -314,7 +314,7 @@ impl Storage { fn fs_relative_path(path: &str) -> crate::Result> { // A `file://` / `file:/` URL is already in scheme-relative form. if let Some(stripped) = path.strip_prefix("file:/") { - return Ok(if stripped.contains('\\') { + return Ok(if cfg!(windows) && stripped.contains('\\') { Cow::Owned(stripped.replace('\\', "/")) } else { Cow::Borrowed(stripped) @@ -561,6 +561,15 @@ mod fs_relative_path_tests { assert_eq!(rel("file:///tmp/wh"), "//tmp/wh"); } + #[cfg(not(windows))] + #[test] + fn file_scheme_preserves_posix_backslashes() { + assert_eq!( + rel(r"file:///tmp/consumer/consumer-id\part"), + r"//tmp/consumer/consumer-id\part" + ); + } + #[test] fn windows_drive_path_keeps_drive_and_normalizes_separators() { // The historical bug dropped the drive letter (`C:\wh` -> `:\wh`); we diff --git a/crates/paimon/src/table/consumer_manager.rs b/crates/paimon/src/table/consumer_manager.rs new file mode 100644 index 000000000..35cb055b9 --- /dev/null +++ b/crates/paimon/src/table/consumer_manager.rs @@ -0,0 +1,209 @@ +// 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. + +//! Consumer progress manager for Java-compatible consumer files. + +use std::collections::HashSet; +use std::time::Duration; + +use crate::io::FileIO; +use futures::{stream, StreamExt, TryStreamExt}; +use serde::Deserialize; + +const CONSUMER_DIR: &str = "consumer"; +const CONSUMER_PREFIX: &str = "consumer-"; +const CONSUMER_READ_CONCURRENCY: usize = 32; +const READ_RETRIES: usize = 10; +const RETRY_DELAY: Duration = Duration::from_millis(200); + +/// Reads consumer progress files stored under a table or branch. +#[derive(Debug, Clone)] +pub struct ConsumerManager { + file_io: FileIO, + table_path: String, +} + +impl ConsumerManager { + pub fn new(file_io: FileIO, table_path: String) -> Self { + Self { + file_io, + table_path, + } + } + + /// Create a manager for a branch of this table. + pub fn with_branch(&self, branch_name: &str) -> Self { + Self::new( + self.file_io.clone(), + format!("{}/branch/branch-{branch_name}", self.table_path), + ) + } + + /// Read one consumer's next snapshot id. + pub async fn get(&self, consumer_id: &str) -> crate::Result> { + let ids = HashSet::from([consumer_id]); + Ok(self + .list_matching(Some(&ids)) + .await? + .into_iter() + .next() + .map(|(_, next_snapshot)| next_snapshot)) + } + + async fn read_path(&self, consumer_id: &str, path: &str) -> crate::Result> { + let input = self.file_io.new_input(path)?; + for attempt in 0..READ_RETRIES { + let bytes = match input.read().await { + Ok(bytes) => bytes, + Err(crate::Error::IoUnexpected { ref source, .. }) + if source.kind() == opendal::ErrorKind::NotFound => + { + return Ok(None); + } + Err(error) => return Err(error), + }; + match serde_json::from_slice::(&bytes) { + Ok(consumer) => return Ok(Some(consumer.next_snapshot)), + Err(_) if attempt + 1 < READ_RETRIES => tokio::time::sleep(RETRY_DELAY).await, + Err(error) => { + return Err(crate::Error::DataInvalid { + message: format!("consumer '{consumer_id}' JSON invalid: {error}"), + source: Some(Box::new(error)), + }); + } + } + } + unreachable!("READ_RETRIES is non-zero") + } + + /// List consumer ids and their next snapshot ids, sorted by consumer id. + pub async fn list_all(&self) -> crate::Result> { + self.list_matching(None).await + } + + /// List the requested consumers that have exactly matching directory entries. + pub async fn list_by_ids(&self, consumer_ids: &[String]) -> crate::Result> { + let ids = consumer_ids.iter().map(String::as_str).collect(); + self.list_matching(Some(&ids)).await + } + + async fn list_matching( + &self, + consumer_ids: Option<&HashSet<&str>>, + ) -> crate::Result> { + let directory = format!("{}/{}", self.table_path, CONSUMER_DIR); + let statuses = match self.file_io.list_status(&directory).await { + Ok(statuses) => statuses, + Err(crate::Error::IoUnexpected { ref source, .. }) + if source.kind() == opendal::ErrorKind::NotFound => + { + return Ok(Vec::new()); + } + Err(error) => return Err(error), + }; + + let mut consumers = stream::iter(statuses.into_iter().filter_map(|status| { + if status.is_dir { + return None; + } + let name = status.path.rsplit('/').next().unwrap_or(&status.path); + let id = name.strip_prefix(CONSUMER_PREFIX)?; + if consumer_ids.is_some_and(|ids| !ids.contains(id)) { + return None; + } + Some((id.to_string(), status.path)) + })) + .map(|(id, path)| async move { + self.read_path(&id, &path) + .await + .map(|next_snapshot| next_snapshot.map(|next_snapshot| (id, next_snapshot))) + }) + .buffer_unordered(CONSUMER_READ_CONCURRENCY) + .try_collect::>() + .await? + .into_iter() + .flatten() + .collect::>(); + consumers.sort_unstable_by(|a, b| a.0.cmp(&b.0)); + Ok(consumers) + } +} + +#[derive(Deserialize)] +struct Consumer { + #[serde(rename = "nextSnapshot")] + next_snapshot: i64, +} + +#[cfg(test)] +mod tests { + use bytes::Bytes; + + use super::*; + use crate::io::FileIOBuilder; + + #[tokio::test] + async fn retries_a_consumer_being_overwritten() { + let file_io = FileIOBuilder::new("memory").build().unwrap(); + let table_path = "memory:/consumer-retry"; + let directory = format!("{table_path}/{CONSUMER_DIR}"); + let path = format!("{directory}/{CONSUMER_PREFIX}job"); + file_io.mkdirs(&directory).await.unwrap(); + file_io + .new_output(&path) + .unwrap() + .write(Bytes::from_static(b"{")) + .await + .unwrap(); + + let manager = ConsumerManager::new(file_io.clone(), table_path.to_string()); + let mut read = Box::pin(manager.get("job")); + tokio::select! { + result = &mut read => panic!("invalid JSON returned before retry: {result:?}"), + _ = tokio::time::sleep(Duration::from_millis(50)) => {} + } + file_io + .new_output(&path) + .unwrap() + .write(Bytes::from_static(br#"{"nextSnapshot":5}"#)) + .await + .unwrap(); + + assert_eq!(read.await.unwrap(), Some(5)); + } + + #[cfg(not(windows))] + #[tokio::test] + async fn lists_posix_consumer_id_with_backslash() { + let tmp = tempfile::tempdir().unwrap(); + let table_path = tmp.path().join("table"); + let directory = table_path.join(CONSUMER_DIR); + std::fs::create_dir_all(&directory).unwrap(); + std::fs::write(directory.join(r"consumer-id\part"), r#"{"nextSnapshot":5}"#).unwrap(); + + let manager = ConsumerManager::new( + FileIOBuilder::new("file").build().unwrap(), + format!("file://{}", table_path.display()), + ); + + assert_eq!( + manager.list_all().await.unwrap(), + vec![(r"id\part".to_string(), 5)] + ); + assert_eq!(manager.get(r"id\part").await.unwrap(), Some(5)); + } +} diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index 92c73a6fb..32484924a 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -33,6 +33,7 @@ mod bucket_assigner_fixed; mod bucket_filter; mod bucket_function; mod commit_message; +mod consumer_manager; pub(crate) mod cow_writer; mod data_evolution_reader; pub mod data_evolution_writer; @@ -113,6 +114,7 @@ pub use audit_log_table::AuditLogTable; pub use blob_resolver::{BlobReader, BlobStream}; pub use branch_manager::BranchManager; pub use commit_message::CommitMessage; +pub use consumer_manager::ConsumerManager; pub use cow_writer::{CopyOnWriteMergeWriter, FileInfo}; pub use data_evolution_writer::{DataEvolutionDeleteWriter, DataEvolutionWriter}; #[cfg(feature = "fulltext")] @@ -322,6 +324,15 @@ impl Table { } } + pub fn consumer_manager(&self) -> ConsumerManager { + let manager = ConsumerManager::new(self.file_io.clone(), self.location.clone()); + if self.is_main_branch() { + manager + } else { + manager.with_branch(&self.branch) + } + } + /// Get the REST environment, if this table was loaded from a REST catalog. pub fn rest_env(&self) -> Option<&RESTEnv> { self.rest_env.as_ref() diff --git a/docs/src/sql.md b/docs/src/sql.md index c53bacbf5..c37ac51cd 100644 --- a/docs/src/sql.md +++ b/docs/src/sql.md @@ -1922,6 +1922,21 @@ Columns: | `watermark` | BIGINT | Watermark | | `next_row_id` | BIGINT | Next row id | +### $consumers + +View the stored progress of streaming consumers: + +```sql +SELECT * FROM paimon.default.my_table$consumers; +``` + +Columns: + +| Column | Type | Description | +|---|---|---| +| `consumer_id` | STRING | Consumer identifier | +| `next_snapshot_id` | BIGINT | Next snapshot the consumer will read | + ### $tags View all named tags of a table: