mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-3.0' into cassandra-3.11
This commit is contained in:
commit
ffd3d412c2
|
|
@ -30,6 +30,7 @@ import java.util.function.Predicate;
|
|||
|
||||
import com.google.common.annotations.VisibleForTesting;
|
||||
import com.google.common.util.concurrent.Runnables;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
|
|
@ -109,6 +110,8 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
|
||||
private final Tracker tracker;
|
||||
private final LogFile txnFile;
|
||||
// We need an explicit lock because the transaction tidier cannot store a reference to the transaction
|
||||
private final Object lock;
|
||||
private final Ref<LogTransaction> selfRef;
|
||||
// Deleting sstables is tricky because the mmapping might not have been finalized yet,
|
||||
// and delete will fail (on Windows) until it is (we only force the unmapping on SUN VMs).
|
||||
|
|
@ -125,7 +128,8 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
{
|
||||
this.tracker = tracker;
|
||||
this.txnFile = new LogFile(opType, UUIDGen.getTimeUUID());
|
||||
this.selfRef = new Ref<>(this, new TransactionTidier(txnFile));
|
||||
this.lock = new Object();
|
||||
this.selfRef = new Ref<>(this, new TransactionTidier(txnFile, lock));
|
||||
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Created transaction logs with id {}", txnFile.id());
|
||||
|
|
@ -136,7 +140,13 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
**/
|
||||
void trackNew(SSTable table)
|
||||
{
|
||||
txnFile.add(table);
|
||||
synchronized (lock)
|
||||
{
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Track NEW sstable {} in {}", table.getFilename(), txnFile.toString());
|
||||
|
||||
txnFile.add(table);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -161,25 +171,34 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
*/
|
||||
SSTableTidier obsoleted(SSTableReader reader, LogRecord logRecord)
|
||||
{
|
||||
if (txnFile.contains(Type.ADD, reader, logRecord))
|
||||
synchronized (lock)
|
||||
{
|
||||
if (txnFile.contains(Type.REMOVE, reader, logRecord))
|
||||
throw new IllegalArgumentException();
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Track OLD sstable {} in {}", reader.getFilename(), txnFile.toString());
|
||||
|
||||
return new SSTableTidier(reader, true, this);
|
||||
if (txnFile.contains(Type.ADD, reader, logRecord))
|
||||
{
|
||||
if (txnFile.contains(Type.REMOVE, reader, logRecord))
|
||||
throw new IllegalArgumentException();
|
||||
|
||||
return new SSTableTidier(reader, true, this);
|
||||
}
|
||||
|
||||
txnFile.addRecord(logRecord);
|
||||
|
||||
if (tracker != null)
|
||||
tracker.notifyDeleting(reader);
|
||||
|
||||
return new SSTableTidier(reader, false, this);
|
||||
}
|
||||
|
||||
txnFile.addRecord(logRecord);
|
||||
|
||||
if (tracker != null)
|
||||
tracker.notifyDeleting(reader);
|
||||
|
||||
return new SSTableTidier(reader, false, this);
|
||||
}
|
||||
|
||||
Map<SSTable, LogRecord> makeRemoveRecords(Iterable<SSTableReader> sstables)
|
||||
{
|
||||
return txnFile.makeRecords(Type.REMOVE, sstables);
|
||||
synchronized (lock)
|
||||
{
|
||||
return txnFile.makeRecords(Type.REMOVE, sstables);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
|
@ -249,13 +268,15 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
private static class TransactionTidier implements RefCounted.Tidy, Runnable
|
||||
{
|
||||
private final LogFile data;
|
||||
private final Object lock;
|
||||
|
||||
TransactionTidier(LogFile data)
|
||||
TransactionTidier(LogFile data, Object lock)
|
||||
{
|
||||
this.data = data;
|
||||
this.lock = lock;
|
||||
}
|
||||
|
||||
public void tidy() throws Exception
|
||||
public void tidy()
|
||||
{
|
||||
run();
|
||||
}
|
||||
|
|
@ -267,34 +288,35 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
|
||||
public void run()
|
||||
{
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Removing files for transaction log {}", data);
|
||||
|
||||
// this happens if we forget to close a txn and the garbage collector closes it for us
|
||||
// or if the transaction journal was never properly created in the first place
|
||||
if (!data.completed())
|
||||
{
|
||||
logger.error("Transaction log {} indicates txn was not completed, trying to abort it now", data);
|
||||
Throwable err = Throwables.perform((Throwable)null, data::abort);
|
||||
if (err != null)
|
||||
logger.error("Failed to abort transaction log {}", data, err);
|
||||
}
|
||||
|
||||
Throwable err = data.removeUnfinishedLeftovers(null);
|
||||
|
||||
if (err != null)
|
||||
{
|
||||
logger.info("Failed deleting files for transaction log {}, we'll retry after GC and on on server restart",
|
||||
data,
|
||||
err);
|
||||
failedDeletions.add(this);
|
||||
}
|
||||
else
|
||||
synchronized (lock)
|
||||
{
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Closing transaction log {}", data);
|
||||
logger.trace("Removing files for transaction {}", name());
|
||||
|
||||
data.close();
|
||||
// this happens if we forget to close a txn and the garbage collector closes it for us
|
||||
// or if the transaction journal was never properly created in the first place
|
||||
if (!data.completed())
|
||||
{
|
||||
logger.error("{} was not completed, trying to abort it now", data);
|
||||
Throwable err = Throwables.perform((Throwable) null, data::abort);
|
||||
if (err != null)
|
||||
logger.error("Failed to abort {}", data, err);
|
||||
}
|
||||
|
||||
Throwable err = data.removeUnfinishedLeftovers(null);
|
||||
|
||||
if (err != null)
|
||||
{
|
||||
logger.info("Failed deleting files for transaction {}, we'll retry after GC and on on server restart", name(), err);
|
||||
failedDeletions.add(this);
|
||||
}
|
||||
else
|
||||
{
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Closing file transaction {}", name());
|
||||
|
||||
data.close();
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -324,7 +346,9 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
private final long sizeOnDisk;
|
||||
private final Tracker tracker;
|
||||
private final boolean wasNew;
|
||||
private final Object lock;
|
||||
private final Ref<LogTransaction> parentRef;
|
||||
private final UUID txnId;
|
||||
|
||||
public SSTableTidier(SSTableReader referent, boolean wasNew, LogTransaction parent)
|
||||
{
|
||||
|
|
@ -332,7 +356,12 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
this.sizeOnDisk = referent.bytesOnDisk();
|
||||
this.tracker = parent.tracker;
|
||||
this.wasNew = wasNew;
|
||||
this.lock = parent.lock;
|
||||
this.parentRef = parent.selfRef.tryRef();
|
||||
this.txnId = parent.id();
|
||||
|
||||
if (this.parentRef == null)
|
||||
throw new IllegalStateException("Transaction already completed");
|
||||
}
|
||||
|
||||
public void run()
|
||||
|
|
@ -340,36 +369,45 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
if (tracker != null && !tracker.isDummy())
|
||||
SystemKeyspace.clearSSTableReadMeter(desc.ksname, desc.cfname, desc.generation);
|
||||
|
||||
try
|
||||
synchronized (lock)
|
||||
{
|
||||
// If we can't successfully delete the DATA component, set the task to be retried later: see TransactionTidier
|
||||
File datafile = new File(desc.filenameFor(Component.DATA));
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Tidier running for old sstable {}", desc.baseFilename());
|
||||
|
||||
if (datafile.exists())
|
||||
delete(datafile);
|
||||
else if (!wasNew)
|
||||
logger.error("SSTableTidier ran with no existing data file for an sstable that was not new");
|
||||
try
|
||||
{
|
||||
// If we can't successfully delete the DATA component, set the task to be retried later: see TransactionTidier
|
||||
File datafile = new File(desc.filenameFor(Component.DATA));
|
||||
|
||||
// let the remainder be cleaned up by delete
|
||||
SSTable.delete(desc, SSTable.discoverComponentsFor(desc));
|
||||
if (datafile.exists())
|
||||
delete(datafile);
|
||||
else if (!wasNew)
|
||||
logger.error("SSTableTidier ran with no existing data file for an sstable that was not new");
|
||||
|
||||
// let the remainder be cleaned up by delete
|
||||
SSTable.delete(desc, SSTable.discoverComponentsFor(desc));
|
||||
}
|
||||
catch (Throwable t)
|
||||
{
|
||||
logger.error("Failed deletion for {}, we'll retry after GC and on server restart", desc);
|
||||
failedDeletions.add(this);
|
||||
return;
|
||||
}
|
||||
|
||||
if (tracker != null && tracker.cfstore != null && !wasNew)
|
||||
tracker.cfstore.metric.totalDiskSpaceUsed.dec(sizeOnDisk);
|
||||
|
||||
// release the referent to the parent so that the all transaction files can be released
|
||||
parentRef.release();
|
||||
}
|
||||
catch (Throwable t)
|
||||
{
|
||||
logger.error("Failed deletion for {}, we'll retry after GC and on server restart", desc);
|
||||
failedDeletions.add(this);
|
||||
return;
|
||||
}
|
||||
|
||||
if (tracker != null && tracker.cfstore != null && !wasNew)
|
||||
tracker.cfstore.metric.totalDiskSpaceUsed.dec(sizeOnDisk);
|
||||
|
||||
// release the referent to the parent so that the all transaction files can be released
|
||||
parentRef.release();
|
||||
}
|
||||
|
||||
public void abort()
|
||||
{
|
||||
parentRef.release();
|
||||
synchronized (lock)
|
||||
{
|
||||
parentRef.release();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -392,6 +430,10 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
@VisibleForTesting
|
||||
Throwable complete(Throwable accumulate)
|
||||
{
|
||||
if (logger.isTraceEnabled())
|
||||
logger.trace("Completing txn {} with last record {}",
|
||||
txnFile.toString(), txnFile.getLastRecord());
|
||||
|
||||
try
|
||||
{
|
||||
accumulate = selfRef.ensureReleased(accumulate);
|
||||
|
|
@ -406,12 +448,18 @@ class LogTransaction extends Transactional.AbstractTransactional implements Tran
|
|||
|
||||
protected Throwable doCommit(Throwable accumulate)
|
||||
{
|
||||
return complete(Throwables.perform(accumulate, txnFile::commit));
|
||||
synchronized (lock)
|
||||
{
|
||||
return complete(Throwables.perform(accumulate, txnFile::commit));
|
||||
}
|
||||
}
|
||||
|
||||
protected Throwable doAbort(Throwable accumulate)
|
||||
{
|
||||
return complete(Throwables.perform(accumulate, txnFile::abort));
|
||||
synchronized (lock)
|
||||
{
|
||||
return complete(Throwables.perform(accumulate, txnFile::abort));
|
||||
}
|
||||
}
|
||||
|
||||
protected void doPrepare() { }
|
||||
|
|
|
|||
Loading…
Reference in New Issue