Merge branch 'cassandra-3.0' into cassandra-3.11

This commit is contained in:
Marcus Eriksson 2021-03-18 09:37:50 +01:00
commit 8257868f32
3 changed files with 126 additions and 21 deletions

View File

@ -1,6 +1,7 @@
3.11.11
* Upgrade jackson-databind to 2.9.10.8 (CASSANDRA-16462)
Merged from 3.0:
* Ignore trailing zeros in hint files (CASSANDRA-16523)
* Refuse DROP COMPACT STORAGE if some 2.x sstables are in use (CASSANDRA-15897)
* Fix ColumnFilter::toString not returning a valid CQL fragment (CASSANDRA-16483)
* Fix ColumnFilter behaviour to prevent digest mitmatches during upgrades (CASSANDRA-16415)

View File

@ -36,7 +36,6 @@ import org.apache.cassandra.io.FSReadError;
import org.apache.cassandra.service.StorageService;
import org.apache.cassandra.utils.AbstractIterator;
import org.apache.cassandra.utils.ByteBufferUtil;
import org.apache.cassandra.utils.NativeLibrary;
/**
* A paged non-compressed hints reader that provides two iterators:
@ -212,6 +211,14 @@ class HintsReader implements AutoCloseable, Iterable<HintsReader.Page>
input.resetLimit();
int size = input.readInt();
if (size == 0)
{
// Avoid throwing IOException when a hint file ends with a run of zeros - this
// can happen when hard-rebooting unresponsive machines.
if (!verifyAllZeros(input))
throw new IOException("Corrupt hint file found");
throw new EOFException("Unexpected end of file (size == 0)");
}
// if we cannot corroborate the size via crc, then we cannot safely skip this hint
if (!input.checkCrc())
@ -309,6 +316,14 @@ class HintsReader implements AutoCloseable, Iterable<HintsReader.Page>
input.resetLimit();
int size = input.readInt();
if (size == 0)
{
// Avoid throwing IOException when a hint file ends with a run of zeros - this
// can happen when hard-rebooting unresponsive machines.
if (!verifyAllZeros(input))
throw new IOException("Corrupt hint file found");
throw new EOFException("Unexpected end of file (size == 0)");
}
// if we cannot corroborate the size via crc, then we cannot safely skip this hint
if (!input.checkCrc())
@ -335,4 +350,14 @@ class HintsReader implements AutoCloseable, Iterable<HintsReader.Page>
return null;
}
}
private static boolean verifyAllZeros(ChecksummedDataInput input) throws IOException
{
while (!input.isEOF())
{
if (input.readByte() != 0)
return false;
}
return true;
}
}

View File

@ -22,9 +22,11 @@ import java.io.File;
import java.io.IOException;
import java.nio.ByteBuffer;
import java.nio.file.Files;
import java.nio.file.StandardOpenOption;
import java.util.Iterator;
import java.util.UUID;
import java.util.concurrent.TimeUnit;
import java.util.function.Function;
import com.google.common.collect.Iterables;
import org.junit.BeforeClass;
@ -35,8 +37,11 @@ import org.apache.cassandra.config.CFMetaData;
import org.apache.cassandra.config.Schema;
import org.apache.cassandra.db.Mutation;
import org.apache.cassandra.db.RowUpdateBuilder;
import org.apache.cassandra.db.UnknownColumnFamilyException;
import org.apache.cassandra.db.rows.Cell;
import org.apache.cassandra.db.rows.Row;
import org.apache.cassandra.io.FSReadError;
import org.apache.cassandra.io.util.DataInputBuffer;
import org.apache.cassandra.io.util.FileUtils;
import org.apache.cassandra.schema.KeyspaceParams;
@ -91,43 +96,117 @@ public class HintsReaderTest
}
}
private void readHints(int num, int numTable) throws IOException
private void readHints(int num, int numTable)
{
readAndVerify(num, numTable, HintsReader.Page::hintsIterator);
readAndVerify(num, numTable, this::deserializePageBuffers);
}
private void readAndVerify(int num, int numTable, Function<HintsReader.Page, Iterator<Hint>> getHints)
{
long baseTimestamp = descriptor.timestamp;
int index = 0;
try (HintsReader reader = HintsReader.open(new File(directory, descriptor.fileName())))
{
for (HintsReader.Page page : reader)
{
Iterator<Hint> hints = page.hintsIterator();
Iterator<Hint> hints = getHints.apply(page);
while (hints.hasNext())
{
int i = index / numTable;
Hint hint = hints.next();
long timestamp = baseTimestamp + i;
Mutation mutation = hint.mutation;
assertEquals(timestamp, hint.creationTime);
assertEquals(dk(bytes(i)), mutation.key());
Row row = mutation.getPartitionUpdates().iterator().next().iterator().next();
assertEquals(1, Iterables.size(row.cells()));
assertEquals(bytes(i), row.clustering().get(0));
Cell cell = row.cells().iterator().next();
assertNotNull(cell);
assertEquals(bytes(i), cell.value());
assertEquals(timestamp * 1000, cell.timestamp());
index++;
if (hint != null)
{
verifyHint(hint, baseTimestamp, i);
index++;
}
}
}
}
assertEquals(index, num);
}
private void verifyHint(Hint hint, long baseTimestamp, int i)
{
long timestamp = baseTimestamp + i;
Mutation mutation = hint.mutation;
assertEquals(timestamp, hint.creationTime);
assertEquals(dk(bytes(i)), mutation.key());
Row row = mutation.getPartitionUpdates().iterator().next().iterator().next();
assertEquals(1, Iterables.size(row.cells()));
assertEquals(bytes(i), row.clustering().get(0));
Cell cell = row.cells().iterator().next();
assertNotNull(cell);
assertEquals(bytes(i), cell.value());
assertEquals(timestamp * 1000, cell.timestamp());
}
private Iterator<Hint> deserializePageBuffers(HintsReader.Page page)
{
final Iterator<ByteBuffer> buffers = page.buffersIterator();
return new Iterator<Hint>()
{
public boolean hasNext()
{
return buffers.hasNext();
}
public Hint next()
{
try
{
return Hint.serializer.deserialize(new DataInputBuffer(buffers.next(), false),
descriptor.messagingVersion());
}
catch (UnknownColumnFamilyException e)
{
return null; // ignore
}
catch (IOException e)
{
throw new RuntimeException("Unexpected error deserializing hint", e);
}
};
};
}
@Test
public void corruptFile() throws IOException
{
corruptFileHelper(new byte[100], "corruptFile");
}
@Test(expected = FSReadError.class)
public void corruptFileNotAllZeros() throws IOException
{
byte [] bs = new byte[100];
bs[50] = 1;
corruptFileHelper(bs, "corruptFileNotAllZeros");
}
private void corruptFileHelper(byte[] toAppend, String ks) throws IOException
{
SchemaLoader.createKeyspace(ks,
KeyspaceParams.simple(1),
SchemaLoader.standardCFMD(ks, CF_STANDARD1),
SchemaLoader.standardCFMD(ks, CF_STANDARD2));
int numTable = 2;
directory = Files.createTempDirectory(null).toFile();
try
{
generateHints(3, ks);
File hintFile = new File(directory, descriptor.fileName());
Files.write(hintFile.toPath(), toAppend, StandardOpenOption.APPEND);
readHints(3 * numTable, numTable);
}
finally
{
directory.delete();
}
}
@Test
public void testNormalRead() throws IOException
{