Add proper error handling to stream receiver

patch by Paulo Motta; reviewed by yukim for CASSANDRA-10774
This commit is contained in:
Paulo Motta 2015-11-27 16:37:37 -08:00 committed by Yuki Morishita
parent 7650fc1963
commit 5ba69a3259
2 changed files with 61 additions and 49 deletions

View File

@ -1,4 +1,5 @@
2.1.12
* Add proper error handling to stream receiver (CASSANDRA-10774)
* Warn or fail when changing cluster topology live (CASSANDRA-10243)
* Status command in debian/ubuntu init script doesn't work (CASSANDRA-10213)
* Some DROP ... IF EXISTS incorrectly result in exceptions on non-existing KS (CASSANDRA-10658)

View File

@ -40,6 +40,7 @@ import org.apache.cassandra.dht.Token;
import org.apache.cassandra.io.sstable.SSTableReader;
import org.apache.cassandra.io.sstable.SSTableWriter;
import org.apache.cassandra.utils.FBUtilities;
import org.apache.cassandra.utils.JVMStabilityInspector;
import org.apache.cassandra.utils.Pair;
import org.apache.cassandra.utils.concurrent.Refs;
@ -116,63 +117,73 @@ public class StreamReceiveTask extends StreamTask
public void run()
{
Pair<String, String> kscf = Schema.instance.getCF(task.cfId);
if (kscf == null)
try
{
// schema was dropped during streaming
for (SSTableWriter writer : task.sstables)
writer.abort();
task.sstables.clear();
return;
}
ColumnFamilyStore cfs = Keyspace.open(kscf.left).getColumnFamilyStore(kscf.right);
File lockfiledir = cfs.directories.getWriteableLocationAsFile(task.sstables.size() * 256L);
if (lockfiledir == null)
throw new IOError(new IOException("All disks full"));
StreamLockfile lockfile = new StreamLockfile(lockfiledir, UUID.randomUUID());
lockfile.create(task.sstables);
List<SSTableReader> readers = new ArrayList<>();
for (SSTableWriter writer : task.sstables)
readers.add(writer.closeAndOpenReader());
lockfile.delete();
task.sstables.clear();
try (Refs<SSTableReader> refs = Refs.ref(readers))
{
// add sstables and build secondary indexes
cfs.addSSTables(readers);
cfs.indexManager.maybeBuildSecondaryIndexes(readers, cfs.indexManager.allIndexesNames());
//invalidate row and counter cache
if (cfs.isRowCacheEnabled() || cfs.metadata.isCounter())
Pair<String, String> kscf = Schema.instance.getCF(task.cfId);
if (kscf == null)
{
List<Bounds<Token>> boundsToInvalidate = new ArrayList<>(readers.size());
for (SSTableReader sstable : readers)
boundsToInvalidate.add(new Bounds<Token>(sstable.first.getToken(), sstable.last.getToken()));
Set<Bounds<Token>> nonOverlappingBounds = Bounds.getNonOverlappingBounds(boundsToInvalidate);
// schema was dropped during streaming
for (SSTableWriter writer : task.sstables)
writer.abort();
task.sstables.clear();
task.session.taskCompleted(task);
return;
}
ColumnFamilyStore cfs = Keyspace.open(kscf.left).getColumnFamilyStore(kscf.right);
if (cfs.isRowCacheEnabled())
{
int invalidatedKeys = cfs.invalidateRowCache(nonOverlappingBounds);
if (invalidatedKeys > 0)
logger.debug("[Stream #{}] Invalidated {} row cache entries on table {}.{} after stream " +
"receive task completed.", task.session.planId(), invalidatedKeys,
cfs.keyspace.getName(), cfs.getColumnFamilyName());
}
File lockfiledir = cfs.directories.getWriteableLocationAsFile(task.sstables.size() * 256L);
if (lockfiledir == null)
throw new IOError(new IOException("All disks full"));
StreamLockfile lockfile = new StreamLockfile(lockfiledir, UUID.randomUUID());
lockfile.create(task.sstables);
List<SSTableReader> readers = new ArrayList<>();
for (SSTableWriter writer : task.sstables)
readers.add(writer.closeAndOpenReader());
lockfile.delete();
task.sstables.clear();
if (cfs.metadata.isCounter())
try (Refs<SSTableReader> refs = Refs.ref(readers))
{
// add sstables and build secondary indexes
cfs.addSSTables(readers);
cfs.indexManager.maybeBuildSecondaryIndexes(readers, cfs.indexManager.allIndexesNames());
//invalidate row and counter cache
if (cfs.isRowCacheEnabled() || cfs.metadata.isCounter())
{
int invalidatedKeys = cfs.invalidateCounterCache(nonOverlappingBounds);
if (invalidatedKeys > 0)
logger.debug("[Stream #{}] Invalidated {} counter cache entries on table {}.{} after stream " +
"receive task completed.", task.session.planId(), invalidatedKeys,
cfs.keyspace.getName(), cfs.getColumnFamilyName());
List<Bounds<Token>> boundsToInvalidate = new ArrayList<>(readers.size());
for (SSTableReader sstable : readers)
boundsToInvalidate.add(new Bounds<Token>(sstable.first.getToken(), sstable.last.getToken()));
Set<Bounds<Token>> nonOverlappingBounds = Bounds.getNonOverlappingBounds(boundsToInvalidate);
if (cfs.isRowCacheEnabled())
{
int invalidatedKeys = cfs.invalidateRowCache(nonOverlappingBounds);
if (invalidatedKeys > 0)
logger.debug("[Stream #{}] Invalidated {} row cache entries on table {}.{} after stream " +
"receive task completed.", task.session.planId(), invalidatedKeys,
cfs.keyspace.getName(), cfs.getColumnFamilyName());
}
if (cfs.metadata.isCounter())
{
int invalidatedKeys = cfs.invalidateCounterCache(nonOverlappingBounds);
if (invalidatedKeys > 0)
logger.debug("[Stream #{}] Invalidated {} counter cache entries on table {}.{} after stream " +
"receive task completed.", task.session.planId(), invalidatedKeys,
cfs.keyspace.getName(), cfs.getColumnFamilyName());
}
}
}
}
task.session.taskCompleted(task);
task.session.taskCompleted(task);
}
catch (Throwable t)
{
logger.error("Error applying streamed data: ", t);
JVMStabilityInspector.inspectThrowable(t);
task.session.onError(t);
}
}
}