From d6464b138dba9f85f311a4dd157329b2f70639b7 Mon Sep 17 00:00:00 2001 From: yafeng Date: Sat, 18 Jul 2026 20:15:15 -0700 Subject: [PATCH 1/4] CASSSIDECAR-484: Fix the spurious oldest segment age in CdcRawDirectorySpaceCleaner --- .../tasks/CdcRawDirectorySpaceCleaner.java | 72 +++++++---- .../CdcRawDirectorySpaceCleanerTest.java | 122 ++++++++++++++++++ 2 files changed, 170 insertions(+), 24 deletions(-) diff --git a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java index 03ed40a14..b1775c21d 100644 --- a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java +++ b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java @@ -232,8 +232,13 @@ protected void cleanUpCdcRawDirectory(File cdcRawDirectory) Collections.sort(segmentFiles); long nowInMillis = timeProvider.currentTimeMillis(); - // track the age of the oldest commit log segment to give indication of the time-window buffer available - cdcMetrics.oldestSegmentAge.metric.setValue((int) MILLISECONDS.toSeconds(nowInMillis - segmentFiles.get(0).lastModified())); + // track the age of the oldest commit log segment to give indication of the time-window buffer available. + // Skip emission if lastModified is 0 (file was deleted before we could snapshot its lastModified time). + long oldestLastModified = segmentFiles.get(0).lastModified(); + if (oldestLastModified > 0) + { + cdcMetrics.oldestSegmentAge.metric.setValue((int) MILLISECONDS.toSeconds(nowInMillis - oldestLastModified)); + } LOGGER.debug("Cdc data cleaner directorySizeBytes={} maxedUsageBytes={} upperLimitBytes={}", directorySizeBytes, maxUsageBytes, upperLimitBytes); @@ -254,31 +259,45 @@ protected void cleanUpCdcRawDirectory(File cdcRawDirectory) while (i < segmentFiles.size() - 1 && directorySizeBytes > upperLimitBytes) { CdcRawSegmentFile segment = segmentFiles.get(i); - long ageMillis = nowInMillis - segment.lastModified(); - - if (ageMillis < criticalMillis) - { - LOGGER.error("Insufficient Cdc buffer size to maintain {}-minute window segment={} maxSize={} ageMinutes={}", - MILLISECONDS.toMinutes(criticalMillis), segment, upperLimitBytes, - MILLISECONDS.toMinutes(ageMillis)); - cdcMetrics.criticalCdcRawSpace.metric.update(1); - } - else if (ageMillis < lowMillis) - { - LOGGER.warn("Insufficient Cdc buffer size to maintain {}-minute window segment={} maxSize={} ageMinutes={}", - MILLISECONDS.toMinutes(lowMillis), segment, upperLimitBytes, - MILLISECONDS.toMinutes(ageMillis)); - cdcMetrics.lowCdcRawSpace.metric.update(1); - } + long segmentLastModified = segment.lastModified(); long length = 0; - try + // When lastModified is 0 the segment was already reclaimed between the directory + // Skip the buffer-window alerts and the metrics update but still discount the cached size + // from the local budget so the outer loop stops at the right point instead of over-cleaning + // subsequent live segments. + if (segmentLastModified > 0) { - length = deleteSegment(segment); - cdcMetrics.deletedSegment.metric.update(length); + long ageMillis = nowInMillis - segmentLastModified; + + if (ageMillis < criticalMillis) + { + LOGGER.error("Insufficient Cdc buffer size to maintain {}-minute window segment={} maxSize={} ageMinutes={}", + MILLISECONDS.toMinutes(criticalMillis), segment, upperLimitBytes, + MILLISECONDS.toMinutes(ageMillis)); + cdcMetrics.criticalCdcRawSpace.metric.update(1); + } + else if (ageMillis < lowMillis) + { + LOGGER.warn("Insufficient Cdc buffer size to maintain {}-minute window segment={} maxSize={} ageMinutes={}", + MILLISECONDS.toMinutes(lowMillis), segment, upperLimitBytes, + MILLISECONDS.toMinutes(ageMillis)); + cdcMetrics.lowCdcRawSpace.metric.update(1); + } + + try + { + length = deleteSegment(segment); + cdcMetrics.deletedSegment.metric.update(length); + } + catch (IOException e) + { + LOGGER.warn("Failed to delete cdc segment", e); + } } - catch (IOException e) + else { - LOGGER.warn("Failed to delete cdc segment", e); + LOGGER.debug("Skipping delete for already-reclaimed cdc segment {}", segment); + length = segment.length(); } directorySizeBytes -= length; i++; @@ -398,6 +417,10 @@ protected static class CdcRawSegmentFile implements Comparable Date: Sat, 18 Jul 2026 20:29:21 -0700 Subject: [PATCH 2/4] add changes --- CHANGES.txt | 1 + 1 file changed, 1 insertion(+) diff --git a/CHANGES.txt b/CHANGES.txt index d9caa34bb..bc4eee4fe 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 0.5.0 ----- + * Fix the spurious oldest segment age in CdcRawDirectorySpaceCleaner (CASSSIDECAR-484) * Wire CDC configs in configs table to SidecarCdcOptions/SidecarStatePersister (CASSSIDECAR-483) * Implement durable operational job tracker (CASSSIDECAR-374) * Remove filesystem path from Http response (CASSSIDECAR-477) From a3d0be1381ab3694354590001dd4670fe40bb1e8 Mon Sep 17 00:00:00 2001 From: yafeng Date: Sat, 18 Jul 2026 21:33:52 -0700 Subject: [PATCH 3/4] change log to warn level --- .../cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java | 4 +++- .../sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java | 2 +- 2 files changed, 4 insertions(+), 2 deletions(-) diff --git a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java index b1775c21d..fca6c370a 100644 --- a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java +++ b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java @@ -296,7 +296,9 @@ else if (ageMillis < lowMillis) } else { - LOGGER.debug("Skipping delete for already-reclaimed cdc segment {}", segment); + LOGGER.warn("Skipping delete for already-reclaimed cdc segment {}; a concurrent reclaim " + + "(e.g. Cassandra reclaiming cdc_raw) occurred between the directory scan and " + + "the deletion pass", segment); length = segment.length(); } directorySizeBytes -= length; diff --git a/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java b/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java index b79e0eb35..ff2f7c54c 100644 --- a/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java +++ b/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java @@ -182,7 +182,7 @@ void testOldestSegmentAgeNotPoisonedWhenTimestampUnavailable(@TempDir Path tempD // Allow +/-1s to tolerate filesystems that store mtime at second granularity. cleaner.routineCleanUp(); int ageAfterValid = cdcMetrics.oldestSegmentAge.metric.getValue(); - assertThat(ageAfterValid).isBetween(118, 122); + assertThat(ageAfterValid).isBetween(119, 121); // Simulate the oldest segment's timestamp being unavailable (concurrent delete/rotation makes // File.lastModified() return 0). Skip on filesystems that do not support a 0 mtime. From 473a6953b81c48d9a2c1e4a1070bfe8beaa8889b Mon Sep 17 00:00:00 2001 From: yafeng Date: Mon, 27 Jul 2026 10:23:35 -0700 Subject: [PATCH 4/4] use Files.readAttributes for len and mtime --- .../tasks/CdcRawDirectorySpaceCleaner.java | 25 +++++++++++++++---- 1 file changed, 20 insertions(+), 5 deletions(-) diff --git a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java index fca6c370a..efe3f618b 100644 --- a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java +++ b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java @@ -22,6 +22,7 @@ import java.io.IOException; import java.nio.file.Files; import java.nio.file.Path; +import java.nio.file.attribute.BasicFileAttributes; import java.util.Arrays; import java.util.Collections; import java.util.HashMap; @@ -419,9 +420,10 @@ protected static class CdcRawSegmentFile implements Comparable 0 + // with lastModified == 0). A later delete before the cleanup loop is still handled by + // the lastModified > 0 guards. private final long lastModified; CdcRawSegmentFile(File logFile) @@ -429,9 +431,22 @@ protected static class CdcRawSegmentFile implements Comparable 0 guards skip metric emission and deletion. + } + this.len = size; + this.lastModified = mtime; } public boolean exists()