Single attempt of CoW DELETE execution.
(
ctx: &SQLContext,
delete: &Delete,
table: &Table,
table_ref: &str,
)
| 88 | |
| 89 | /// Single attempt of CoW DELETE execution. |
| 90 | async fn execute_cow_delete_once( |
| 91 | ctx: &SQLContext, |
| 92 | delete: &Delete, |
| 93 | table: &Table, |
| 94 | table_ref: &str, |
| 95 | ) -> DFResult<DataFrame> { |
| 96 | let where_str = delete.selection.as_ref().map(|e| e.to_string()); |
| 97 | let partition_set = |
| 98 | build_partition_set_from_where(ctx, table, table_ref, where_str.as_deref()).await?; |
| 99 | |
| 100 | let mut writer = CopyOnWriteMergeWriter::new(table, vec![], partition_set) |
| 101 | .await |
| 102 | .map_err(to_datafusion_error)?; |
| 103 | |
| 104 | let mut temp_tracker = TempTableTracker::new(ctx); |
| 105 | let (has_data, cow_table_name) = |
| 106 | register_cow_target_table(ctx, table, &writer, &mut temp_tracker).await?; |
| 107 | if !has_data { |
| 108 | return ok_result(ctx.ctx(), 0); |
| 109 | } |
| 110 | |
| 111 | let cow_target_qualified = cow_table_name; |
| 112 | let result = |
| 113 | execute_cow_delete_inner(ctx.ctx(), &cow_target_qualified, delete, &mut writer).await; |
| 114 | let total_count = result?; |
| 115 | |
| 116 | let messages = writer.prepare_commit().await.map_err(to_datafusion_error)?; |
| 117 | if !messages.is_empty() { |
| 118 | let commit = table.new_write_builder().new_commit(); |
| 119 | commit.commit(messages).await.map_err(to_datafusion_error)?; |
| 120 | } |
| 121 | |
| 122 | ok_result(ctx.ctx(), total_count) |
| 123 | } |
| 124 | |
| 125 | async fn execute_cow_delete_inner( |
| 126 | ctx: &SessionContext, |
no test coverage detected