Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
190 changes: 190 additions & 0 deletions crates/integrations/datafusion/src/system_tables/consumers.rs
Original file line number Diff line number Diff line change
@@ -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<Arc<dyn TableProvider>> {
Ok(Arc::new(ConsumersTable { table }))
}

fn consumers_schema() -> SchemaRef {
static SCHEMA: OnceLock<SchemaRef> = 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<usize>>,
filters: &[Expr],
_limit: Option<usize>,
) -> DFResult<Arc<dyn ExecutionPlan>> {
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<Vec<TableProviderFilterPushDown>> {
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<Vec<String>> {
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::<Vec<_>>();
ids.sort_unstable();
Some(ids)
}

fn consumer_ids_from_filter(filter: &Expr) -> Option<HashSet<String>> {
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<String> {
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<String> {
match expr {
Expr::Literal(
ScalarValue::Utf8(Some(value))
| ScalarValue::LargeUtf8(Some(value))
| ScalarValue::Utf8View(Some(value)),
_,
) => Some(value.clone()),
_ => None,
}
}
6 changes: 6 additions & 0 deletions crates/integrations/datafusion/src/system_tables/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ use paimon::table::Table;
use crate::error::to_datafusion_error;

mod branches;
mod consumers;
mod files;
mod manifests;
mod options;
Expand All @@ -49,6 +50,7 @@ type Builder = fn(Table) -> DFResult<Arc<dyn TableProvider>>;
// metadata via `Catalog::list_partitions`).
const TABLES: &[(&str, Builder)] = &[
("branches", branches::build),
("consumers", consumers::build),
("files", files::build),
("manifests", manifests::build),
("options", options::build),
Expand All @@ -62,6 +64,7 @@ const TABLES: &[(&str, Builder)] = &[

const SYSTEM_TABLE_NAMES: &[&str] = &[
"branches",
"consumers",
"files",
"manifests",
"options",
Expand Down Expand Up @@ -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"));
Expand Down
87 changes: 87 additions & 0 deletions crates/integrations/datafusion/tests/system_tables.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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::<usize>(),
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::<StringArray>()
.expect("consumer_id is Utf8");
let snapshots = batch
.column(1)
.as_any()
.downcast_ref::<Int64Array>()
.expect("next_snapshot_id is Int64");
assert_eq!(ids.iter().flatten().collect::<Vec<_>>(), 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::<StringArray>()
.unwrap();
let snapshots = batch
.column(1)
.as_any()
.downcast_ref::<Int64Array>()
.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::<Vec<_>>()
);
}
}

#[tokio::test]
async fn test_branches_system_table_empty_when_no_branch_dir() {
let (ctx, _catalog, _tmp) = create_context().await;
Expand Down
11 changes: 10 additions & 1 deletion crates/paimon/src/io/storage.rs
Original file line number Diff line number Diff line change
Expand Up @@ -314,7 +314,7 @@ impl Storage {
fn fs_relative_path(path: &str) -> crate::Result<Cow<'_, str>> {
// 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)
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading