Skip to content
Draft
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
2 changes: 1 addition & 1 deletion acceptance/src/data.rs
Original file line number Diff line number Diff line change
Expand Up @@ -115,7 +115,7 @@ pub async fn assert_scan_metadata(
test_case: &TestCaseInfo,
) -> TestResult<()> {
let table_root = test_case.table_root()?;
let snapshot = Snapshot::builder_for(table_root).build(engine.as_ref())?;
let snapshot = Snapshot::builder_for(table_root).build(engine.as_ref()).await?;
let scan = snapshot.scan_builder().build()?;
let mut schema = None;
let batches: Vec<RecordBatch> = scan
Expand Down
4 changes: 2 additions & 2 deletions acceptance/src/meta.rs
Original file line number Diff line number Diff line change
Expand Up @@ -103,13 +103,13 @@ impl TestCaseInfo {
let engine = engine.as_ref();
let (latest, versions) = self.versions().await?;

let snapshot = Snapshot::builder_for(self.table_root()?).build(engine)?;
let snapshot = Snapshot::builder_for(self.table_root()?).build(engine).await?;
self.assert_snapshot_meta(&latest, &snapshot)?;

for table_version in versions {
let snapshot = Snapshot::builder_for(self.table_root()?)
.at_version(table_version.version)
.build(engine)?;
.build(engine).await?;
self.assert_snapshot_meta(&table_version, &snapshot)?;
}

Expand Down
1 change: 1 addition & 0 deletions ffi/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@ release = false
crate-type = ["lib", "cdylib", "staticlib"]

[dependencies]
futures = "0.3"
tracing = "0.1"
tracing-core = { version = "0.1", optional = true }
tracing-subscriber = { version = "0.3", optional = true, features = [ "json" ] }
Expand Down
2 changes: 1 addition & 1 deletion ffi/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -608,7 +608,7 @@ fn snapshot_impl(
} else {
builder
};
let snapshot = builder.build(extern_engine.engine().as_ref())?;
let snapshot = futures::executor::block_on(builder.build(extern_engine.engine().as_ref()))?;
Ok(snapshot.into())
}

Expand Down
6 changes: 4 additions & 2 deletions ffi/src/transaction/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,9 @@ fn transaction_impl(
url: DeltaResult<Url>,
extern_engine: &dyn ExternEngine,
) -> DeltaResult<Handle<ExclusiveTransaction>> {
let snapshot = Snapshot::builder_for(url?).build(extern_engine.engine().as_ref())?;
let snapshot = futures::executor::block_on(
Snapshot::builder_for(url?).build(extern_engine.engine().as_ref())
)?;
let transaction = snapshot.transaction();
Ok(Box::new(transaction?).into())
}
Expand Down Expand Up @@ -372,7 +374,7 @@ mod tests {

// Confirm that the data matches what we appended
let test_batch = ArrowEngineData::from(batch);
test_read(&test_batch, &table_url, unsafe { engine.as_ref().engine() })?;
test_read(&test_batch, &table_url, unsafe { engine.as_ref().engine() }).await?;

unsafe { free_schema(write_schema) };
unsafe { free_write_context(write_context) };
Expand Down
3 changes: 1 addition & 2 deletions kernel/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ pre-release-hook = [
delta_kernel_derive = { path = "../derive-macros", version = "0.16.0" }
bytes = "1.10"
chrono = "0.4.41"
futures = "0.3"
indexmap = "2.10.0"
itertools = "0.14"
roaring = "0.11.2"
Expand All @@ -56,7 +57,6 @@ uuid = { version = "1.18.0", features = ["v4", "fast-rng"] }
z85 = "3.0.6"

# optional deps
futures = { version = "0.3", optional = true }
# Used for fetching direct urls (like pre-signed urls)
reqwest = { version = "0.12.23", default-features = false, optional = true }
# optionally used with default engine (though not required)
Expand Down Expand Up @@ -116,7 +116,6 @@ catalog-managed = []
default-engine-base = [
"arrow-conversion",
"arrow-expression",
"futures",
"need-arrow",
"tokio",
]
Expand Down
1 change: 1 addition & 0 deletions kernel/examples/inspect-table/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ delta_kernel = { path = "../../../kernel", features = [
"internal-api",
] }
env_logger = "0.11.8"
tokio = { version = "1", features = ["rt", "macros"] }

# for cargo-release
[package.metadata.release]
Expand Down
7 changes: 4 additions & 3 deletions kernel/examples/inspect-table/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -51,7 +51,8 @@ enum Commands {

fn main() -> ExitCode {
env_logger::init();
match try_main() {
let runtime = tokio::runtime::Runtime::new().unwrap();
match runtime.block_on(try_main()) {
Ok(()) => ExitCode::SUCCESS,
Err(e) => {
println!("{e:#?}");
Expand Down Expand Up @@ -177,12 +178,12 @@ fn print_scan_file(
);
}

fn try_main() -> DeltaResult<()> {
async fn try_main() -> DeltaResult<()> {
let cli = Cli::parse();

let url = delta_kernel::try_parse_uri(&cli.location_args.path)?;
let engine = common::get_engine(&url, &cli.location_args)?;
let snapshot = Snapshot::builder_for(url).build(&engine)?;
let snapshot = Snapshot::builder_for(url).build(&engine).await?;

match cli.command {
Commands::TableVersion => {
Expand Down
2 changes: 2 additions & 0 deletions kernel/examples/read-table-changes/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -16,4 +16,6 @@ delta_kernel = { path = "../../../kernel", features = [
"default-engine-rustls",
"internal-api",
] }
futures = "0.3"
itertools = "0.14"
tokio = { version = "1", features = ["rt", "macros"] }
10 changes: 6 additions & 4 deletions kernel/examples/read-table-changes/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@ use delta_kernel::arrow::{compute::filter_record_batch, util::pretty::print_batc
use delta_kernel::engine::arrow_data::ArrowEngineData;
use delta_kernel::table_changes::TableChanges;
use delta_kernel::DeltaResult;
use futures::{StreamExt, TryStreamExt};
use itertools::Itertools;

#[derive(Parser)]
Expand All @@ -24,15 +25,16 @@ struct Cli {
end_version: Option<u64>,
}

fn main() -> DeltaResult<()> {
#[tokio::main]
async fn main() -> DeltaResult<()> {
let cli = Cli::parse();
let url = delta_kernel::try_parse_uri(cli.location_args.path.as_str())?;
let engine = common::get_engine(&url, &cli.location_args)?;
let table_changes = TableChanges::try_new(url, &engine, cli.start_version, cli.end_version)?;
let table_changes = TableChanges::try_new(url, &engine, cli.start_version, cli.end_version).await?;

let table_changes_scan = table_changes.into_scan_builder().build()?;
let batches: Vec<RecordBatch> = table_changes_scan
.execute(Arc::new(engine))?
.execute(Arc::new(engine)).await?
.map(|scan_result| -> DeltaResult<_> {
let scan_result = scan_result?;
let mask = scan_result.full_mask();
Expand All @@ -48,7 +50,7 @@ fn main() -> DeltaResult<()> {
Ok(record_batch)
}
})
.try_collect()?;
.try_collect().await?;
print_batches(&batches)?;
Ok(())
}
1 change: 1 addition & 0 deletions kernel/examples/read-table-multi-threaded/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ delta_kernel = { path = "../../../kernel", features = [
] }
env_logger = "0.11.8"
spmc = "0.3.0"
tokio = { version = "1", features = ["rt", "macros"] }
url = "2"

# for cargo-release
Expand Down
7 changes: 4 additions & 3 deletions kernel/examples/read-table-multi-threaded/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,8 @@ struct Cli {

fn main() -> ExitCode {
env_logger::init();
match try_main() {
let runtime = tokio::runtime::Runtime::new().unwrap();
match runtime.block_on(try_main()) {
Ok(()) => ExitCode::SUCCESS,
Err(e) => {
println!("{e:#?}");
Expand Down Expand Up @@ -93,13 +94,13 @@ struct ScanState {
logical_schema: SchemaRef,
}

fn try_main() -> DeltaResult<()> {
async fn try_main() -> DeltaResult<()> {
let cli = Cli::parse();

let url = delta_kernel::try_parse_uri(&cli.location_args.path)?;
println!("Reading {url}");
let engine = common::get_engine(&url, &cli.location_args)?;
let snapshot = Snapshot::builder_for(url).build(&engine)?;
let snapshot = Snapshot::builder_for(url).build(&engine).await?;
let Some(scan) = common::get_scan(snapshot, &cli.scan_args)? else {
return Ok(());
};
Expand Down
1 change: 1 addition & 0 deletions kernel/examples/read-table-single-threaded/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ delta_kernel = { path = "../../../kernel", features = [
] }
env_logger = "0.11.8"
itertools = "0.14"
tokio = { version = "1", features = ["rt", "macros"] }

# for cargo-release
[package.metadata.release]
Expand Down
7 changes: 4 additions & 3 deletions kernel/examples/read-table-single-threaded/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,8 @@ struct Cli {

fn main() -> ExitCode {
env_logger::init();
match try_main() {
let runtime = tokio::runtime::Runtime::new().unwrap();
match runtime.block_on(try_main()) {
Ok(()) => ExitCode::SUCCESS,
Err(e) => {
println!("{e:#?}");
Expand All @@ -38,12 +39,12 @@ fn main() -> ExitCode {
}
}

fn try_main() -> DeltaResult<()> {
async fn try_main() -> DeltaResult<()> {
let cli = Cli::parse();
let url = delta_kernel::try_parse_uri(&cli.location_args.path)?;
println!("Reading {url}");
let engine = common::get_engine(&url, &cli.location_args)?;
let snapshot = Snapshot::builder_for(url).build(&engine)?;
let snapshot = Snapshot::builder_for(url).build(&engine).await?;
let Some(scan) = common::get_scan(snapshot, &cli.scan_args)? else {
return Ok(());
};
Expand Down
6 changes: 3 additions & 3 deletions kernel/examples/write-table/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -125,7 +125,7 @@ async fn create_or_get_base_snapshot(
schema_str: &str,
) -> DeltaResult<SnapshotRef> {
// Check if table already exists
match Snapshot::builder_for(url.clone()).build(engine) {
match Snapshot::builder_for(url.clone()).build(engine).await {
Ok(snapshot) => {
println!("✓ Found existing table at version {}", snapshot.version());
Ok(snapshot)
Expand All @@ -135,7 +135,7 @@ async fn create_or_get_base_snapshot(
println!("Creating new Delta table...");
let schema = parse_schema(schema_str)?;
create_table(url, &schema).await?;
Snapshot::builder_for(url.clone()).build(engine)
Snapshot::builder_for(url.clone()).build(engine).await
}
}
}
Expand Down Expand Up @@ -294,7 +294,7 @@ async fn read_and_display_data(
table_url: &Url,
engine: DefaultEngine<TokioBackgroundExecutor>,
) -> DeltaResult<()> {
let snapshot = Snapshot::builder_for(table_url.clone()).build(&engine)?;
let snapshot = Snapshot::builder_for(table_url.clone()).build(&engine).await?;
let scan = snapshot.scan_builder().build()?;

let batches: Vec<RecordBatch> = scan
Expand Down
26 changes: 13 additions & 13 deletions kernel/src/actions/set_transaction.rs
Original file line number Diff line number Diff line change
Expand Up @@ -105,15 +105,15 @@ mod tests {
use crate::arrow::array::StringArray;
use itertools::Itertools;

fn get_latest_transactions(
async fn get_latest_transactions(
path: &str,
app_id: &str,
) -> (SetTransactionMap, Option<SetTransaction>) {
let path = std::fs::canonicalize(PathBuf::from(path)).unwrap();
let url = url::Url::from_directory_path(path).unwrap();
let engine = SyncEngine::new();

let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
let snapshot = Snapshot::builder_for(url).build(&engine).await.unwrap();
let log_segment = snapshot.log_segment();

(
Expand All @@ -122,13 +122,13 @@ mod tests {
)
}

#[test]
fn test_txn() {
let (txns, txn) = get_latest_transactions("./tests/data/basic_partitioned/", "test");
#[tokio::test]
async fn test_txn() {
let (txns, txn) = get_latest_transactions("./tests/data/basic_partitioned/", "test").await;
assert!(txn.is_none());
assert_eq!(txns.len(), 0);

let (txns, txn) = get_latest_transactions("./tests/data/app-txn-no-checkpoint/", "my-app");
let (txns, txn) = get_latest_transactions("./tests/data/app-txn-no-checkpoint/", "my-app").await;
assert!(txn.is_some());
assert_eq!(txns.len(), 2);
assert_eq!(txns.get("my-app"), txn.as_ref());
Expand All @@ -142,7 +142,7 @@ mod tests {
.as_ref()
);

let (txns, txn) = get_latest_transactions("./tests/data/app-txn-checkpoint/", "my-app");
let (txns, txn) = get_latest_transactions("./tests/data/app-txn-checkpoint/", "my-app").await;
assert!(txn.is_some());
assert_eq!(txns.len(), 2);
assert_eq!(txns.get("my-app"), txn.as_ref());
Expand All @@ -157,13 +157,13 @@ mod tests {
);
}

#[test]
fn test_replay_for_app_ids() {
#[tokio::test]
async fn test_replay_for_app_ids() {
let path = std::fs::canonicalize(PathBuf::from("./tests/data/parquet_row_group_skipping/"));
let url = url::Url::from_directory_path(path.unwrap()).unwrap();
let engine = SyncEngine::new();

let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
let snapshot = Snapshot::builder_for(url).build(&engine).await.unwrap();
let log_segment = snapshot.log_segment();

// The checkpoint has five parts, each containing one action. There are two app ids.
Expand All @@ -174,13 +174,13 @@ mod tests {
assert_eq!(data.len(), 2);
}

#[test]
fn test_txn_retention_filtering() {
#[tokio::test]
async fn test_txn_retention_filtering() {
let path = std::fs::canonicalize(PathBuf::from("./tests/data/app-txn-with-last-updated/"));
let url = url::Url::from_directory_path(path.unwrap()).unwrap();
let engine = SyncEngine::new();

let snapshot = Snapshot::builder_for(url).build(&engine).unwrap();
let snapshot = Snapshot::builder_for(url).build(&engine).await.unwrap();
let log_segment = snapshot.log_segment();

// Test with no retention (should get all transactions)
Expand Down
Loading
Loading