Skip to content
Open
Show file tree
Hide file tree
Changes from 2 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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

1 change: 1 addition & 0 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ async-trait = "0.1"
backon = { version = "1.5" }
bincode = { version = "2", features = ["serde"] }
bitstream-io = "4.5"
bytes = "1.0"
chrono = { version = "0.4", default-features = false }
clap = { version = "4", features = ["derive"] }
crc32fast = "1"
Expand Down
1 change: 1 addition & 0 deletions src/moonlink/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ backon = { workspace = true }
base64 = { version = "0.22", optional = true }
bincode = { workspace = true }
bitstream-io = { workspace = true }
bytes = { workspace = true }
chrono = { workspace = true }
crc32fast = { workspace = true }
fastbloom = { workspace = true }
Expand Down
1 change: 1 addition & 0 deletions src/moonlink/src/storage/mooncake_table.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
mod batch_id_counter;
mod batch_ingestion;
mod data_batches;
pub(crate) mod delete_vector;
mod disk_slice;
Expand Down
59 changes: 59 additions & 0 deletions src/moonlink/src/storage/mooncake_table/batch_ingestion.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
use crate::{
create_data_file, storage::mooncake_table::transaction_stream::TransactionStreamCommit,
};

use super::*;

use bytes::Bytes;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;

impl MooncakeTable {
/// Batch ingestion the given [`parquet_files`] into mooncake table.
///
/// TODO(hjiang):
/// 1. Record table events.
/// 2. It involves IO operations, should be placed at background thread.s
pub(crate) async fn batch_ingest(&mut self, parquet_files: Vec<String>, lsn: u64) {
let mut disk_files = hashbrown::HashMap::with_capacity(parquet_files.len());

// TODO(hjiang): Parallel IO and error handling.
// Construct disk file from the given disk files.
for cur_file in parquet_files.into_iter() {
let content = tokio::fs::read(&cur_file).await.unwrap();
let content = Bytes::from(content);
let file_size = content.len();
let builder = ParquetRecordBatchReaderBuilder::try_new(content).unwrap();
let num_rows: usize = builder.metadata().file_metadata().num_rows() as usize;

let cur_file_id = self.next_file_id;
self.next_file_id += 1;
let mooncake_data_file = create_data_file(cur_file_id as u64, cur_file);
let disk_file_entry = DiskFileEntry {
cache_handle: None,
num_rows,
file_size,
batch_deletion_vector: BatchDeletionVector::new(num_rows),
puffin_deletion_blob: None,
};

assert!(disk_files
.insert(mooncake_data_file, disk_file_entry)
.is_none());
}

// let mut stream_state = TransactionStreamState::new(
// self.metadata.schema.clone(),
// self.metadata.config.batch_size,
// self.metadata.identity.clone(),
// Arc::clone(&self.streaming_batch_id_counter),
// );
// stream_state.flushed_files = disk_files;
// stream_state.commit_lsn = Some(lsn);

// Commit the current crafted streaming transaction.
let commit = TransactionStreamCommit::from_disk_files(disk_files, lsn);
self.next_snapshot_task
.new_streaming_xact
.push(TransactionStreamOutput::Commit(commit));
}
}
22 changes: 20 additions & 2 deletions src/moonlink/src/storage/mooncake_table/transaction_stream.rs
Original file line number Diff line number Diff line change
Expand Up @@ -28,7 +28,7 @@ pub(crate) struct TransactionStreamState {
/// Only safe to remove transaction stream state when there are no pending flushes.
ongoing_flush_count: u32,
/// Commit LSN for this transaction, set when transaction is committed.
commit_lsn: Option<u64>,
pub(crate) commit_lsn: Option<u64>,
}

/// Determines the state of a transaction stream.
Expand Down Expand Up @@ -65,6 +65,24 @@ pub struct TransactionStreamCommit {
}

