fix(inkless:retention): bound retention enforcement boundary scan to O(cap) - #722
Merged
Merged
Conversation
Contributor
There was a problem hiding this comment.
Pull request overview
This PR improves Inkless retention enforcement scalability by bounding the retention “boundary scan” to O(max_batches_per_request) rather than O(partition depth), preventing long-running scans from exceeding the control plane JDBC socket timeout on very deep partitions. It also flips the default retention.enforcement.max.batches.per.request from 0 (unbounded) to 1000 to enable the bounded behavior by default.
Changes:
- Update
enforce_retention_v2to scan only the oldest(max_batches_per_request + 1)batches (or unbounded when the cap is0), keeping per-pass work bounded. - Change the default retention enforcement cap to
1000and update docs/tests accordingly. - Refresh benchmark guidance and regenerate jOOQ artifacts for schema version 22.
Reviewed changes
Copilot reviewed 6 out of 124 changed files in this pull request and generated no comments.
Show a summary per file
| File | Description |
|---|---|
| storage/inkless/src/test/java/io/aiven/inkless/control_plane/postgres/EnforceRetentionCommitBlockingBenchmarkTest.java | Updates benchmark documentation and allows overriding the per-pass cap via a system property. |
| storage/inkless/src/test/java/io/aiven/inkless/control_plane/AbstractControlPlaneTest.java | Adds a test to lock in correctness when the retention boundary falls within the bounded scan window. |
| storage/inkless/src/test/java/io/aiven/inkless/config/InklessConfigTest.java | Updates expected default for maxBatchesPerEnforcementRequest() to 1000. |
| storage/inkless/src/main/resources/db/migration/V22__Retention_enforcement_bounded_scan.sql | Implements the bounded boundary scan behavior for enforce_retention_v2. |
| storage/inkless/src/main/jooq/org/jooq/generated/UDTs.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/RepairDisklessLogResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/RepairDisklessLogRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/RepairDisklessLogResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/RepairDisklessLogRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/PruneBatchesBelowHighestTieredOffsetResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/PruneBatchesBelowHighestTieredOffsetRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogProducerStateV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchMetadataV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchInfoV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/AdvanceCrossTierLogStartResponseV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/records/AdvanceCrossTierLogStartRequestV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/PruneBatchesBelowHighestTieredOffsetResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/PruneBatchesBelowHighestTieredOffsetRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/RepairDisklessLogResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/RepairDisklessLogRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/PruneBatchesBelowHighestTieredOffsetResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/PruneBatchesBelowHighestTieredOffsetRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogProducerStateV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchMetadataV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchInfoV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/AdvanceCrossTierLogStartResponseV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/AdvanceCrossTierLogStartRequestV1Path.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogProducerStateV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchMetadataV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchInfoV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/AdvanceCrossTierLogStartResponseV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/udt/AdvanceCrossTierLogStartRequestV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/RepairDisklessLogV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/RepairDisklessLogV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/PruneBatchesBelowHighestTieredOffsetV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ProducerStateRecord.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/LogsRecord.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ListOffsetsV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/InitDisklessLogV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV2Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FilesRecord.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV2Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/DeleteRecordsV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/CommitFileV2Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/CommitFileV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/BatchesRecord.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/records/AdvanceCrossTierLogStartV1Record.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/PruneBatchesBelowHighestTieredOffsetV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/ProducerState.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/Logs.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/ListOffsetsV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/InitDisklessLogV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV2.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/Files.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV2.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/DeleteRecordsV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/CommitFileV2.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/CommitFileV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/Batches.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/tables/AdvanceCrossTierLogStartV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/Tables.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/MarkFileToDeleteV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/FlushCommitRunV2.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteTopicV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteFilesV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteBatchV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/routines/BatchTimestamp.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/Routines.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/Keys.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/Indexes.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/PruneBatchesBelowHighestTieredOffsetErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/ListOffsetsResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/InitDisklessLogErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/FindBatchesResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/FileStateT.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/FileReasonT.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/EnforceRetentionResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/DeleteRecordsResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/CommitBatchResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/enums/AdvanceCrossTierLogStartResponseErrorV1.java | Regenerated jOOQ artifact for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/Domains.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/jooq/org/jooq/generated/DefaultSchema.java | Regenerated jOOQ metadata for schema version 22. |
| storage/inkless/src/main/java/io/aiven/inkless/config/InklessConfig.java | Updates config documentation and changes default retention enforcement cap to 1000. |
| docs/inkless/configs.rst | Updates published configuration docs to reflect the new default and behavior. |
…O(cap) The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every deletable batch, O(partition depth). PR #705 moved it out of the log lock so it no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain (retention lowered so almost the whole partition is deletable) a single scan on a multi-million-batch partition runs for seconds-to-minutes and exceeds the enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the still-running backend so the partition never drains. V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The boundary can never be deeper than the cap because the delete is clamped there regardless; the +1 distinguishes an in-window boundary (delete up to it) from a deeper one that falls through to high_watermark and is re-clamped by the existing under-lock cap-probe. Each pass is O(cap), index-only on batches_by_last_offset_covering_idx, and a deep partition drains over successive enforcement cycles. max_batches_per_request = 0 keeps the original unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes oldest-first, so a snapshot-stale boundary can only under-delete). Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the bound is on out of the box; it must stay above the per-interval expiry rate or retention falls behind. Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via -Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms (flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is the sweet spot: short passes at 10x the drain throughput of 100. [KC-354]
jeqo
force-pushed
the
jeqo/retention-enforcer-bound-exec
branch
from
July 28, 2026 08:39
9fcab41 to
508d7fa
Compare
jeqo
marked this pull request as ready for review
July 28, 2026 08:42
viktorsomogyi
self-requested a review
July 28, 2026 08:58
viktorsomogyi
approved these changes
Jul 28, 2026
giuseppelillo
pushed a commit
that referenced
this pull request
Jul 29, 2026
…O(cap) (#722) The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every deletable batch, O(partition depth). PR #705 moved it out of the log lock so it no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain (retention lowered so almost the whole partition is deletable) a single scan on a multi-million-batch partition runs for seconds-to-minutes and exceeds the enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the still-running backend so the partition never drains. V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The boundary can never be deeper than the cap because the delete is clamped there regardless; the +1 distinguishes an in-window boundary (delete up to it) from a deeper one that falls through to high_watermark and is re-clamped by the existing under-lock cap-probe. Each pass is O(cap), index-only on batches_by_last_offset_covering_idx, and a deep partition drains over successive enforcement cycles. max_batches_per_request = 0 keeps the original unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes oldest-first, so a snapshot-stale boundary can only under-delete). Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the bound is on out of the box; it must stay above the per-interval expiry rate or retention falls behind. Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via -Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms (flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is the sweet spot: short passes at 10x the drain throughput of 100. [KC-354]
giuseppelillo
pushed a commit
that referenced
this pull request
Jul 30, 2026
…O(cap) (#722) The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every deletable batch, O(partition depth). PR #705 moved it out of the log lock so it no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain (retention lowered so almost the whole partition is deletable) a single scan on a multi-million-batch partition runs for seconds-to-minutes and exceeds the enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the still-running backend so the partition never drains. V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The boundary can never be deeper than the cap because the delete is clamped there regardless; the +1 distinguishes an in-window boundary (delete up to it) from a deeper one that falls through to high_watermark and is re-clamped by the existing under-lock cap-probe. Each pass is O(cap), index-only on batches_by_last_offset_covering_idx, and a deep partition drains over successive enforcement cycles. max_batches_per_request = 0 keeps the original unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes oldest-first, so a snapshot-stale boundary can only under-delete). Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the bound is on out of the box; it must stay above the per-interval expiry rate or retention falls behind. Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via -Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms (flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is the sweet spot: short passes at 10x the drain throughput of 100. [KC-354]
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
The enforce_retention_v2 boundary scan reverse-aggregates byte_size over every deletable batch, O(partition depth). PR #705 moved it out of the log lock so it no longer blocks commit, but the scan itself stayed O(depth): on a bulk drain (retention lowered so almost the whole partition is deletable) a single scan on a multi-million-batch partition runs for seconds-to-minutes and exceeds the enforcer's socket.timeout.ms (default 5s), aborting the call and orphaning the still-running backend so the partition never drains.
V22 caps the scan to the oldest (max_batches_per_request + 1) batches. The boundary can never be deeper than the cap because the delete is clamped there regardless; the +1 distinguishes an in-window boundary (delete up to it) from a deeper one that falls through to high_watermark and is re-clamped by the existing under-lock cap-probe. Each pass is O(cap), index-only on batches_by_last_offset_covering_idx, and a deep partition drains over successive enforcement cycles. max_batches_per_request = 0 keeps the original unbounded full scan via LIMIT NULL. Correctness is unchanged from V21 (deletes oldest-first, so a snapshot-stale boundary can only under-delete).
Flip retention.enforcement.max.batches.per.request default 0 -> 1000 so the bound is on out of the box; it must stay above the per-interval expiry rate or retention falls behind.
Benchmark (EnforceRetentionCommitBlockingBenchmarkTest, now cap-overridable via -Dinkless.benchmark.maxBatchesPerEnforce). enforce ms vs depth at 200k/400k/800k batches: unbounded 359/729/1465 ms (linear O(depth)); bounded cap=1 21/21/29 ms (flat). A cap sweep shows per-pass cost is flat from cap 1 to 1000, so 1000 is the sweet spot: short passes at 10x the drain throughput of 100.
Diff from previous version without comments:
WHERE topic_id = l_request.topic_id AND partition = l_request.partition ), + limited_batches AS ( + SELECT b.topic_id, b.partition, b.last_offset, b.base_offset, b.byte_size, + batch_timestamp(b.timestamp_type, b.batch_max_timestamp, b.log_append_timestamp) AS effective_timestamp + FROM batches b + WHERE b.topic_id = l_request.topic_id + AND b.partition = l_request.partition + ORDER BY b.topic_id, b.partition, b.last_offset + LIMIT (CASE WHEN max_batches_per_request > 0 THEN max_batches_per_request::bigint + 1 ELSE NULL END) + ), augmented_batches AS ( - SELECT b.topic_id, b.partition, b.last_offset, b.base_offset, + SELECT topic_id, partition, last_offset, base_offset, (SELECT byte_size FROM selected_log) - - SUM(b.byte_size) OVER (ORDER BY b.topic_id, b.partition, b.last_offset) - + b.byte_size + - SUM(byte_size) OVER (ORDER BY topic_id, partition, last_offset) + + byte_size AS reverse_agg_byte_size, - batch_timestamp(b.timestamp_type, b.batch_max_timestamp, b.log_append_timestamp) AS effective_timestamp - FROM batches b - WHERE topic_id = l_request.topic_id - AND partition = l_request.partition - ORDER BY topic_id, partition, last_offset + effective_timestamp + FROM limited_batchesKC-354