Skip to main content

micromegas_analytics/lakehouse/
retire_partition_by_metadata_udf.rs

1use 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/// A scalar UDF that retires a single partition by its metadata.
24///
25/// This function retires partitions by their metadata identifiers (view_set_name,
26/// view_instance_id, begin_insert_time, end_insert_time). This works for both empty
27/// partitions (file_path=NULL) and non-empty partitions.
28///
29/// This is the preferred method for retiring partitions as it uses the partition's
30/// natural identifiers rather than relying on file paths.
31#[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    /// Retires a single partition by its metadata within an existing transaction.
68    ///
69    /// # Arguments
70    /// * `transaction` - Database transaction to use
71    /// * `view_set_name` - The name of the view set
72    /// * `view_instance_id` - The instance ID (e.g., process_id or 'global')
73    /// * `begin_insert_time` - Begin insert time timestamp
74    /// * `end_insert_time` - End insert time timestamp
75    ///
76    /// # Returns
77    /// * `Ok(())` on successful retirement
78    /// * `Err(anyhow::Error)` with descriptive message for any failure
79    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        // First, check if the partition exists and get its details
88        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        // Handle file cleanup if file_path is not NULL
117        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 to temporary files for cleanup (expires in 1 hour)
121            add_file_for_cleanup(transaction, &file_path, file_size).await?;
122        }
123
124        // Remove from active partitions
125        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            // This shouldn't happen since we checked existence above, but handle it gracefully
149            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        // Use a single transaction for the entire batch
229        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        // Process each partition in the batch within the same transaction
238        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            // Convert nanoseconds to DateTime<Utc> for proper sqlx binding
255            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        // Commit the transaction only if there were no errors
286        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
306/// Creates a user-defined function to retire a single partition by its metadata.
307///
308/// This function retires partitions by their metadata identifiers rather than file path,
309/// making it suitable for both empty partitions (file_path=NULL) and non-empty partitions.
310///
311/// # Usage
312/// ```sql
313/// SELECT retire_partition_by_metadata(
314///     'log_entries',
315///     'process_123',
316///     TIMESTAMP '2024-01-01 00:00:00',
317///     TIMESTAMP '2024-01-01 01:00:00'
318/// ) as result;
319/// ```
320///
321/// # Returns
322/// A string message indicating success or failure:
323/// - "SUCCESS: Retired partition \<view_set\>/\<instance\> [\<begin\>, \<end\>)" on successful retirement
324/// - "ERROR: Partition not found: \<view_set\>/\<instance\> [\<begin\>, \<end\>)" if the partition doesn't exist
325/// - "ERROR: Database error: \<details\>" for any database-related failures
326pub 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}