, ,

Delta Kernel (Rust) and the future of the Lakehouse.

I recently spent some time with my old love, that is, Rust; fawning over the old days when I learned to fight the borrow checker without my own two worn-out hands, no LLMs to do the work for me. I still think it’s probably one of the most fun languages to learn, for many reasons. Yes, it can be a bit verbose, but it also teaches you to think about your programs and problems in a different way. It’s fast, way fast, and simply brings joy. It’s also a great tool to write building blocks with, CLIs and other crates that can be consumed by others.

I’ve also been an avid user of the Lakehouse Architecture since before it was cool.

I cut my teeth on Kimball and the classic SQL Server Data Warehouse of yore, and I fought and bloodied myself over the many pitfalls of that technology and ideology. Nothing has done more to bring the beauty of distributed file system storage with ACID … ala … the Lakehouse, than Delta Lake. Of course, because of the zealots, we must mention Apache Iceberg, but Delta is where it all started … I’m sure within a few short years the two will merge into one.

One of the main pain points of running diverse workloads inside the Lakehouse using all the new tooling; DuckDB, Polars, Daft, Spark, etc.- is that the rate of change, say, feature additions, is way too fast for any one open-source project to keep up with. I have personally fought, in production, the pains of tools like Polars and DuckDB being able to work with various new, non-backward-compatible changes to Delta, for example, deletion vectors. It was a real pain and caused real problems for production pipelines.

This is where Delta Kernel comes into play.

It’s an important change and direction that is taking place, pushed by Databricks and the Delta team. It doesn’t get much news or attention, simply because it’s the building block on which other tools are built and the interface they use to talk the Delta Protocol with ease … it’s the new standard. In the same way that Datafusion always takes a backseat, besides odd mentions here and there, Delta Kernel is the power, the One Ring, behind many data tools today, and going forward into the future.

Read up on the link above, but at the end of the day, when you have multiple, separate tools reading the same complicated Delta Protocol, things are going to get messy.

The ACID-based complexity behind Delta Lake is real; there are checkpoints, data files, metadata logs, commits, etc.; these can all live in various remote cloud storage as well. It is not trivial at all to build first-class support for this sort of protocol when someone simply wants to SCAN and READ a Delta Table.

But you can do with ease now with tools like Rust and the Delta Kernel.

I played around with it a little, building my own CLI tool to Import / Export data back and forth from files to remote S3 based Delta Lake tables using this Rust Delta Kernel just to get an idea of how easy integration is to build with this tool.

use std::fs::File;
use std::path::{Path, PathBuf};
use std::sync::Arc;

use arrow_csv::reader::Format;
use arrow_csv::ReaderBuilder;
use clap::{Parser, Subcommand};

use delta_kernel::committer::FileSystemCommitter;
use delta_kernel::engine::arrow_conversion::TryIntoKernel;
use delta_kernel::engine::arrow_data::{ArrowEngineData, EngineDataArrowExt};
use delta_kernel::object_store::aws::AmazonS3Builder;
use delta_kernel::transaction::create_table::create_table;
use delta_kernel::transaction::{CommitResult, RetryableTransaction};
use delta_kernel::Snapshot;

use delta_kernel_default_engine::DefaultEngineBuilder;

use parquet::arrow::ArrowWriter;
use parquet::basic::Compression;
use parquet::file::properties::WriterProperties;

#[derive(Parser, Debug)]
#[command(
    name = "delta-import-export",
    version,
    about = "Import CSV files into Delta Lake and export Delta tables to Parquet"
)]
struct Args {
    #[command(subcommand)]
    command: Command,
}

#[derive(Subcommand, Debug)]
enum Command {
    Import {
        #[arg(short, long)]
        table: String,

        #[arg(short, long)]
        csv: PathBuf,
    },

    Export {
        #[arg(short, long)]
        table: String,

        #[arg(short, long)]
        output: PathBuf,
    },
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let args = Args::parse();

    match args.command {
        Command::Import { table, csv } => {
            import_csv(&table, &csv).await?;
        }

        Command::Export { table, output } => {
            export_delta(&table, &output)?;
        }
    }

    Ok(())
}

