HesandaLiyanage commented on code in PR #3193: URL: https://github.com/apache/james-project/pull/3193#discussion_r4070313249
########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/BlobCompactionAlgorithm.java: ########## @@ -0,0 +1,617 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobReferenceSource; +import org.apache.james.blob.api.BlobStoreDAO; +import org.apache.james.blob.api.ObjectStoreIOException; +import org.apache.james.blob.compaction.BlobReferenceMappingSource.BlobIdMessageIdMapping; +import org.apache.james.blob.compaction.ChunkFormat.BlobSlotContent; +import org.apache.james.blob.compaction.ChunkFormat.ChunkWriteResult; +import org.apache.james.blob.compaction.ChunkFormat.SlotRange; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.base.Preconditions; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + +/** + * Executes object compaction and garbage collection for chunked blob storage in Apache James. + * + * <h3>Crash Safety and Liveness Invariants</h3> + * The compaction algorithms follow a strict step ordering to ensure crash safety and liveness: + * <ul> + * <li><b>Initial Compaction:</b> + * <ol> + * <li>Read candidates and accumulate slots into an immutable chunk.</li> + * <li>Save the new chunk to raw storage (unreferenced yet by source-of-truth tables).</li> + * <li>Update source-of-truth table references (Cassandra) to the new chunk-slot references.</li> + * <li>Delete original standalone objects from raw storage.</li> + * </ol> + * If interrupted before step 3, the saved chunk is an orphan with 0 references and will be + * safely reclaimed by the next {@code gc-compact} run. The original blobs and references remain untouched. + * If interrupted after step 3 but before step 4, the old standalone blobs have 0 references and will be + * cleaned up by standard GC. + * </li> + * <li><b>GC-Compact (Rewrite / Merge / Purge):</b> + * <ol> + * <li>Read existing chunk objects and inspect slot references against the BloomFilter/mapping.</li> + * <li>For chunks with dead slots (or pairs of small chunks to merge), assemble a new chunk with only live slots.</li> + * <li>Save the new chunk to raw storage.</li> + * <li>Update source-of-truth table references to the new chunk slot references.</li> + * <li>Delete the old chunk(s) from raw storage.</li> + * </ol> + * If interrupted before step 4, the newly created chunk is an orphan with no references and will be + * purged on the subsequent GC-compact pass. If interrupted after step 4, the old chunk has 0 references + * and will be purged as an orphan chunk on the next GC-compact pass. + * </li> + * <li><b>Orphan Chunk Purging:</b> + * Chunks with 0 live references (100% dead slots) are identified as orphan chunks and deleted immediately. + * This guarantees self-healing and liveness by construction across process crashes. + * </li> + * </ul> + * + * <h3>Memory Bounds and Operational Characteristics</h3> + * <ul> + * <li><b>Candidate Payload Streaming:</b> {@link #initialCompact(CompactionRequest)} streams candidate blob identifiers + * and partitions them into windows of {@value #DEFAULT_CANDIDATE_BATCH_SIZE} blobs. Candidate payloads are fetched + * and packed chunk-by-chunk. Payloads are persisted and freed window-by-window, ensuring that candidate byte arrays + * are never held in heap for the entire generation simultaneously. Candidate payload heap usage is bounded by + * {@code O(min(candidateBatchSize * avgBlobSize, chunkTargetSize))}.</li> + * <li><b>Reference Mapping Memory Ceiling:</b> {@code loadReferenceMapping()} materializes all live blob-to-messageId + * mappings for the generation into an in-memory multimap. Memory consumption is {@code O(liveGenerationReferences)} + * at approximately ~200 bytes per reference (~200MB heap for 1 million live references; ~2GB heap for 10 million). + * Because {@link BlobReferenceMappingSource} currently exposes a full stream without partition-paged query capabilities, + * this table is loaded per compaction pass. High-scale deployments exceeding tens of millions of live references per + * generation can introduce partition-paged reference lookups in future iterations.</li> + * <li><b>GC Compaction Memory Bounds:</b> {@link #gcCompact(CompactionRequest)} discovers chunks by reading only trailing + * 64KB footers via HTTP ranged reads (metadata-only). Orphan chunks (100% dead slots) are deleted with 0 payload bytes read. + * During chunk purge or merge, surviving live slots are streamed individually via HTTP ranged reads, strictly bounding + * GC payload heap usage to {@code O(maxSlotSize)} (~1MB).</li> + * </ul> + */ +public class BlobCompactionAlgorithm { + public static final int DEFAULT_CANDIDATE_BATCH_SIZE = 1000; + private static final Logger LOGGER = LoggerFactory.getLogger(BlobCompactionAlgorithm.class); + + private final BlobStoreDAO blobStoreDAO; + private final BlobStoreDAO rawStore; + private final BlobReferenceSource referenceSource; + private final BlobReferenceMappingSource mappingSource; + private final BlobIdUpdater blobIdUpdater; + + public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO, + BlobStoreDAO rawStore, + BlobReferenceSource referenceSource, + BlobReferenceMappingSource mappingSource, + BlobIdUpdater blobIdUpdater) { + this.blobStoreDAO = Preconditions.checkNotNull(blobStoreDAO, "'blobStoreDAO' must not be null"); + this.rawStore = Preconditions.checkNotNull(rawStore, "'rawStore' must not be null"); + this.referenceSource = Preconditions.checkNotNull(referenceSource, "'referenceSource' must not be null"); + this.mappingSource = Preconditions.checkNotNull(mappingSource, "'mappingSource' must not be null"); + this.blobIdUpdater = Preconditions.checkNotNull(blobIdUpdater, "'blobIdUpdater' must not be null"); + } + + public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO, + BlobStoreDAO rawStore, + BlobReferenceMappingSource mappingSource, + BlobIdUpdater blobIdUpdater) { + this(blobStoreDAO, rawStore, + () -> Flux.from(mappingSource.listBlobIdMessageIdMappings()).map(BlobIdMessageIdMapping::blobId), + mappingSource, blobIdUpdater); + } + + public Mono<CompactionResult> compact(CompactionRequest request) { + return initialCompact(request) + .flatMap(initialResult -> gcCompact(request).map(initialResult::combine)); + } + + public Mono<CompactionResult> initialCompact(CompactionRequest request) { + Preconditions.checkNotNull(request, "'request' must not be null"); + + return loadReferenceMapping() + .flatMap(mapping -> { + if (mapping.isEmpty()) { + LOGGER.info("No blob references found in mapping source; skipping initial compaction for generation {}", request.generation()); + return Mono.just(CompactionResult.NONE); + } + + return Flux.from(rawStore.listBlobs(request.bucketName())) + .filter(blobId -> matchesGenerationAndFamily(blobId.asString(), request.generation(), request.family())) Review Comment: Done. Added generation prefix pushdown when listing blobs using `BlobStoreDAO.listBlobs(bucket, prefix)` with fallback to post-filtering when prefix listing is unsupported. ########## server/blob/blob-compaction/src/main/java/org/apache/james/blob/compaction/BlobCompactionAlgorithm.java: ########## @@ -0,0 +1,617 @@ +/**************************************************************** + * Licensed to the Apache Software Foundation (ASF) under one * + * or more contributor license agreements. See the NOTICE file * + * distributed with this work for additional information * + * regarding copyright ownership. The ASF licenses this file * + * to you under the Apache License, Version 2.0 (the * + * "License"); you may not use this file except in compliance * + * with the License. You may obtain a copy of the License at * + * * + * http://www.apache.org/licenses/LICENSE-2.0 * + * * + * Unless required by applicable law or agreed to in writing, * + * software distributed under the License is distributed on an * + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY * + * KIND, either express or implied. See the License for the * + * specific language governing permissions and limitations * + * under the License. * + ****************************************************************/ + +package org.apache.james.blob.compaction; + +import java.io.IOException; +import java.util.ArrayList; +import java.util.Collection; +import java.util.HashMap; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Optional; +import java.util.Set; + +import org.apache.james.blob.api.BlobId; +import org.apache.james.blob.api.BlobReferenceSource; +import org.apache.james.blob.api.BlobStoreDAO; +import org.apache.james.blob.api.ObjectStoreIOException; +import org.apache.james.blob.compaction.BlobReferenceMappingSource.BlobIdMessageIdMapping; +import org.apache.james.blob.compaction.ChunkFormat.BlobSlotContent; +import org.apache.james.blob.compaction.ChunkFormat.ChunkWriteResult; +import org.apache.james.blob.compaction.ChunkFormat.SlotRange; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import com.google.common.base.Preconditions; + +import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + +/** + * Executes object compaction and garbage collection for chunked blob storage in Apache James. + * + * <h3>Crash Safety and Liveness Invariants</h3> + * The compaction algorithms follow a strict step ordering to ensure crash safety and liveness: + * <ul> + * <li><b>Initial Compaction:</b> + * <ol> + * <li>Read candidates and accumulate slots into an immutable chunk.</li> + * <li>Save the new chunk to raw storage (unreferenced yet by source-of-truth tables).</li> + * <li>Update source-of-truth table references (Cassandra) to the new chunk-slot references.</li> + * <li>Delete original standalone objects from raw storage.</li> + * </ol> + * If interrupted before step 3, the saved chunk is an orphan with 0 references and will be + * safely reclaimed by the next {@code gc-compact} run. The original blobs and references remain untouched. + * If interrupted after step 3 but before step 4, the old standalone blobs have 0 references and will be + * cleaned up by standard GC. + * </li> + * <li><b>GC-Compact (Rewrite / Merge / Purge):</b> + * <ol> + * <li>Read existing chunk objects and inspect slot references against the BloomFilter/mapping.</li> + * <li>For chunks with dead slots (or pairs of small chunks to merge), assemble a new chunk with only live slots.</li> + * <li>Save the new chunk to raw storage.</li> + * <li>Update source-of-truth table references to the new chunk slot references.</li> + * <li>Delete the old chunk(s) from raw storage.</li> + * </ol> + * If interrupted before step 4, the newly created chunk is an orphan with no references and will be + * purged on the subsequent GC-compact pass. If interrupted after step 4, the old chunk has 0 references + * and will be purged as an orphan chunk on the next GC-compact pass. + * </li> + * <li><b>Orphan Chunk Purging:</b> + * Chunks with 0 live references (100% dead slots) are identified as orphan chunks and deleted immediately. + * This guarantees self-healing and liveness by construction across process crashes. + * </li> + * </ul> + * + * <h3>Memory Bounds and Operational Characteristics</h3> + * <ul> + * <li><b>Candidate Payload Streaming:</b> {@link #initialCompact(CompactionRequest)} streams candidate blob identifiers + * and partitions them into windows of {@value #DEFAULT_CANDIDATE_BATCH_SIZE} blobs. Candidate payloads are fetched + * and packed chunk-by-chunk. Payloads are persisted and freed window-by-window, ensuring that candidate byte arrays + * are never held in heap for the entire generation simultaneously. Candidate payload heap usage is bounded by + * {@code O(min(candidateBatchSize * avgBlobSize, chunkTargetSize))}.</li> + * <li><b>Reference Mapping Memory Ceiling:</b> {@code loadReferenceMapping()} materializes all live blob-to-messageId + * mappings for the generation into an in-memory multimap. Memory consumption is {@code O(liveGenerationReferences)} + * at approximately ~200 bytes per reference (~200MB heap for 1 million live references; ~2GB heap for 10 million). + * Because {@link BlobReferenceMappingSource} currently exposes a full stream without partition-paged query capabilities, + * this table is loaded per compaction pass. High-scale deployments exceeding tens of millions of live references per + * generation can introduce partition-paged reference lookups in future iterations.</li> + * <li><b>GC Compaction Memory Bounds:</b> {@link #gcCompact(CompactionRequest)} discovers chunks by reading only trailing + * 64KB footers via HTTP ranged reads (metadata-only). Orphan chunks (100% dead slots) are deleted with 0 payload bytes read. + * During chunk purge or merge, surviving live slots are streamed individually via HTTP ranged reads, strictly bounding + * GC payload heap usage to {@code O(maxSlotSize)} (~1MB).</li> + * </ul> + */ +public class BlobCompactionAlgorithm { + public static final int DEFAULT_CANDIDATE_BATCH_SIZE = 1000; + private static final Logger LOGGER = LoggerFactory.getLogger(BlobCompactionAlgorithm.class); + + private final BlobStoreDAO blobStoreDAO; + private final BlobStoreDAO rawStore; + private final BlobReferenceSource referenceSource; + private final BlobReferenceMappingSource mappingSource; + private final BlobIdUpdater blobIdUpdater; + + public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO, + BlobStoreDAO rawStore, + BlobReferenceSource referenceSource, + BlobReferenceMappingSource mappingSource, + BlobIdUpdater blobIdUpdater) { + this.blobStoreDAO = Preconditions.checkNotNull(blobStoreDAO, "'blobStoreDAO' must not be null"); + this.rawStore = Preconditions.checkNotNull(rawStore, "'rawStore' must not be null"); + this.referenceSource = Preconditions.checkNotNull(referenceSource, "'referenceSource' must not be null"); + this.mappingSource = Preconditions.checkNotNull(mappingSource, "'mappingSource' must not be null"); + this.blobIdUpdater = Preconditions.checkNotNull(blobIdUpdater, "'blobIdUpdater' must not be null"); + } + + public BlobCompactionAlgorithm(BlobStoreDAO blobStoreDAO, + BlobStoreDAO rawStore, + BlobReferenceMappingSource mappingSource, + BlobIdUpdater blobIdUpdater) { + this(blobStoreDAO, rawStore, + () -> Flux.from(mappingSource.listBlobIdMessageIdMappings()).map(BlobIdMessageIdMapping::blobId), + mappingSource, blobIdUpdater); + } + + public Mono<CompactionResult> compact(CompactionRequest request) { + return initialCompact(request) + .flatMap(initialResult -> gcCompact(request).map(initialResult::combine)); + } + + public Mono<CompactionResult> initialCompact(CompactionRequest request) { + Preconditions.checkNotNull(request, "'request' must not be null"); + + return loadReferenceMapping() + .flatMap(mapping -> { + if (mapping.isEmpty()) { + LOGGER.info("No blob references found in mapping source; skipping initial compaction for generation {}", request.generation()); + return Mono.just(CompactionResult.NONE); + } + + return Flux.from(rawStore.listBlobs(request.bucketName())) + .filter(blobId -> matchesGenerationAndFamily(blobId.asString(), request.generation(), request.family())) + .filter(blobId -> !ChunkId.isChunkRef(blobId)) + .filter(mapping::containsKey) + .window(DEFAULT_CANDIDATE_BATCH_SIZE) Review Comment: Done. Windowing is now based on cumulative byte size up to `chunkTargetSize` rather than pure item count. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
