From 8a6f311c6bdcdcf88638d7a2333019298c0b079b Mon Sep 17 00:00:00 2001 From: Bjoern Gruening Date: Fri, 5 Jun 2026 04:42:52 +0200 Subject: [PATCH] try to optimize the purge_historyless_hdas case --- scripts/cleanup_datasets/pgcleanup.py | 42 ++++++++++++++++----------- 1 file changed, 25 insertions(+), 17 deletions(-) diff --git a/scripts/cleanup_datasets/pgcleanup.py b/scripts/cleanup_datasets/pgcleanup.py index c88a44e9c32..6921911ac03 100755 --- a/scripts/cleanup_datasets/pgcleanup.py +++ b/scripts/cleanup_datasets/pgcleanup.py @@ -215,35 +215,38 @@ class Action: args["days"] = self._days return args - def _collect_row_results(self, row, results, primary_key): - primary_key = row._fields[0] - primary = getattr(row, primary_key) - if primary not in results: - results[primary] = [set() for x in range(len(self.causals))] + def _collect_row_results(self, row, results): rowgetter = partial(getattr, row) for i, causal in enumerate(self.causals): vals = tuple(map(rowgetter, causal)) if any(vals[1:]): - results[primary][i].add(vals) + results[i].add(vals) - def _log_results(self, results, primary_key): - for primary in sorted(results.keys()): - self.log.info(f"{primary_key}: {primary}") - for causal, s in zip(self.causals, results[primary]): - for r in sorted(s): - secondaries = ", ".join(f"{x[0]}: {x[1]}" for x in zip(causal[1:], r[1:])) - self.log.info(f"{causal[0]} {r[0]} caused {secondaries}") + def _log_results(self, primary, results, primary_key): + self.log.info(f"{primary_key}: {primary}") + for causal, s in zip(self.causals, results): + for r in sorted(s): + secondaries = ", ".join(f"{x[0]}: {x[1]}" for x in zip(causal[1:], r[1:])) + self.log.info(f"{causal[0]} {r[0]} caused {secondaries}") def handle_results(self, cur): - results = {} primary_key = self.primary_key + primary = None + results = None for row in cur: if primary_key is None: primary_key = row._fields[0] - self._collect_row_results(row, results, primary_key) + row_primary = getattr(row, primary_key) + if results is None or row_primary != primary: + if results is not None: + self._log_results(primary, results, primary_key) + primary = row_primary + results = [set() for x in range(len(self.causals))] + self._collect_row_results(row, results) for method in self.__row_methods: method(row) - self._log_results(results, primary_key) + if results is not None: + self._log_results(primary, results, primary_key) for method in self.__post_methods: method() @@ -261,6 +264,7 @@ class RemovesObjects: """Base class for mixins that remove objects from object stores.""" requires_objectstore = True + object_removal_batch_size = 1000 def _init(self): super()._init() @@ -276,10 +280,13 @@ class RemovesObjects: object_uuid = str(uuid.UUID(object_uuid)) if object_id: self.objects_to_remove.add(self.object_class(object_id, row.object_store_id, object_uuid)) + if len(self.objects_to_remove) >= self.object_removal_batch_size: + self.remove_objects() def remove_objects(self): for object_to_remove in sorted(self.objects_to_remove): self.remove_object(object_to_remove) + self.objects_to_remove.clear() def remove_from_object_store(self, object_to_remove, object_store_kwargs, entire_dir=False, check_exists=False): # only remove the "object store path" - if it's at an external_filename, that file will be untouched anyway @@ -905,7 +912,8 @@ class PurgeHistorylessHDAs(PurgesHDAs, RemovesMetadataFiles, RequiresDiskUsageRe AS ( UPDATE history_dataset_association SET purged = true, deleted = true{update_time_sql} FROM dataset - WHERE history_id IS NULL{force_retry_sql}{object_store_id_sql} + WHERE history_dataset_association.dataset_id = dataset.id + AND history_id IS NULL{force_retry_sql}{object_store_id_sql} AND history_dataset_association.update_time < (NOW() AT TIME ZONE 'utc' - (:days * interval '1 day')) RETURNING history_dataset_association.id as id, history_dataset_association.history_id as history_id),