impl TransactionStreamCommit {
/// Create a stream commit from disk files.
pub(crate) fn from_disk_files(
disk_files: hashbrown::HashMap<MooncakeDataFileRef, DiskFileEntry>,
lsn: u64,
) -> Self {
Self {
xact_id: 0, // Unused.
commit_lsn: lsn,
flushed_file_index: MooncakeIndex {
in_memory_index: HashSet::new(),
file_indices: Vec::new(),
},
flushed_files: disk_files,
local_deletions: Vec::new(),
pending_deletions: Vec::new(),
}
}

/// Get flushed data files for the current streaming commit.
pub(crate) fn get_flushed_data_files(&self) -> Vec<MooncakeDataFileRef> {
self.flushed_files.keys().cloned().collect::<Vec<_>>()
Expand All @@ -91,7 +109,7 @@ impl TransactionStreamCommit {
}

impl TransactionStreamState {
fn new(
pub(crate) fn new(
schema: Arc<Schema>,
batch_size: usize,
identity: IdentityProp,
Expand Down
7 changes: 7 additions & 0 deletions src/moonlink/src/table_handler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,13 @@ impl TableHandler {
Self::process_cdc_table_event(event, &mut table, &mut table_handler_state)
.await;
}
// ==============================
// Bulk ingestion events
// ==============================
TableEvent::LoadFiles { files, lsn } => {
table.batch_ingest(files, lsn).await;
}

// ==============================
// Interactive blocking events
// ==============================
Expand Down
28 changes: 27 additions & 1 deletion src/moonlink/src/table_handler/test_utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,12 +22,13 @@ use crate::{
WalConfig, WalTransactionState,
};

use arrow_array::RecordBatch;
use arrow_array::{Int32Array, RecordBatch, StringArray};
use futures::StreamExt;
use iceberg::io::FileIOBuilder;
use iceberg::io::FileRead;
use more_asserts as ma;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::arrow::AsyncArrowWriter;
use std::sync::Arc;
use tempfile::{tempdir, TempDir};
use tokio::sync::{mpsc, watch};
Expand Down Expand Up @@ -337,6 +338,10 @@ impl TestEnvironment {
.await;
}

pub async fn bulk_upload_files(&self, files: Vec<String>, lsn: u64) {
self.send_event(TableEvent::LoadFiles { files, lsn }).await;
}

// --- LSN and Verification Helpers ---

/// Sets both table commit and replication LSN to the same value.
Expand Down Expand Up @@ -613,3 +618,24 @@ pub(crate) async fn load_one_arrow_batch(filepath: &str) -> RecordBatch {
.unwrap()
.expect("Should have one batch")
}

/// Test util function to generate a parquet under the given [`tempdir`].
pub(crate) async fn generate_parquet_file(tempdir: &TempDir) -> String {
let schema = create_test_arrow_schema();
let ids = Int32Array::from(vec![1, 2, 3]);
let names = StringArray::from(vec!["Alice", "Bob", "Charlie"]);
let ages = Int32Array::from(vec![10, 20, 30]);
let batch = RecordBatch::try_new(
schema.clone(),
vec![Arc::new(ids), Arc::new(names), Arc::new(ages)],
)
.unwrap();
let file_path = tempdir.path().join("test.parquet");
let file_path_str = file_path.to_str().unwrap().to_string();
let file = tokio::fs::File::create(file_path).await.unwrap();
let mut writer: AsyncArrowWriter<tokio::fs::File> =
AsyncArrowWriter::try_new(file, schema, /*props=*/ None).unwrap();
writer.write(&batch).await.unwrap();
writer.close().await.unwrap();
file_path_str
}
11 changes: 11 additions & 0 deletions src/moonlink/src/table_handler/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2114,3 +2114,14 @@ async fn test_commit_flush_streaming_transaction_with_deletes() {

env.shutdown().await;
}

/// Testing scenario: batch upload parquet files into mooncake table.
#[tokio::test]
async fn test_batch_ingestion() {
let env = TestEnvironment::default().await;
let disk_file = generate_parquet_file(&env.temp_dir).await;
env.bulk_upload_files(vec![disk_file], /*lsn=*/ 10).await;

env.set_readable_lsn(10);
env.verify_snapshot(101, &[1, 2, 3]).await;
}
14 changes: 14 additions & 0 deletions src/moonlink/src/table_notify.rs
Original file line number Diff line number Diff line change
Expand Up @@ -103,6 +103,18 @@ pub enum TableEvent {
/// Result for mem slice flush.
flush_result: Option<Result<DiskSliceWriter>>,
},

/// ==============================
/// Bulk ingestion events
/// ==============================
///
LoadFiles {
/// Parquet files to directly load into mooncake table, without schema validation, index construction, etc.
files: Vec<String>,
/// LSN for the bulk upload operation.
lsn: u64,
},

/// ==============================
/// Test events
/// ==============================
Expand All @@ -118,6 +130,7 @@ pub enum TableEvent {
xact_id: u32,
is_recovery: bool,
},

/// ==============================
/// Interactive blocking events
/// ==============================
Expand Down Expand Up @@ -149,6 +162,7 @@ pub enum TableEvent {
FinishInitialCopy {
start_lsn: u64,
},

/// ==============================
/// Table internal events
/// ==============================
Expand Down
2 changes: 1 addition & 1 deletion src/moonlink_connectors/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -17,7 +17,7 @@ arrow-schema = { workspace = true }
async-trait = { workspace = true }
bigdecimal = { version = "0.4", default-features = false, features = ["std"] }
byteorder = "1.5"
bytes = "1.0"
bytes = { workspace = true }
chrono = { workspace = true }
chrono-tz = "0.10"
futures = { workspace = true }
Expand Down