try to optimize the purge_historyless_hdas case

This commit is contained in:
Bjoern Gruening
2026-06-05 04:42:52 +02:00
parent faee1c1a58
commit 8a6f311c6b
+25 -17
View File
@@ -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),