async fn import_csv(
    table: &str,
    csv: &Path,
) -> Result<(), Box<dyn std::error::Error>> {
    println!("CSV: {}", csv.display());
    println!("Delta table: {}", table);

    let table_url = delta_kernel::try_parse_uri(table)?;

    let object_store = AmazonS3Builder::from_env()
        .with_url(table)
        .build()?;

    let engine =
        DefaultEngineBuilder::new(Arc::new(object_store))
            .build();

    let schema_file = File::open(csv).map_err(|e| {
        format!(
            "Failed to open CSV file '{}': {}",
            csv.display(),
            e
        )
    })?;

    let format = Format::default().with_header(true);

    let (arrow_schema, _) =
        format.infer_schema(schema_file, Some(1000))?;

    let kernel_schema =
        Arc::new((&arrow_schema).try_into_kernel()?);

    let snapshot_result =
        Snapshot::builder_for(&table_url)
            .build(&engine);

    let snapshot = match snapshot_result {
        Ok(snapshot) => {
            println!(
                "Existing Delta table found at version {}",
                snapshot.version()
            );

            if snapshot.schema().as_ref()
                != kernel_schema.as_ref()
            {
                return Err(
                    format!(
                        "CSV schema does not match Delta table schema.\nCSV: {:#?}\nDelta: {:#?}",
                        kernel_schema,
                        snapshot.schema()
                    )
                    .into()
                );
            }

            snapshot
        }

        Err(_) => {
            println!("Delta table not found.");
            println!("Creating new Delta table...");

            let create_result =
                create_table(
                    table,
                    kernel_schema.clone(),
                    "delta-import-export/0.1.0",
                )
                .build(
                    &engine,
                    Box::new(
                        FileSystemCommitter::new()
                    ),
                )?
                .commit(&engine)?;

            match create_result {
                CommitResult::Committed(committed) => {
                    println!(
                        "Created Delta table at version {}",
                        committed.commit_version()
                    );
                }

                CommitResult::Conflicted(conflicted) => {
                    return Err(
                        format!(
                            "Create table conflict at version {}",
                            conflicted.conflict_version()
                        )
                        .into()
                    );
                }

                CommitResult::Retryable(_) => {
                    return Err(
                        "Create table transaction requires retry"
                            .into()
                    );
                }
            }

            Snapshot::builder_for(&table_url)
                .build(&engine)?
        }
    };

    let mut transaction = snapshot
        .transaction(
            Box::new(FileSystemCommitter::new()),
            &engine,
        )?
        .with_operation("WRITE".to_string())
        .with_engine_info(
            "delta-import-export/0.1.0"
        )
        .with_data_change(true)
        .with_blind_append(true);

    let write_context = transaction
        .write_state()?
        .write_context_builder()
        .build()?;

    let csv_file = File::open(csv).map_err(|e| {
        format!(
            "Failed to open CSV file '{}': {}",
            csv.display(),
            e
        )
    })?;

    let csv_reader =
        ReaderBuilder::new(
            Arc::new(arrow_schema)
        )
        .with_header(true)
        .with_batch_size(8192)
        .build(csv_file)?;

    let mut total_rows = 0usize;
    let mut batch_count = 0usize;

    for batch in csv_reader {
        let batch = batch?;

        total_rows += batch.num_rows();
        batch_count += 1;

        println!(
            "Writing batch {}: {} rows",
            batch_count,
            batch.num_rows()
        );

        let data =
            ArrowEngineData::new(batch);

        let file_metadata = engine
            .write_parquet(
                &data,
                &write_context,
            )
            .await?;

        transaction.add_files(file_metadata);
    }

    let mut retries = 0;

    let committed = loop {
        if retries > 5 {
            return Err(
                "Exceeded maximum commit retries"
                    .into()
            );
        }

        transaction =
            match transaction.commit(&engine)? {
                CommitResult::Committed(
                    committed
                ) => {
                    break committed;
                }

                CommitResult::Conflicted(
                    conflicted
                ) => {
                    return Err(
                        format!(
                            "Transaction conflict at version {}",
                            conflicted.conflict_version()
                        )
                        .into()
                    );
                }

                CommitResult::Retryable(
                    RetryableTransaction {
                        transaction,
                        error,
                    },
                ) => {
                    println!(
                        "Retrying commit: {error}"
                    );

                    transaction
                }
            };

        retries += 1;
    };

    println!(
        "Committed Delta version: {}",
        committed.commit_version()
    );

    println!("Rows written: {}", total_rows);
    println!("Batches written: {}", batch_count);

    Ok(())
}

