From c163d0bc365239f4960ab2e19fb72a0ff785afa8 Mon Sep 17 00:00:00 2001 From: Stefania Alborghetti Date: Wed, 9 Sep 2015 11:26:14 +0800 Subject: [PATCH] Refactor TransactionLog code and fix order of cleanup bug on Windows patch by stefania; reviewed by benedict for CASSANDRA-10286 --- .../compress/CompressedSequentialWriter.java | 3 +-- .../io/compress/CompressionMetadata.java | 4 +-- .../io/sstable/SSTableTxnWriter.java | 4 +-- .../io/sstable/format/big/BigTableWriter.java | 8 ++---- .../io/util/ChecksummedSequentialWriter.java | 5 +--- .../cassandra/io/util/SequentialWriter.java | 27 +------------------ .../utils/concurrent/Transactional.java | 5 ++-- .../CompressedSequentialWriterTest.java | 1 - .../io/sstable/SSTableLoaderTest.java | 10 +++++-- .../util/ChecksummedSequentialWriterTest.java | 1 - .../io/util/SequentialWriterTest.java | 1 - 11 files changed, 20 insertions(+), 49 deletions(-) diff --git a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java index 8e1ebff018..bbec6f57f3 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java +++ b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java @@ -263,7 +263,7 @@ public class CompressedSequentialWriter extends SequentialWriter @Override protected Throwable doCommit(Throwable accumulate) { - return metadataWriter.commit(accumulate); + return super.doCommit(metadataWriter.commit(accumulate)); } @Override @@ -278,7 +278,6 @@ public class CompressedSequentialWriter extends SequentialWriter syncInternal(); if (descriptor != null) crcMetadata.writeFullChecksum(descriptor); - releaseFileHandle(); sstableMetadataCollector.addCompressionRatio(compressedSize, uncompressedSize); metadataWriter.finalizeLength(current(), chunkCount).prepareToCommit(); } diff --git a/src/java/org/apache/cassandra/io/compress/CompressionMetadata.java b/src/java/org/apache/cassandra/io/compress/CompressionMetadata.java index 1681b0c8ce..04ef2d30f0 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressionMetadata.java +++ b/src/java/org/apache/cassandra/io/compress/CompressionMetadata.java @@ -410,7 +410,7 @@ public class CompressionMetadata count = chunkIndex; } - protected Throwable doPreCleanup(Throwable failed) + protected Throwable doPostCleanup(Throwable failed) { return offsets.close(failed); } @@ -422,7 +422,7 @@ public class CompressionMetadata protected Throwable doAbort(Throwable accumulate) { - return FileUtils.deleteWithConfirm(filePath, false, accumulate); + return accumulate; } } diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableTxnWriter.java b/src/java/org/apache/cassandra/io/sstable/SSTableTxnWriter.java index 6e1ac380b7..5d65a30a31 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableTxnWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableTxnWriter.java @@ -69,13 +69,13 @@ public class SSTableTxnWriter extends Transactional.AbstractTransactional implem protected Throwable doAbort(Throwable accumulate) { - return writer.abort(txn.abort(accumulate)); + return txn.abort(writer.abort(accumulate)); } protected void doPrepare() { - txn.prepareToCommit(); writer.prepareToCommit(); + txn.prepareToCommit(); } public Collection finish(boolean openResult) diff --git a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java index 06dd508fe1..d2500b46ab 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java +++ b/src/java/org/apache/cassandra/io/sstable/format/big/BigTableWriter.java @@ -81,10 +81,6 @@ public class BigTableWriter extends SSTableWriter dbuilder = SegmentedFile.getBuilder(DatabaseDescriptor.getDiskAccessMode(), false); } iwriter = new IndexWriter(keyCount, dataFile); - - // txnLogs will delete if safe to do so (early readers) - iwriter.indexFile.deleteFile(false); - dataFile.deleteFile(false); } public void mark() @@ -322,7 +318,7 @@ public class BigTableWriter extends SSTableWriter } @Override - protected Throwable doPreCleanup(Throwable accumulate) + protected Throwable doPostCleanup(Throwable accumulate) { accumulate = dbuilder.close(accumulate); return accumulate; @@ -485,7 +481,7 @@ public class BigTableWriter extends SSTableWriter } @Override - protected Throwable doPreCleanup(Throwable accumulate) + protected Throwable doPostCleanup(Throwable accumulate) { accumulate = summary.close(accumulate); accumulate = bf.close(accumulate); diff --git a/src/java/org/apache/cassandra/io/util/ChecksummedSequentialWriter.java b/src/java/org/apache/cassandra/io/util/ChecksummedSequentialWriter.java index 8203a377b5..fd88151124 100644 --- a/src/java/org/apache/cassandra/io/util/ChecksummedSequentialWriter.java +++ b/src/java/org/apache/cassandra/io/util/ChecksummedSequentialWriter.java @@ -50,7 +50,7 @@ public class ChecksummedSequentialWriter extends SequentialWriter @Override protected Throwable doCommit(Throwable accumulate) { - return crcWriter.commit(accumulate); + return super.doCommit(crcWriter.commit(accumulate)); } @Override @@ -66,9 +66,6 @@ public class ChecksummedSequentialWriter extends SequentialWriter if (descriptor != null) crcMetadata.writeFullChecksum(descriptor); crcWriter.setDescriptor(descriptor).prepareToCommit(); - // we must cleanup our file handles during prepareCommit for Windows compatibility as we cannot rename an open file; - // TODO: once we stop file renaming, remove this for clarity - releaseFileHandle(); } } diff --git a/src/java/org/apache/cassandra/io/util/SequentialWriter.java b/src/java/org/apache/cassandra/io/util/SequentialWriter.java index 6000f952fc..5bdc15a110 100644 --- a/src/java/org/apache/cassandra/io/util/SequentialWriter.java +++ b/src/java/org/apache/cassandra/io/util/SequentialWriter.java @@ -68,8 +68,6 @@ public class SequentialWriter extends BufferedDataOutputStreamPlus implements Tr // due to lack of multiple-inheritance, we proxy our transactional implementation protected class TransactionalProxy extends AbstractTransactional { - private boolean deleteFile = true; - @Override protected Throwable doPreCleanup(Throwable accumulate) { @@ -90,9 +88,6 @@ public class SequentialWriter extends BufferedDataOutputStreamPlus implements Tr protected void doPrepare() { syncInternal(); - // we must cleanup our file handles during prepareCommit for Windows compatibility as we cannot rename an open file; - // TODO: once we stop file renaming, remove this for clarity - releaseFileHandle(); } protected Throwable doCommit(Throwable accumulate) @@ -102,10 +97,7 @@ public class SequentialWriter extends BufferedDataOutputStreamPlus implements Tr protected Throwable doAbort(Throwable accumulate) { - if (deleteFile) - return FileUtils.deleteWithConfirm(filePath, false, accumulate); - else - return accumulate; + return accumulate; } } @@ -409,23 +401,6 @@ public class SequentialWriter extends BufferedDataOutputStreamPlus implements Tr return new TransactionalProxy(); } - public void deleteFile(boolean val) - { - txnProxy.deleteFile = val; - } - - public void releaseFileHandle() - { - try - { - channel.close(); - } - catch (IOException e) - { - throw new FSWriteError(e, filePath); - } - } - /** * Class to hold a mark to the position of the file */ diff --git a/src/java/org/apache/cassandra/utils/concurrent/Transactional.java b/src/java/org/apache/cassandra/utils/concurrent/Transactional.java index 02562ce5be..d142f06dd2 100644 --- a/src/java/org/apache/cassandra/utils/concurrent/Transactional.java +++ b/src/java/org/apache/cassandra/utils/concurrent/Transactional.java @@ -88,7 +88,8 @@ public interface Transactional extends AutoCloseable // Transactional objects will perform cleanup in the commit() or abort() calls /** - * perform an exception-safe pre-abort cleanup; this will still be run *after* commit + * perform an exception-safe pre-abort/commit cleanup; + * this will be run after prepareToCommit (so before commit), and before abort */ protected Throwable doPreCleanup(Throwable accumulate){ return accumulate; } @@ -113,7 +114,6 @@ public interface Transactional extends AutoCloseable if (state != State.READY_TO_COMMIT) throw new IllegalStateException("Cannot commit unless READY_TO_COMMIT; state is " + state); accumulate = doCommit(accumulate); - accumulate = doPreCleanup(accumulate); accumulate = doPostCleanup(accumulate); state = State.COMMITTED; return accumulate; @@ -171,6 +171,7 @@ public interface Transactional extends AutoCloseable throw new IllegalStateException("Cannot prepare to commit unless IN_PROGRESS; state is " + state); doPrepare(); + maybeFail(doPreCleanup(null)); state = State.READY_TO_COMMIT; } diff --git a/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java b/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java index 1bdc591061..56c83dae82 100644 --- a/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java +++ b/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java @@ -222,7 +222,6 @@ public class CompressedSequentialWriterTest extends SequentialWriterTest protected void assertAborted() throws Exception { super.assertAborted(); - Assert.assertFalse(offsetsFile.exists()); } void cleanup() diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableLoaderTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableLoaderTest.java index faa9c3e049..ad7523d567 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableLoaderTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableLoaderTest.java @@ -131,11 +131,14 @@ public class SSTableLoaderTest writer.addRow("key1", "col1", "100"); } + ColumnFamilyStore cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1); + cfs.forceBlockingFlush(); // wait for sstables to be on disk else we won't be able to stream them + final CountDownLatch latch = new CountDownLatch(1); SSTableLoader loader = new SSTableLoader(dataDir, new TestClient(), new OutputHandler.SystemOutput(false, false)); loader.stream(Collections.emptySet(), completionStreamListener(latch)).get(); - List partitions = Util.getAll(Util.cmd(Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD1)).build()); + List partitions = Util.getAll(Util.cmd(cfs).build()); assertEquals(1, partitions.size()); assertEquals("key1", AsciiType.instance.getString(partitions.get(0).partitionKey().getKey())); @@ -175,6 +178,9 @@ public class SSTableLoaderTest writer.addRow(String.format("key%d", i), String.format("col%d", j), "100"); } + ColumnFamilyStore cfs = Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD2); + cfs.forceBlockingFlush(); // wait for sstables to be on disk else we won't be able to stream them + //make sure we have some tables... assertTrue(dataDir.listFiles().length > 0); @@ -183,7 +189,7 @@ public class SSTableLoaderTest SSTableLoader loader = new SSTableLoader(dataDir, new TestClient(), new OutputHandler.SystemOutput(false, false)); loader.stream(Collections.emptySet(), completionStreamListener(latch)).get(); - List partitions = Util.getAll(Util.cmd(Keyspace.open(KEYSPACE1).getColumnFamilyStore(CF_STANDARD2)).build()); + List partitions = Util.getAll(Util.cmd(cfs).build()); assertTrue(partitions.size() > 0 && partitions.size() < NB_PARTITIONS); diff --git a/test/unit/org/apache/cassandra/io/util/ChecksummedSequentialWriterTest.java b/test/unit/org/apache/cassandra/io/util/ChecksummedSequentialWriterTest.java index 9731a8d320..bea3aac4cd 100644 --- a/test/unit/org/apache/cassandra/io/util/ChecksummedSequentialWriterTest.java +++ b/test/unit/org/apache/cassandra/io/util/ChecksummedSequentialWriterTest.java @@ -85,7 +85,6 @@ public class ChecksummedSequentialWriterTest extends SequentialWriterTest protected void assertAborted() throws Exception { super.assertAborted(); - Assert.assertFalse(crcFile.exists()); } } diff --git a/test/unit/org/apache/cassandra/io/util/SequentialWriterTest.java b/test/unit/org/apache/cassandra/io/util/SequentialWriterTest.java index fd38427977..f5a366ebca 100644 --- a/test/unit/org/apache/cassandra/io/util/SequentialWriterTest.java +++ b/test/unit/org/apache/cassandra/io/util/SequentialWriterTest.java @@ -102,7 +102,6 @@ public class SequentialWriterTest extends AbstractTransactionalTest protected void assertAborted() throws Exception { Assert.assertFalse(writer.isOpen()); - Assert.assertFalse(file.exists()); } protected void assertCommitted() throws Exception