diff --git a/conf/storage-conf.xml b/conf/storage-conf.xml index 562200784a..1ddaea46c5 100644 --- a/conf/storage-conf.xml +++ b/conf/storage-conf.xml @@ -190,29 +190,41 @@ - - 256 + + 32 + 8 + + + 64 - 32 - + 64 - 0.01 + 0.1 + Increase ConcurrentWrites to the number of clients writing + at once if you enable CommitLogSync + CommitLogSyncDelay. --> 8 32 diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index bdc6219d95..68aa4a14f4 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -67,6 +67,9 @@ public class DatabaseDescriptor private static int consistencyThreads_ = 4; // not configurable private static int concurrentReaders_ = 8; private static int concurrentWriters_ = 32; + + private static int flushDataBufferSizeInMB_ = 32; + private static int flushIndexBufferSizeInMB_ = 32; private static List tables_ = new ArrayList(); private static Set applicationColumnFamilies_ = new HashSet(); @@ -224,6 +227,17 @@ public class DatabaseDescriptor concurrentWriters_ = Integer.parseInt(rawWriters); } + String rawFlushData = xmlUtils.getNodeValue("/Storage/FlushDataBufferSizeInMB"); + if (rawFlushData != null) + { + flushDataBufferSizeInMB_ = Integer.parseInt(rawFlushData); + } + String rawFlushIndex = xmlUtils.getNodeValue("/Storage/FlushIndexBufferSizeInMB"); + if (rawFlushIndex != null) + { + flushIndexBufferSizeInMB_ = Integer.parseInt(rawFlushIndex); + } + /* TCP port on which the storage system listens */ String port = xmlUtils.getNodeValue("/Storage/StoragePort"); if ( port != null ) @@ -909,4 +923,14 @@ public class DatabaseDescriptor { return commitLogSync_; } + + public static int getFlushDataBufferSizeInMB() + { + return flushDataBufferSizeInMB_; + } + + public static int getFlushIndexBufferSizeInMB() + { + return flushIndexBufferSizeInMB_; + } } diff --git a/src/java/org/apache/cassandra/io/SSTableWriter.java b/src/java/org/apache/cassandra/io/SSTableWriter.java index 2cc902a4b4..9098c11718 100644 --- a/src/java/org/apache/cassandra/io/SSTableWriter.java +++ b/src/java/org/apache/cassandra/io/SSTableWriter.java @@ -11,6 +11,7 @@ import org.apache.log4j.Logger; import org.apache.cassandra.dht.IPartitioner; import org.apache.cassandra.utils.BloomFilter; +import org.apache.cassandra.config.DatabaseDescriptor; import com.reardencommerce.kernel.collections.shared.evictable.ConcurrentLinkedHashMap; public class SSTableWriter extends SSTable @@ -26,8 +27,8 @@ public class SSTableWriter extends SSTable public SSTableWriter(String filename, int keyCount, IPartitioner partitioner) throws IOException { super(filename, partitioner); - dataFile = new BufferedRandomAccessFile(path, "rw", 4 * 1024 * 1024); - indexFile = new BufferedRandomAccessFile(indexFilename(), "rw", 1024 * 1024); + dataFile = new BufferedRandomAccessFile(path, "rw", DatabaseDescriptor.getFlushDataBufferSizeInMB() * 1024 * 1024); + indexFile = new BufferedRandomAccessFile(indexFilename(), "rw", DatabaseDescriptor.getFlushIndexBufferSizeInMB() * 1024 * 1024); bf = new BloomFilter(keyCount, 15); } diff --git a/test/system/stress.py b/test/system/stress.py index 86768c939b..16cc68ed93 100644 --- a/test/system/stress.py +++ b/test/system/stress.py @@ -30,7 +30,7 @@ class Inserter(Thread): self.count = 0 client = get_client(port=9160) client.transport.open() - for i in xrange(0, 1000): + for i in xrange(0, 200): data = md5(str(i)).hexdigest() for j in xrange(0, 1000): key = '%s.%s.%s' % (time.time(), id, j)