feat(inkless): POD-2395 Prune consolidated diskless offsets - #587
Conversation
There was a problem hiding this comment.
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 fromReplicaManagerto 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.
bc2a937 to
e4525a8
Compare
e4525a8 to
1d2bba3
Compare
1a67b69 to
f813f46
Compare
7f6ecd5 to
a5347e6
Compare
0e3b86b to
7d817a1
Compare
a5347e6 to
69b7334
Compare
291cecc to
331715d
Compare
jeqo
left a comment
There was a problem hiding this comment.
Few small suggestions -- mostly LGTM
| */ | ||
| def isAtMinIsr: Boolean = leaderLogIfLocal.exists { partitionState.isr.size == effectiveMinIsr(_) } | ||
|
|
||
| def maybeUpdateDisklessLogStartOffset(newDisklessLogStartOffset: Long): Boolean = { |
There was a problem hiding this comment.
nit: should this be synchronized? I know it's called from a single thread, but worth checking/considering
There was a problem hiding this comment.
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); |
There was a problem hiding this comment.
only used in tests -- worth keeping?
There was a problem hiding this comment.
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 ({}).") |
There was a problem hiding this comment.
missing args?
| 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) | |
There was a problem hiding this comment.
What if I just moved this logging inside Partition? That way we could log both offsets.
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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>
331715d to
5f47cc2
Compare
|
@jeqo I have applied your suggestions + rebased the commit on the latest main. The most recent changes have a few other updates too:
|
* 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
* 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
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:
The pruning will be invoked as a periodic task in ReplicaManager
with a configurable cleanup period.