(
ctx: &SessionContext,
cow_table_name: &str,
delete: &Delete,
writer: &mut CopyOnWriteMergeWriter,
)
| 123 | } |
| 124 | |
| 125 | async 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 | } |
no test coverage detected