fn export_delta(
    table: &str,
    output: &Path,
) -> Result<(), Box<dyn std::error::Error>> {
    println!("Delta table: {}", table);
    println!("Output file: {}", output.display());

    let table_url =
        delta_kernel::try_parse_uri(table)?;

    let object_store =
        AmazonS3Builder::from_env()
            .with_url(table)
            .build()?;

    let engine =
        DefaultEngineBuilder::new(
            Arc::new(object_store)
        )
        .build();

    let snapshot =
        Snapshot::builder_for(&table_url)
            .build(&engine)?;

    println!(
        "Delta version: {}",
        snapshot.version()
    );

    let scan =
        snapshot
            .scan_builder()
            .build()?;

    let batches =
        scan.execute(Arc::new(engine))?;

    let mut parquet_writer = None;
    let mut total_rows = 0usize;
    let mut batch_count = 0usize;

    for engine_data in batches {
        let batch =
            EngineDataArrowExt::try_into_record_batch(
                engine_data
            )?;

        if parquet_writer.is_none() {
            let file =
                File::create(output)?;

            let properties =
                WriterProperties::builder()
                    .set_compression(
                        Compression::SNAPPY
                    )
                    .build();

            parquet_writer =
                Some(
                    ArrowWriter::try_new(
                        file,
                        batch.schema(),
                        Some(properties),
                    )?
                );
        }

        if let Some(writer) =
            parquet_writer.as_mut()
        {
            writer.write(&batch)?;
        }

        total_rows += batch.num_rows();
        batch_count += 1;

        println!(
            "Reading batch {}: {} rows | total: {}",
            batch_count,
            batch.num_rows(),
            total_rows
        );
    }

    if let Some(writer) =
        parquet_writer
    {
        writer.close()?;
    }

    println!("Rows exported: {}", total_rows);
    println!("Batches exported: {}", batch_count);
    println!("Output: {}", output.display());

    Ok(())
}

and a lib.rs

use clap::{Parser, Subcommand};
use std::path::PathBuf;

#[derive(Parser, Debug)]
#[command(
    name = "delta-import-export",
    version,
    about = "Import and export Delta Lake tables"
)]
pub struct Args {
    #[command(subcommand)]
    pub command: Command,
}

#[derive(Subcommand, Debug)]
pub enum Command {
    Export {
        #[arg(short, long)]
        table: String,

        #[arg(short, long)]
        output: PathBuf,
    },

    Import {
        #[arg(short, long)]
        table: String,

        #[arg(short, long)]
        parquet: PathBuf,
    },
}

This might seem like a lot of code, but it really isn’t.

The Delta Kernel provides a set of APIs to interact with Delta Lake tables, both to READ and WRITE. Obviously, WRITE is the trickier one. The Kernelmakes it easy to build an Enginethat will handle all the types of Deltainteraction a tool might be looking for. Even when that Deltatable sits in remote S3storage.

let table_url = delta_kernel::try_parse_uri(table)?;

let object_store = AmazonS3Builder::from_env()
        .with_url(table)
        .build()?;

let engine =
        DefaultEngineBuilder::new(Arc::new(object_store))
            .build();

It’s super simple to get the current state of a Delta table and Scan the table, build a query, and get results. First-class support for Arrow as well.

let snapshot =
        Snapshot::builder_for(&table_url)
            .build(&engine)?;

    println!(
        "Delta version: {}",
        snapshot.version()
    );

let scan =
        snapshot
            .scan_builder()
            .build()?;

let batches =
        scan.execute(Arc::new(engine))?;

Even Transaction support is made simple, to WRITE to a Delta table, which is no mean feat.

let mut transaction = snapshot
        .transaction(
            Box::new(FileSystemCommitter::new()),
            &engine,
        )?
        .with_operation("WRITE".to_string())
        .with_engine_info(
            "delta-import-export/0.1.0"
        )
        .with_data_change(true)
        .with_blind_append(true);

let write_context = transaction
        .write_state()?
        .write_context_builder()
        .build()?;

If you are looking to build custom programs and applications on top of your Lakehouse in the form of Delta , look no further than the Delta Kernel.

This is a major step forward in the open-source and general openness of the Delta Lakehouse architecture. Anyone can build whatever they please with the Kernel and now there are crates in Rust to interact with Unity Catalog , this will bring these features into full Produciton mode.