Fix MAX_SEGMENT_SIZE < chunkSize in MmappedRegions::updateState

patch by samlightfoot, Szymon Mi<C4><99><C5><BC>a<C5><82>; reviewed by Ariel Weisberg, Michael Semb Wever for CASSANDRA-20636
This commit is contained in:
samlightfoot 2025-06-09 13:22:45 -04:00 committed by Ariel Weisberg
parent f07510a303
commit 53ccfd58f2
3 changed files with 64 additions and 18 deletions

View File

@ -1,4 +1,5 @@
5.0.5
* Fix MAX_SEGMENT_SIZE < chunkSize in MmappedRegions::updateState (CASSANDRA-20636)
* Full Java 17 support (CASSANDRA-20681)
* Ensure replica filtering protection does not trigger unnecessary short read protection reads (CASSANDRA-20639)
* Unified Compaction does not properly validate min and target sizes (CASSANDRA-20398)

View File

@ -183,7 +183,7 @@ public class MmappedRegions extends SharedCloseableImpl
private void updateState(long length, int chunkSize)
{
// make sure the regions span whole chunks
long maxSize = (long) (MAX_SEGMENT_SIZE / chunkSize) * chunkSize;
long maxSize = Math.max(chunkSize, (long) (MAX_SEGMENT_SIZE / chunkSize) * chunkSize);
state.length = length;
long pos = state.getPosition();
while (pos < length)

View File

@ -30,12 +30,14 @@ import java.util.concurrent.Future;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import java.util.stream.Collectors;
import java.util.stream.IntStream;
import java.util.stream.Stream;
import com.google.common.collect.ImmutableList;
import com.google.common.collect.Sets;
import org.junit.After;
import org.junit.Assume;
import org.junit.BeforeClass;
import org.junit.Rule;
@ -88,6 +90,9 @@ import org.apache.cassandra.schema.CompressionParams;
import org.apache.cassandra.schema.KeyspaceParams;
import org.apache.cassandra.service.CacheService;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.Throwables;
import org.apache.cassandra.utils.concurrent.Ref;
import org.apache.cassandra.utils.concurrent.SelfRefCounted;
import org.mockito.Mockito;
import static java.lang.String.format;
@ -97,10 +102,12 @@ import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
import static org.junit.Assert.fail;
import static org.junit.Assume.assumeTrue;
public class SSTableReaderTest
{
public static final String KEYSPACE1 = "SSTableReaderTest";
public static final String CF_STANDARD = "Standard1";
public static final String CF_STANDARD2 = "Standard2";
@ -111,6 +118,7 @@ public class SSTableReaderTest
public static final String CF_STANDARD_LOW_INDEX_INTERVAL = "StandardLowIndexInterval";
public static final String CF_STANDARD_SMALL_BLOOM_FILTER = "StandardSmallBloomFilter";
private final List<Ref<?>> refsToRelease = new ArrayList<>();
private IPartitioner partitioner;
Token t(int i)
@ -143,6 +151,28 @@ public class SSTableReaderTest
CompactionManager.instance.disableAutoCompaction();
}
@After
public void teardown()
{
Throwable exceptions = null;
for (Ref<?> ref : refsToRelease)
{
try
{
ref.release();
}
catch (Throwable exc)
{
exceptions = Throwables.merge(exceptions, exc);
}
}
if (exceptions != null)
fail("Unable to release all tracked references " + exceptions);
refsToRelease.clear();
}
@Test
public void testGetPositionsForRanges()
{
@ -310,7 +340,6 @@ public class SSTableReaderTest
}
}
@Test
public void testOnDiskSizeCompressedBoundaries()
{
@ -344,7 +373,6 @@ public class SSTableReaderTest
onDiskSizeForRanges(sstable, Collections.singleton(new Range<>(partitioner.getMinimumToken(), t0(k - 1))))); // inclusive end
}
long onDiskSizeForRanges(SSTableReader sstable, Collection<Range<Token>> ranges)
{
return sstable.onDiskSizeForPartitionPositions(sstable.getPositionsForRanges(ranges));
@ -375,12 +403,22 @@ public class SSTableReaderTest
return s.substring(0, s.length() - n);
}
@Test
public void testSpannedIndexPositions() throws IOException
{
doTestSpannedIndexPositions(PageAware.PAGE_SIZE);
}
@Test
public void testSpannedIndexPositionsWithMaxSegmentSizeSmallerThanPageSize() throws IOException
{
doTestSpannedIndexPositions(PageAware.PAGE_SIZE - 1);
}
public void doTestSpannedIndexPositions(int maxSegmentSize) throws IOException
{
int originalMaxSegmentSize = MmappedRegions.MAX_SEGMENT_SIZE;
MmappedRegions.MAX_SEGMENT_SIZE = PageAware.PAGE_SIZE;
MmappedRegions.MAX_SEGMENT_SIZE = maxSegmentSize;
try
{
@ -655,7 +693,6 @@ public class SSTableReaderTest
assertEquals(fpCount + 2, sstable.getFilterTracker().getFalsePositiveCount());
}
@Test
public void testGetPositionsListenerCalls()
{
@ -775,12 +812,12 @@ public class SSTableReaderTest
Util.flush(store);
CompactionManager.instance.performMaximal(store, false);
SSTableReaderWithFilter sstable = (SSTableReaderWithFilter) store.getLiveSSTables().iterator().next();
sstable = (SSTableReaderWithFilter) sstable.cloneWithNewStart(dk(3));
return sstable;
return (SSTableReaderWithFilter) trackReleaseableRef(() -> {
SSTableReader reader = store.getLiveSSTables().iterator().next();
return reader.cloneWithNewStart(dk(3));
});
}
@Test
public void testOpeningSSTable() throws Exception
{
@ -895,7 +932,8 @@ public class SSTableReaderTest
components.remove(Components.SUMMARY);
target = SSTableReader.open(store, desc, components, store.metadata);
try {
try
{
assertTrue("Bloomfilter was not recreated", bloomModified < bloomFile.lastModified());
assertTrue("Summary was not recreated", summaryModified < summaryFile.lastModified());
}
@ -1054,9 +1092,8 @@ public class SSTableReaderTest
if (sstable instanceof IndexSummarySupport<?>)
{
new IndexSummaryComponent(((IndexSummarySupport<?>) sstable).getIndexSummary(), sstable.getFirst(), sstable.getLast()).save(sstable.descriptor.fileFor(Components.SUMMARY), true);
SSTableReader reopened = SSTableReader.open(store, sstable.descriptor);
SSTableReader reopened = trackReleaseableRef(() -> SSTableReader.open(store, sstable.descriptor));
assert reopened.getFirst().getToken() instanceof LocalToken;
reopened.selfRef().release();
}
}
}
@ -1065,7 +1102,7 @@ public class SSTableReaderTest
* see CASSANDRA-5407
*/
@Test
public void testGetScannerForNoIntersectingRanges() throws Exception
public void testGetScannerForNoIntersectingRanges()
{
ColumnFamilyStore store = discardSSTables(KEYSPACE1, CF_STANDARD);
partitioner = store.getPartitioner();
@ -1240,7 +1277,7 @@ public class SSTableReaderTest
txn.update(replacement, true);
txn.finish();
}
R reopen = (R) SSTableReader.open(store, sstable.descriptor);
R reopen = (R) trackReleaseableRef(() -> SSTableReader.open(store, sstable.descriptor));
assert reopen.getIndexSummary().getSamplingLevel() == sstable.getIndexSummary().getSamplingLevel() + 1;
}
@ -1311,9 +1348,9 @@ public class SSTableReaderTest
{
Keyspace keyspace = Keyspace.open(KEYSPACE1);
ColumnFamilyStore cfs = keyspace.getColumnFamilyStore(CF_MOVE_AND_OPEN);
SSTableReader sstable = getNewSSTable(cfs);
SSTableReader sstable = trackReleaseableRef(() -> getNewSSTable(cfs));
cfs.clearUnsafe();
sstable.selfRef().release();
File tmpdir = new File(Files.createTempDirectory("testMoveAndOpen"));
tmpdir.deleteOnExit();
SSTableId id = SSTableIdFactory.instance.defaultBuilder().generator(Stream.empty()).get();
@ -1325,7 +1362,8 @@ public class SSTableReaderTest
assertFalse(f.exists());
assertTrue(sstable.descriptor.fileFor(c).exists());
}
SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, notLiveDesc, sstable.components, false);
trackReleaseableRef(() -> SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, notLiveDesc, sstable.components, false));
// make sure the files were moved:
for (Component c : sstable.components)
{
@ -1435,4 +1473,11 @@ public class SSTableReaderTest
cfs.discardSSTables(System.currentTimeMillis());
return cfs;
}
private <T extends SelfRefCounted<T>> T trackReleaseableRef(Supplier<T> refSupplier)
{
T ref = refSupplier.get();
refsToRelease.add(ref.selfRef());
return ref;
}
}