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

Function execute_cow_delete_once

crates/integrations/datafusion/src/delete.rs:90–123  ·  view source on GitHub ↗

Single attempt of CoW DELETE execution.

(
    ctx: &SQLContext,
    delete: &Delete,
    table: &Table,
    table_ref: &str,
)

Source from the content-addressed store, hash-verified

88
89/// Single attempt of CoW DELETE execution.
90async 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
125async fn execute_cow_delete_inner(
126 ctx: &SessionContext,

Callers 1

execute_cow_deleteFunction · 0.85

Calls 10

execute_cow_delete_innerFunction · 0.85
ctxMethod · 0.80
new_commitMethod · 0.80
new_write_builderMethod · 0.80
ok_resultFunction · 0.70
prepare_commitMethod · 0.45
is_emptyMethod · 0.45
commitMethod · 0.45

Tested by

no test coverage detected