feat: skip checkpoint cleanup for nodes with running jobs

This commit is contained in:
Aditya Ujeniya
2026-07-08 14:13:15 +02:00
parent 1e416a6dd0
commit 58c71da0b8
2 changed files with 130 additions and 5 deletions
+29 -5
View File
@@ -94,16 +94,32 @@ var ErrNoNewArchiveData error = errors.New("all data already archived")
// CleanupCheckpoints deletes or archives all checkpoint files older than `from`. // CleanupCheckpoints deletes or archives all checkpoint files older than `from`.
// When archiving, consolidates all hosts per cluster into a single Parquet file. // When archiving, consolidates all hosts per cluster into a single Parquet file.
//
// Hosts currently in use by running jobs (per the MemoryStore's NodeProvider)
// are skipped entirely; their files are handled by a later cycle once the job
// has ended. If the provider query fails, the whole cycle is aborted: deleting
// without knowing which nodes are in use risks data loss. Without a provider,
// everything older than `from` is cleaned (old behavior).
func CleanupCheckpoints(checkpointsDir, cleanupDir string, from int64, deleteInstead bool) (int, error) { func CleanupCheckpoints(checkpointsDir, cleanupDir string, from int64, deleteInstead bool) (int, error) {
if deleteInstead { var usedNodes map[string][]string
return deleteCheckpoints(checkpointsDir, from) if ms := GetMemoryStore(); ms != nil && ms.nodeProvider != nil {
un, err := ms.nodeProvider.GetUsedNodes(from)
if err != nil {
return 0, fmt.Errorf("aborting cleanup, could not determine nodes in use: %w", err)
}
usedNodes = un
} }
return archiveCheckpoints(checkpointsDir, cleanupDir, from) if deleteInstead {
return deleteCheckpoints(checkpointsDir, from, usedNodes)
}
return archiveCheckpoints(checkpointsDir, cleanupDir, from, usedNodes)
} }
// deleteCheckpoints removes checkpoint files older than `from` across all clusters/hosts. // deleteCheckpoints removes checkpoint files older than `from` across all clusters/hosts.
func deleteCheckpoints(checkpointsDir string, from int64) (int, error) { // Hosts present in usedNodes are skipped.
func deleteCheckpoints(checkpointsDir string, from int64, usedNodes map[string][]string) (int, error) {
entries1, err := os.ReadDir(checkpointsDir) entries1, err := os.ReadDir(checkpointsDir)
if err != nil { if err != nil {
return 0, err return 0, err
@@ -157,6 +173,10 @@ func deleteCheckpoints(checkpointsDir string, from int64) (int, error) {
} }
for _, de2 := range entries2 { for _, de2 := range entries2 {
if isNodeUsed(usedNodes, de1.Name(), de2.Name()) {
cclog.Debugf("[METRICSTORE]> skipping delete for %s/%s: node in use by running job", de1.Name(), de2.Name())
continue
}
work <- workItem{ work <- workItem{
dir: filepath.Join(checkpointsDir, de1.Name(), de2.Name()), dir: filepath.Join(checkpointsDir, de1.Name(), de2.Name()),
cluster: de1.Name(), cluster: de1.Name(),
@@ -182,7 +202,7 @@ func deleteCheckpoints(checkpointsDir string, from int64) (int, error) {
// Workers load checkpoint files from disk and send CheckpointFile trees on a // Workers load checkpoint files from disk and send CheckpointFile trees on a
// back-pressured channel. The main thread streams each tree directly to Parquet // back-pressured channel. The main thread streams each tree directly to Parquet
// rows without materializing all rows in memory. // rows without materializing all rows in memory.
func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64) (int, error) { func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64, usedNodes map[string][]string) (int, error) {
cclog.Info("[METRICSTORE]> start archiving checkpoints to parquet") cclog.Info("[METRICSTORE]> start archiving checkpoints to parquet")
startTime := time.Now() startTime := time.Now()
@@ -252,6 +272,10 @@ func archiveCheckpoints(checkpointsDir, cleanupDir string, from int64) (int, err
if !hostEntry.IsDir() { if !hostEntry.IsDir() {
continue continue
} }
if isNodeUsed(usedNodes, cluster, hostEntry.Name()) {
cclog.Debugf("[METRICSTORE]> skipping archive for %s/%s: node in use by running job", cluster, hostEntry.Name())
continue
}
dir := filepath.Join(checkpointsDir, cluster, hostEntry.Name()) dir := filepath.Join(checkpointsDir, cluster, hostEntry.Name())
work <- struct { work <- struct {
dir, host string dir, host string
+101
View File
@@ -161,3 +161,104 @@ func TestFromCheckpointProviderErrorFallsBack(t *testing.T) {
t.Fatalf("expected fallback to cutoff load (4 files), got %d", n) t.Fatalf("expected fallback to cutoff load (4 files), got %d", n)
} }
} }
func TestDeleteCheckpointsSkipsUsedNodes(t *testing.T) {
oldWorkers := Keys.NumWorkers
Keys.NumWorkers = 2
t.Cleanup(func() { Keys.NumWorkers = oldWorkers })
dir := t.TempDir()
writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000)
writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node002"), 1000)
used := map[string][]string{"fritz": {"node001"}}
n, err := deleteCheckpoints(dir, 3000, used)
if err != nil {
t.Fatal(err)
}
if n != 1 {
t.Fatalf("expected 1 file deleted, got %d", n)
}
if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil {
t.Errorf("used node's checkpoint must survive: %v", err)
}
if _, err := os.Stat(filepath.Join(dir, "fritz", "node002", "1000.json")); !os.IsNotExist(err) {
t.Error("unused node's checkpoint must be deleted")
}
}
func TestArchiveCheckpointsSkipsUsedNodes(t *testing.T) {
oldWorkers := Keys.NumWorkers
Keys.NumWorkers = 2
t.Cleanup(func() { Keys.NumWorkers = oldWorkers })
cpDir := t.TempDir()
archiveDir := t.TempDir()
writeTestCheckpoint(t, filepath.Join(cpDir, "fritz", "node001"), 1000)
writeTestCheckpoint(t, filepath.Join(cpDir, "fritz", "node002"), 1000)
used := map[string][]string{"fritz": {"node001"}}
n, err := archiveCheckpoints(cpDir, archiveDir, 3000, used)
if err != nil {
t.Fatal(err)
}
if n != 1 {
t.Fatalf("expected 1 file archived, got %d", n)
}
if _, err := os.Stat(filepath.Join(cpDir, "fritz", "node001", "1000.json")); err != nil {
t.Errorf("used node's checkpoint must survive: %v", err)
}
if _, err := os.Stat(filepath.Join(cpDir, "fritz", "node002", "1000.json")); !os.IsNotExist(err) {
t.Error("unused node's checkpoint must be removed after archiving")
}
if _, err := os.Stat(filepath.Join(archiveDir, "fritz", "3000.parquet")); err != nil {
t.Errorf("parquet archive must exist: %v", err)
}
}
func TestCleanupCheckpointsAbortsOnProviderError(t *testing.T) {
oldWorkers := Keys.NumWorkers
Keys.NumWorkers = 2
t.Cleanup(func() { Keys.NumWorkers = oldWorkers })
oldMS := msInstance
msInstance = newTestStore()
msInstance.SetNodeProvider(&fakeNodeProvider{err: errTestProvider})
t.Cleanup(func() { msInstance = oldMS })
dir := t.TempDir()
writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000)
if _, err := CleanupCheckpoints(dir, "", 3000, true); err == nil {
t.Fatal("expected error when GetUsedNodes fails")
}
if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil {
t.Errorf("no files may be deleted when provider errors: %v", err)
}
}
func TestCleanupCheckpointsUsedNodesSurvive(t *testing.T) {
oldWorkers := Keys.NumWorkers
Keys.NumWorkers = 2
t.Cleanup(func() { Keys.NumWorkers = oldWorkers })
oldMS := msInstance
msInstance = newTestStore()
msInstance.SetNodeProvider(&fakeNodeProvider{nodes: map[string][]string{"fritz": {"node001"}}})
t.Cleanup(func() { msInstance = oldMS })
dir := t.TempDir()
writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node001"), 1000)
writeTestCheckpoint(t, filepath.Join(dir, "fritz", "node002"), 1000)
n, err := CleanupCheckpoints(dir, "", 3000, true)
if err != nil {
t.Fatal(err)
}
if n != 1 {
t.Fatalf("expected 1 file deleted, got %d", n)
}
if _, err := os.Stat(filepath.Join(dir, "fritz", "node001", "1000.json")); err != nil {
t.Errorf("used node's checkpoint must survive: %v", err)
}
}