mirror of https://github.com/apache/cassandra
Fix for CASSANDRA-9 JIRA.
git-svn-id: https://svn.apache.org/repos/asf/incubator/cassandra/trunk@758044 13f79535-47bb-0310-9956-ffa450edef68
This commit is contained in:
parent
5506a44d0a
commit
74dcfb263c
|
|
@ -18,6 +18,7 @@
|
|||
|
||||
package org.apache.cassandra.db;
|
||||
|
||||
import java.io.FileOutputStream;
|
||||
import java.io.IOException;
|
||||
import java.util.*;
|
||||
import java.util.concurrent.Callable;
|
||||
|
|
@ -53,6 +54,7 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
private static Logger logger_ = Logger.getLogger( Memtable.class );
|
||||
private static Map<String, ExecutorService> apartments_ = new HashMap<String, ExecutorService>();
|
||||
public static final String flushKey_ = "FlushKey";
|
||||
|
||||
public static void shutdown()
|
||||
{
|
||||
Set<String> names = apartments_.keySet();
|
||||
|
|
@ -156,6 +158,28 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
columnFamilies_.put(key_, columnFamily_);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Flushes the current memtable to disk.
|
||||
*
|
||||
* @author alakshman
|
||||
*
|
||||
*/
|
||||
class Flusher implements Runnable
|
||||
{
|
||||
private CommitLog.CommitLogContext cLogCtx_;
|
||||
|
||||
Flusher(CommitLog.CommitLogContext cLogCtx)
|
||||
{
|
||||
cLogCtx_ = cLogCtx;
|
||||
}
|
||||
|
||||
public void run()
|
||||
{
|
||||
ColumnFamilyStore cfStore = Table.open(table_).getColumnFamilyStore(cfName_);
|
||||
MemtableManager.instance().submit(cfName_, Memtable.this, cLogCtx_);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Compares two Memtable based on creation time.
|
||||
|
|
@ -235,7 +259,10 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
if (!isFrozen_)
|
||||
{
|
||||
isFrozen_ = true;
|
||||
MemtableManager.instance().submit(cfStore.getColumnFamilyName(), this, cLogCtx);
|
||||
/* Submit this Memtable to be flushed. */
|
||||
Runnable flusher = new Flusher(cLogCtx);
|
||||
apartments_.get(cfName_).submit(flusher);
|
||||
// MemtableManager.instance().submit(cfStore.getColumnFamilyName(), this, cLogCtx);
|
||||
cfStore.switchMemtable(key, columnFamily, cLogCtx);
|
||||
}
|
||||
else
|
||||
|
|
@ -280,8 +307,6 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
}
|
||||
}
|
||||
|
||||
|
||||
|
||||
private void resolve(String key, ColumnFamily columnFamily)
|
||||
{
|
||||
ColumnFamily oldCf = columnFamilies_.get(key);
|
||||
|
|
@ -314,7 +339,6 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
resolve(key, columnFamily);
|
||||
}
|
||||
|
||||
|
||||
ColumnFamily getLocalCopy(String key, String cfName, IFilter filter)
|
||||
{
|
||||
String[] values = RowMutation.getColumnAndColumnFamily(cfName);
|
||||
|
|
@ -458,7 +482,7 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
if ( columnFamily != null )
|
||||
{
|
||||
/* serialize the cf with column indexes */
|
||||
ColumnFamily.serializer2().serialize( columnFamily, buffer );
|
||||
ColumnFamily.serializer2().serialize( columnFamily, buffer );
|
||||
/* Now write the key and value to disk */
|
||||
ssTable.append(pKey.key(), pKey.hash(), buffer);
|
||||
bf.fill(pKey.key());
|
||||
|
|
@ -485,7 +509,7 @@ public class Memtable implements MemtableMBean, Comparable<Memtable>
|
|||
if ( columnFamily != null )
|
||||
{
|
||||
/* serialize the cf with column indexes */
|
||||
ColumnFamily.serializer2().serialize( columnFamily, buffer );
|
||||
ColumnFamily.serializer2().serialize( columnFamily, buffer );
|
||||
/* Now write the key and value to disk */
|
||||
ssTable.append(key, buffer);
|
||||
bf.fill(key);
|
||||
|
|
|
|||
Loading…
Reference in New Issue