Skip to content

feat(inkless): POD-2395 Prune consolidated diskless offsets - #587

Merged
jeqo merged 19 commits into
mainfrom
svv/ts-unification-delete-tiered-diskless
May 27, 2026
Merged

feat(inkless): POD-2395 Prune consolidated diskless offsets#587
jeqo merged 19 commits into
mainfrom
svv/ts-unification-delete-tiered-diskless

Conversation

@viktorsomogyi

Copy link
Copy Markdown
Contributor

Diskless logs which have been already consolidated to the remote
tier should be removed from the coordinator and the WAL. This commit
adds the functionality to do that:

  • control plane (in-memory and postgres) implementation
  • postgres routine
  • wiring into ReplicaManager

The pruning will be invoked as a periodic task in ReplicaManager
with a configurable cleanup period.

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Adds support for pruning diskless WAL batch metadata once it has been consolidated to remote storage, wiring the pruning into the broker so diskless logs can advance their effective start offset and drop already-tiered batch metadata from the control plane.

Changes:

  • Introduces a Postgres routine + control-plane job/API for pruning batches below the highest tiered offset and updating logs.diskless_start_offset.
  • Adds a periodic broker-side pruner (ConsolidatedDisklessLogPruner) scheduled from ReplicaManager to invoke the control-plane pruning and update in-memory partition state.
  • Updates fetch-path log start offset handling and expands test coverage for pruning + fetch overlay behavior.

Reviewed changes

Copilot reviewed 19 out of 136 changed files in this pull request and generated 4 comments.

