MCPcopy Create free account
hub / github.com/apache/paimon-rust / execute_cow_delete_inner

Function execute_cow_delete_inner

crates/integrations/datafusion/src/delete.rs:125–159  ·  view source on GitHub ↗
(
    ctx: &SessionContext,
    cow_table_name: &str,
    delete: &Delete,
    writer: &mut CopyOnWriteMergeWriter,
)

Source from the content-addressed store, hash-verified

123}
124
125async fn execute_cow_delete_inner(
126 ctx: &SessionContext,
127 cow_table_name: &str,
128 delete: &Delete,
129 writer: &mut CopyOnWriteMergeWriter,
130) -> DFResult<u64> {
131 let where_clause = match &delete.selection {
132 Some(expr) => format!(" WHERE {expr}"),
133 None => String::new(),
134 };
135
136 // Safety: where_clause comes from sqlparser AST to_string(), not raw user input.
137 let query_sql = format!(
138 "SELECT \"__paimon_file_idx\", \"__paimon_row_offset\" FROM {cow_table_name}{where_clause}"
139 );
140 let batches = ctx.sql(&query_sql).await?.collect().await?;
141
142 let mut total_count: u64 = 0;
143 for batch in &batches {
144 if batch.num_rows() == 0 {
145 continue;
146 }
147
148 let (file_idx_col, row_offset_col) = extract_tracking_columns(batch)?;
149
150 for row in 0..batch.num_rows() {
151 let file_idx = file_idx_col.value(row) as usize;
152 let row_offset = row_offset_col.value(row) as usize;
153 writer.add_matched_delete(file_idx, row_offset);
154 total_count += 1;
155 }
156 }
157
158 Ok(total_count)
159}

Callers 1

execute_cow_delete_onceFunction · 0.85

Calls 4

extract_tracking_columnsFunction · 0.85
num_rowsMethod · 0.80
add_matched_deleteMethod · 0.80
sqlMethod · 0.45

Tested by

no test coverage detected