micromegas_analytics/lakehouse/
retire_partition_by_metadata_udf.rs1use anyhow::{Context, Result};
2use async_trait::async_trait;
3use chrono::{DateTime, Utc};
4use datafusion::{
5 arrow::{
6 array::{Array, StringArray, StringBuilder, TimestampNanosecondArray},
7 datatypes::{DataType, TimeUnit},
8 },
9 common::internal_err,
10 error::DataFusionError,
11 logical_expr::{
12 ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl, Signature, Volatility,
13 async_udf::AsyncScalarUDFImpl,
14 },
15};
16use micromegas_ingestion::data_lake_connection::DataLakeConnection;
17use micromegas_tracing::prelude::*;
18use sqlx::Row;
19use std::sync::Arc;
20
21use super::write_partition::add_file_for_cleanup;
22
23#[derive(Debug)]
32pub struct RetirePartitionByMetadata {
33 signature: Signature,
34 lake: Arc<DataLakeConnection>,
35}
36
37impl PartialEq for RetirePartitionByMetadata {
38 fn eq(&self, other: &Self) -> bool {
39 self.signature == other.signature
40 }
41}
42
43impl Eq for RetirePartitionByMetadata {}
44
45impl std::hash::Hash for RetirePartitionByMetadata {
46 fn hash<H: std::hash::Hasher>(&self, state: &mut H) {
47 self.signature.hash(state);
48 }
49}
50
51impl RetirePartitionByMetadata {
52 pub fn new(lake: Arc<DataLakeConnection>) -> Self {
53 Self {
54 signature: Signature::exact(
55 vec![
56 DataType::Utf8,
57 DataType::Utf8,
58 DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
59 DataType::Timestamp(TimeUnit::Nanosecond, Some("+00:00".into())),
60 ],
61 Volatility::Volatile,
62 ),
63 lake,
64 }
65 }
66
67 async fn retire_partition_in_transaction(
80 &self,
81 transaction: &mut sqlx::Transaction<'_, sqlx::Postgres>,
82 view_set_name: &str,
83 view_instance_id: &str,
84 begin_insert_time: DateTime<Utc>,
85 end_insert_time: DateTime<Utc>,
86 ) -> Result<()> {
87 let partition_query = instrument_named!(
89 sqlx::query(
90 "SELECT file_path, file_size FROM lakehouse_partitions
91 WHERE view_set_name = $1
92 AND view_instance_id = $2
93 AND begin_insert_time = $3
94 AND end_insert_time = $4",
95 )
96 .bind(view_set_name)
97 .bind(view_instance_id)
98 .bind(begin_insert_time)
99 .bind(end_insert_time)
100 .fetch_optional(&mut **transaction),
101 "sql_select_partition_by_metadata"
102 )
103 .await
104 .with_context(|| {
105 format!(
106 "querying partition {view_set_name}/{view_instance_id} [{begin_insert_time}, {end_insert_time})"
107 )
108 })?;
109
110 let Some(partition_row) = partition_query else {
111 anyhow::bail!(
112 "Partition not found: {view_set_name}/{view_instance_id} [{begin_insert_time}, {end_insert_time})"
113 );
114 };
115
116 let file_path_opt: Option<String> = partition_row.try_get("file_path")?;
118 if let Some(file_path) = file_path_opt {
119 let file_size: i64 = partition_row.try_get("file_size")?;
120 add_file_for_cleanup(transaction, &file_path, file_size).await?;
122 }
123
124 let delete_result = instrument_named!(
126 sqlx::query(
127 "DELETE FROM lakehouse_partitions
128 WHERE view_set_name = $1
129 AND view_instance_id = $2
130 AND begin_insert_time = $3
131 AND end_insert_time = $4",
132 )
133 .bind(view_set_name)
134 .bind(view_instance_id)
135 .bind(begin_insert_time)
136 .bind(end_insert_time)
137 .execute(&mut **transaction),
138 "sql_delete_partition_by_metadata"
139 )
140 .await
141 .with_context(|| {
142 format!(
143 "deleting partition {view_set_name}/{view_instance_id} [{begin_insert_time}, {end_insert_time})"
144 )
145 })?;
146
147 if delete_result.rows_affected() == 0 {
148 anyhow::bail!(
150 "Partition not found during deletion: {view_set_name}/{view_instance_id} [{begin_insert_time}, {end_insert_time})"
151 );
152 }
153
154 info!(
155 "Successfully retired partition: {}/{} [{}, {})",
156 view_set_name, view_instance_id, begin_insert_time, end_insert_time
157 );
158 Ok(())
159 }
160}
161
162impl ScalarUDFImpl for RetirePartitionByMetadata {
163 fn name(&self) -> &str {
164 "retire_partition_by_metadata"
165 }
166
167 fn signature(&self) -> &Signature {
168 &self.signature
169 }
170
171 fn return_type(&self, _arg_types: &[DataType]) -> datafusion::error::Result<DataType> {
172 Ok(DataType::Utf8)
173 }
174
175 fn invoke_with_args(
176 &self,
177 _args: ScalarFunctionArgs,
178 ) -> datafusion::error::Result<ColumnarValue> {
179 Err(DataFusionError::NotImplemented(
180 "retire_partition_by_metadata can only be called from async contexts".into(),
181 ))
182 }
183}
184
185#[async_trait]
186impl AsyncScalarUDFImpl for RetirePartitionByMetadata {
187 async fn invoke_async_with_args(
188 &self,
189 args: ScalarFunctionArgs,
190 ) -> datafusion::error::Result<ColumnarValue> {
191 let args = ColumnarValue::values_to_arrays(&args.args)?;
192 if args.len() != 4 {
193 return internal_err!(
194 "retire_partition_by_metadata expects exactly 4 arguments: view_set_name, view_instance_id, begin_insert_time, end_insert_time"
195 );
196 }
197
198 let view_set_names: &StringArray =
199 args[0].as_any().downcast_ref::<_>().ok_or_else(|| {
200 DataFusionError::Execution(
201 "error casting view_set_name argument as StringArray".into(),
202 )
203 })?;
204
205 let view_instance_ids: &StringArray =
206 args[1].as_any().downcast_ref::<_>().ok_or_else(|| {
207 DataFusionError::Execution(
208 "error casting view_instance_id argument as StringArray".into(),
209 )
210 })?;
211
212 let begin_insert_times: &TimestampNanosecondArray =
213 args[2].as_any().downcast_ref::<_>().ok_or_else(|| {
214 DataFusionError::Execution(
215 "error casting begin_insert_time argument as TimestampNanosecondArray".into(),
216 )
217 })?;
218
219 let end_insert_times: &TimestampNanosecondArray =
220 args[3].as_any().downcast_ref::<_>().ok_or_else(|| {
221 DataFusionError::Execution(
222 "error casting end_insert_time argument as TimestampNanosecondArray".into(),
223 )
224 })?;
225
226 let mut builder = StringBuilder::with_capacity(view_set_names.len(), 64);
227
228 let mut transaction =
230 self.lake.db_pool.begin().await.map_err(|e| {
231 DataFusionError::Execution(format!("Failed to begin transaction: {e}"))
232 })?;
233
234 let mut success_count = 0;
235 let mut has_errors = false;
236
237 for index in 0..view_set_names.len() {
239 if view_set_names.is_null(index)
240 || view_instance_ids.is_null(index)
241 || begin_insert_times.is_null(index)
242 || end_insert_times.is_null(index)
243 {
244 builder.append_value("ERROR: all arguments must be non-null");
245 has_errors = true;
246 continue;
247 }
248
249 let view_set_name = view_set_names.value(index);
250 let view_instance_id = view_instance_ids.value(index);
251 let begin_insert_time_nanos = begin_insert_times.value(index);
252 let end_insert_time_nanos = end_insert_times.value(index);
253
254 let begin_insert_time = DateTime::from_timestamp_nanos(begin_insert_time_nanos);
256 let end_insert_time = DateTime::from_timestamp_nanos(end_insert_time_nanos);
257
258 match self
259 .retire_partition_in_transaction(
260 &mut transaction,
261 view_set_name,
262 view_instance_id,
263 begin_insert_time,
264 end_insert_time,
265 )
266 .await
267 {
268 Ok(()) => {
269 success_count += 1;
270 builder.append_value(format!(
271 "SUCCESS: Retired partition {view_set_name}/{view_instance_id} [{begin_insert_time}, {end_insert_time})"
272 ));
273 }
274 Err(e) => {
275 error!(
276 "Failed to retire partition {}/{} [{}, {}): {:?}",
277 view_set_name, view_instance_id, begin_insert_time, end_insert_time, e
278 );
279 builder.append_value(format!("ERROR: {e:?}"));
280 has_errors = true;
281 }
282 }
283 }
284
285 if has_errors {
287 if let Err(e) = transaction.rollback().await {
288 error!("Failed to rollback transaction after errors: {:?}", e);
289 }
290 info!("Rolled back transaction due to errors in batch retirement");
291 builder.append_value(format!(
292 "ROLLED_BACK: All {} previous changes were reverted due to errors in batch",
293 success_count
294 ));
295 } else {
296 transaction.commit().await.map_err(|e| {
297 DataFusionError::Execution(format!("Failed to commit transaction: {e}"))
298 })?;
299 info!("Successfully retired {} partitions in batch", success_count);
300 }
301
302 Ok(ColumnarValue::Array(Arc::new(builder.finish())))
303 }
304}
305
306pub fn make_retire_partition_by_metadata_udf(
327 lake: Arc<DataLakeConnection>,
328) -> datafusion::logical_expr::async_udf::AsyncScalarUDF {
329 datafusion::logical_expr::async_udf::AsyncScalarUDF::new(Arc::new(
330 RetirePartitionByMetadata::new(lake),
331 ))
332}