Show a summary per file
File Description
storage/src/main/java/org/apache/kafka/storage/internals/log/UnifiedLog.java Exposes highestOffsetInRemoteStorage() publicly for cross-module access.
storage/inkless/src/test/java/io/aiven/inkless/control_plane/postgres/PruneBatchesBelowHighestTieredOffsetV1Test.java New PG integration tests for pruning routine semantics.
storage/inkless/src/main/resources/db/migration/V12__Prune_diskless_batches.sql Adds V12 types + pruning function in Postgres.
storage/inkless/src/main/jooq/org/jooq/generated/UDTs.java jOOQ generated updates: adds prune UDT references and schema v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/ReleaseFileMergeWorkItemResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ReleaseFileMergeWorkItemResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ListOffsetsRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/InitDisklessLogProducerStateV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FindBatchesRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseFileV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseBatchV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/EnforceRetentionRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/DeleteRecordsRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitFileMergeWorkItemResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitFileMergeWorkItemBatchV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchResponseV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitBatchRequestV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchMetadataV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/records/BatchInfoV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ReleaseFileMergeWorkItemResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ListOffsetsRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/InitDisklessLogProducerStateV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FindBatchesRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseFileV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseBatchV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/EnforceRetentionRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/DeleteRecordsRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitFileMergeWorkItemResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitFileMergeWorkItemBatchV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchResponseV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitBatchRequestV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchMetadataV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/BatchInfoV1Path.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/ListOffsetsRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/InitDisklessLogProducerStateV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FindBatchesRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseFileV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseBatchV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/EnforceRetentionRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/DeleteRecordsRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitFileMergeWorkItemResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitFileMergeWorkItemBatchV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchResponseV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitBatchRequestV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchMetadataV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/udt/BatchInfoV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ProducerStateRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/LogsRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/ListOffsetsV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/InitDisklessLogV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/GetFileMergeWorkItemV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV2Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FindBatchesV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FilesRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FileMergeWorkItemsRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FileMergeWorkItemFilesRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV2Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/EnforceRetentionV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/DeleteRecordsV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/CommitFileV1Record.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/records/BatchesRecord.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/ProducerState.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Logs.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/ListOffsetsV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/InitDisklessLogV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/GetFileMergeWorkItemV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV2.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FindBatchesV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Files.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FileMergeWorkItems.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/FileMergeWorkItemFiles.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV2.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/EnforceRetentionV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/DeleteRecordsV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/CommitFileV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/tables/Batches.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/Tables.java jOOQ generated updates: adds prune table-function wiring and schema v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/ReleaseFileMergeWorkItemV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/MarkFileToDeleteV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteTopicV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteFilesV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/DeleteBatchV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/CommitFileMergeWorkItemV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/routines/BatchTimestamp.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/Routines.java jOOQ generated updates: adds prune routine wiring and schema v12.
storage/inkless/src/main/jooq/org/jooq/generated/Keys.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/Indexes.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/ReleaseFileMergeWorkItemErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/ListOffsetsResponseErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/InitDisklessLogErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FindBatchesResponseErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FileStateT.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/FileReasonT.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/EnforceRetentionResponseErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/DeleteRecordsResponseErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/CommitFileMergeWorkItemErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/enums/CommitBatchResponseErrorV1.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/Domains.java jOOQ generated schema version bump to v12.
storage/inkless/src/main/jooq/org/jooq/generated/DefaultSchema.java jOOQ generated updates: adds prune objects into schema model and schema v12.
storage/inkless/src/main/java/io/aiven/inkless/delete/PruneDisklessLogsResponse.java New control-plane response record for pruning results.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/PruneDisklessLogsRequest.java New control-plane request record for pruning inputs.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/postgres/PruneDisklessLogsJob.java New PG job calling prune routine and mapping results.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/postgres/PostgresControlPlaneMetrics.java Adds metrics tracking for prune job latency.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/postgres/PostgresControlPlane.java Exposes pruneDisklessLogs via Postgres control plane.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/MetadataView.java Adds topic-id->name lookup and consolidating partition enumeration.
storage/inkless/src/main/java/io/aiven/inkless/control_plane/InMemoryControlPlane.java Adds new API method (currently unimplemented).
storage/inkless/src/main/java/io/aiven/inkless/control_plane/ControlPlane.java Adds new pruneDisklessLogs control-plane API.
core/src/test/scala/io/aiven/inkless/consolidation/DisklessLeaderEndPointTest.scala Adjusts/extends fetch overlay tests for logStartOffset + error behavior.
core/src/test/scala/io/aiven/inkless/consolidation/ConsolidatedDisklessLogPrunerTest.scala New unit tests for pruner request building and update behavior.
core/src/test/java/kafka/server/InklessConsolidatedDisklessTopicsTest.java Adds end-to-end assertions that control-plane WAL metadata is pruned post-tiering.
core/src/main/scala/kafka/server/ReplicaManager.scala Schedules periodic consolidated-diskless pruning task.
core/src/main/scala/kafka/server/metadata/InklessMetadataView.scala Implements new MetadataView APIs for topic name and consolidating partitions.
core/src/main/scala/kafka/cluster/Partition.scala Adds volatile diskless start offset state + setters/getters.
core/src/main/scala/io/aiven/inkless/consolidation/DisklessLeaderEndPoint.scala Overlays local partition logStartOffset into fetch response when appropriate.
core/src/main/scala/io/aiven/inkless/consolidation/ConsolidatedDisklessLogPruner.scala New broker-side pruner invoking control plane and updating partitions.

💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.

Comment thread storage/inkless/src/main/resources/db/migration/V12__Prune_diskless_batches.sql Outdated
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch 6 times, most recently from bc2a937 to e4525a8 Compare May 14, 2026 08:54
@viktorsomogyi viktorsomogyi reopened this May 14, 2026
@viktorsomogyi
viktorsomogyi changed the base branch from main to svv/ts-unification-list-offsets May 14, 2026 08:59
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch from e4525a8 to 1d2bba3 Compare May 14, 2026 09:01
@viktorsomogyi
viktorsomogyi requested a review from Copilot May 14, 2026 09:02

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Pull request overview

Copilot reviewed 20 out of 145 changed files in this pull request and generated 6 comments.

@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch 3 times, most recently from 1a67b69 to f813f46 Compare May 14, 2026 12:35
@viktorsomogyi
viktorsomogyi marked this pull request as ready for review May 14, 2026 12:37
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch 3 times, most recently from 7f6ecd5 to a5347e6 Compare May 14, 2026 15:12
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-list-offsets branch from 0e3b86b to 7d817a1 Compare May 19, 2026 12:11
Base automatically changed from svv/ts-unification-list-offsets to main May 19, 2026 13:02
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch from a5347e6 to 69b7334 Compare May 19, 2026 13:10
Comment thread core/src/main/scala/kafka/cluster/Partition.scala Outdated
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch from 291cecc to 331715d Compare May 26, 2026 08:48

@jeqo jeqo left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Few small suggestions -- mostly LGTM

*/
def isAtMinIsr: Boolean = leaderLogIfLocal.exists { partitionState.isr.size == effectiveMinIsr(_) }

def maybeUpdateDisklessLogStartOffset(newDisklessLogStartOffset: Long): Boolean = {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: should this be synchronized? I know it's called from a single thread, but worth checking/considering

@viktorsomogyi viktorsomogyi May 27, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

given the many other methods use the ISR lock, it might be expected by callers that this does it too, so I think it wouldn't hurt adding it.


Uuid getTopicId(String topicName);

Optional<String> getTopicName(Uuid topicId);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

only used in tests -- worth keeping?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Makes sense, I'll remove it.

case Right(partition) =>
val newDisklessLogStart = pruneDisklessLogsResponse.disklessLogStartOffset
if (!partition.maybeUpdateDisklessLogStartOffset(newDisklessLogStart)) {
logger.error("Diskless log start offset is non-monotonic. The old one ({}) is greater than the new ({}).")

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

missing args?

Suggested change
logger.error("Diskless log start offset is non-monotonic. The old one ({}) is greater than the new ({}).")
logger.error("Diskless log start offset is non-monotonic for {}. The new value ({}) is not greater than current.",
pruneDisklessLogsResponse.topicIdPartition.topicPartition,
newDisklessLogStart: java.lang.Long)

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What if I just moved this logging inside Partition? That way we could log both offsets.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thinking more about it, logging an error here doesn't make sense, we should enforce this sooner in the control plane and log an error there if it doesn't work. So I'll probably refactor this.

@viktorsomogyi viktorsomogyi May 27, 2026

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Refactored this and ended up moving this method into Partition. Now both postgres and the in-memory implementations are monotonic, so logging here only means that we prevented a stale update going through (so I recuced it to warn too).

Consolidated diskless log pruning will be introduced on the interface
level too to define a common API for various control plane implementations.
Adds consolidation.cleanup.interval.ms to be able to control the
log pruning interval. This will be used to start a periodic
cleanup process to remove diskless batches that have been moved
to tiered storage already.
ConsolidatedDisklessLogPruner implements the cleanup logic that is
common to all control plane implementations and utilizes the control
plane to start the cleanup process.

# Conflicts:
#	core/src/main/scala/kafka/server/ReplicaManager.scala

# Conflicts:
#	core/src/main/scala/kafka/server/ReplicaManager.scala
Adds the diskless pruner implementation in the Postgres control plane.

� Conflicts:
�	storage/inkless/src/main/jooq/org/jooq/generated/DefaultSchema.java
�	storage/inkless/src/main/jooq/org/jooq/generated/Routines.java
�	storage/inkless/src/main/jooq/org/jooq/generated/Tables.java
�	storage/inkless/src/main/jooq/org/jooq/generated/UDTs.java
�	storage/inkless/src/main/jooq/org/jooq/generated/enums/CommitFileMergeWorkItemErrorV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/enums/ReleaseFileMergeWorkItemErrorV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/routines/CommitFileMergeWorkItemV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/routines/ReleaseFileMergeWorkItemV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/FileMergeWorkItemFiles.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/FileMergeWorkItems.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/GetFileMergeWorkItemV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FileMergeWorkItemFilesRecord.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/records/FileMergeWorkItemsRecord.java
�	storage/inkless/src/main/jooq/org/jooq/generated/tables/records/GetFileMergeWorkItemV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitFileMergeWorkItemBatchV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/CommitFileMergeWorkItemResponseV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseBatchV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseFileV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/FileMergeWorkItemResponseV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/ReleaseFileMergeWorkItemResponseV1.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitFileMergeWorkItemBatchV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/CommitFileMergeWorkItemResponseV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseBatchV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseFileV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/FileMergeWorkItemResponseV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/paths/ReleaseFileMergeWorkItemResponseV1Path.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitFileMergeWorkItemBatchV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/CommitFileMergeWorkItemResponseV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseBatchV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseFileV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/FileMergeWorkItemResponseV1Record.java
�	storage/inkless/src/main/jooq/org/jooq/generated/udt/records/ReleaseFileMergeWorkItemResponseV1Record.java
Adds extra verification to the consolidation test to verify
that diskless pruning happens correctly.
Co-authored-by: Viktor Somogyi-Vass <viktorsomogyi@gmail.com>
@viktorsomogyi
viktorsomogyi force-pushed the svv/ts-unification-delete-tiered-diskless branch from 331715d to 5f47cc2 Compare May 27, 2026 14:57
@viktorsomogyi

Copy link
Copy Markdown
Contributor Author

@jeqo I have applied your suggestions + rebased the commit on the latest main. The most recent changes have a few other updates too:

  • ensure monotonicity of log start offset in postgres
  • update byte_size after pruning
  • address logging problem with diskless start offset
  • construct pruner only when consolidation is enabled

@viktorsomogyi
viktorsomogyi requested a review from jeqo May 27, 2026 15:03
@jeqo
jeqo merged commit 22144c9 into main May 27, 2026
5 checks passed
@jeqo
jeqo deleted the svv/ts-unification-delete-tiered-diskless branch May 27, 2026 19:20
giuseppelillo pushed a commit that referenced this pull request May 29, 2026
* feat(inkless): POD-2395 Control plane interface changes for pruner

Consolidated diskless log pruning will be introduced on the interface
level too to define a common API for various control plane implementations.

* feat(inkless): POD-2395 Add new config to control pruning interval

Adds consolidation.cleanup.interval.ms to be able to control the
log pruning interval. This will be used to start a periodic
cleanup process to remove diskless batches that have been moved
to tiered storage already.

* feat(inkless): POD-2395 Add the diskless pruner implementation

ConsolidatedDisklessLogPruner implements the cleanup logic that is
common to all control plane implementations and utilizes the control
plane to start the cleanup process.

* feat(inkless): POD-2395 Add postgres pruner implementation

Adds the diskless pruner implementation in the Postgres control plane.

* feat(inkless): POD-2395 Add pruner integration test

Adds extra verification to the consolidation test to verify
that diskless pruning happens correctly.

* feat(inkless): POD-2395 Diskless pruning in-memory implementation

* Address comments in the pruner implementation

* Address comments in InMemoryControlPlane

* Address comments in PostgresControlPlane

* Enfore diskless start offset monotonicity in Partition

* Equal offsets should be allowed in maybeUpdateDisklessLogStartOffsets

Co-authored-by: Viktor Somogyi-Vass <viktorsomogyi@gmail.com>

* Remove getTopicName from MetadataView

* Add lock to maybeUpdateDisklessLogStartOffset

* Fine-tune error logging

* Ensure monotonicity for pruning

* Reduce byte size in postgres after pruning

* Construct pruner only when consolidation is enabled

* Update the integration test to accept already pruned log

* Refactor diskless start offset advancing
giuseppelillo pushed a commit that referenced this pull request May 29, 2026
* feat(inkless): POD-2395 Control plane interface changes for pruner

Consolidated diskless log pruning will be introduced on the interface
level too to define a common API for various control plane implementations.

* feat(inkless): POD-2395 Add new config to control pruning interval

Adds consolidation.cleanup.interval.ms to be able to control the
log pruning interval. This will be used to start a periodic
cleanup process to remove diskless batches that have been moved
to tiered storage already.

* feat(inkless): POD-2395 Add the diskless pruner implementation

ConsolidatedDisklessLogPruner implements the cleanup logic that is
common to all control plane implementations and utilizes the control
plane to start the cleanup process.

* feat(inkless): POD-2395 Add postgres pruner implementation

Adds the diskless pruner implementation in the Postgres control plane.

* feat(inkless): POD-2395 Add pruner integration test

Adds extra verification to the consolidation test to verify
that diskless pruning happens correctly.

* feat(inkless): POD-2395 Diskless pruning in-memory implementation

* Address comments in the pruner implementation

* Address comments in InMemoryControlPlane

* Address comments in PostgresControlPlane

* Enfore diskless start offset monotonicity in Partition

* Equal offsets should be allowed in maybeUpdateDisklessLogStartOffsets

Co-authored-by: Viktor Somogyi-Vass <viktorsomogyi@gmail.com>

* Remove getTopicName from MetadataView

* Add lock to maybeUpdateDisklessLogStartOffset

* Fine-tune error logging

* Ensure monotonicity for pruning

* Reduce byte size in postgres after pruning

* Construct pruner only when consolidation is enabled

* Update the integration test to accept already pruned log

* Refactor diskless start offset advancing
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants