From a287f42ceb1d5eac033e2db2201422801b772e99 Mon Sep 17 00:00:00 2001 From: Jonathan Ellis Date: Wed, 6 Apr 2011 16:21:17 +0000 Subject: [PATCH] add a server-wide cap on memtable memory usage git-svn-id: https://svn.apache.org/repos/asf/cassandra/trunk@1089521 13f79535-47bb-0310-9956-ffa450edef68 --- build.xml | 1 + conf/cassandra-env.sh | 3 + conf/cassandra.yaml | 7 ++ lib/jamm-0.2.jar | Bin 0 -> 5935 bytes .../org/apache/cassandra/config/Config.java | 3 +- .../cassandra/config/DatabaseDescriptor.java | 17 +++ .../apache/cassandra/db/BinaryMemtable.java | 2 +- .../cassandra/db/ColumnFamilyStore.java | 68 ++++++++++-- .../org/apache/cassandra/db/Memtable.java | 97 ++++++++++++++-- .../apache/cassandra/db/MeteredFlusher.java | 104 ++++++++++++++++++ .../cassandra/service/StorageService.java | 2 +- .../apache/cassandra/streaming/StreamIn.java | 2 +- test/conf/cassandra.yaml | 1 + .../cassandra/db/MeteredFlusherTest.java | 51 +++++++++ .../org/apache/cassandra/db/DefsTest.java | 2 +- 15 files changed, 333 insertions(+), 27 deletions(-) create mode 100644 lib/jamm-0.2.jar create mode 100644 src/java/org/apache/cassandra/db/MeteredFlusher.java create mode 100644 test/long/org/apache/cassandra/db/MeteredFlusherTest.java diff --git a/build.xml b/build.xml index f801791383..9fe151b970 100644 --- a/build.xml +++ b/build.xml @@ -615,6 +615,7 @@ + diff --git a/conf/cassandra-env.sh b/conf/cassandra-env.sh index f9930e484d..68641a9c54 100644 --- a/conf/cassandra-env.sh +++ b/conf/cassandra-env.sh @@ -91,6 +91,9 @@ JMX_PORT="7199" # performance benefit (around 5%). JVM_OPTS="$JVM_OPTS -ea" +# add the jamm javaagent +JVM_OPTS="$JVM_OPTS -javaagent:$CASSANDRA_HOME/lib/jamm-0.2.jar" + # enable thread priorities, primarily so we can give periodic tasks # a lower priority to avoid interfering with client workload JVM_OPTS="$JVM_OPTS -XX:+UseThreadPriorities" diff --git a/conf/cassandra.yaml b/conf/cassandra.yaml index 4b02273bae..f97580cd31 100644 --- a/conf/cassandra.yaml +++ b/conf/cassandra.yaml @@ -141,6 +141,13 @@ reduce_cache_capacity_to: 0.6 concurrent_reads: 32 concurrent_writes: 32 +# Total memory to use for memtables. Cassandra will flush the largest +# memtable when this much memory is used. Prefer using this to +# the older, per-ColumnFamily memtable flush thresholds. +# If omitted, Cassandra will set it to 1/3 of the heap. +# If set to 0, only the old flush thresholds are used. +# memtable_total_space_in_mb: 2048 + # This sets the amount of memtable flush writer threads. These will # be blocked by disk io, and each one will hold a memtable in memory # while blocked. If you have a large heap and many data directories, diff --git a/lib/jamm-0.2.jar b/lib/jamm-0.2.jar new file mode 100644 index 0000000000000000000000000000000000000000..8af087ec9d4114c986210be7b7f3fe7d4550a56b GIT binary patch literal 5935 zcmai2by!r}`X0KQkw#izKTj##2yz204{9Zaz$+i4&UyJvl zR7K@Et{;X%`-WyF0$d!2)d7)Rue~8($G8RB1v7^4X2r%9TN9%)i`+!3`Q#-Zv<~1X zD1C`*g~yt=t*nIkvG-$mz~hA=Y%0%;hVVtofb*Pj?xhi=LS)4qjvMOfAHTCIDaDS9lad9t^NiChyGKF@Si{@ zOBa{F(3pT9=xBvc!DC1Oz#|j@faN!6DQ90xe~6rmyO+PVtrrvbeQRe+2qZ#l(REsq zP;61i!-*p&jre|9nIIja7~PCYx4NXcnsk+3b474$cz$p`es)j#JGL3%3*4s0li-S512=sj0eDCTk@UP@3-eeN6Ft17O*j^@FEOZ2476ivP2U*V)#Fx zJe_0}QIj|xspY3!97$*htS7)y4QG>S=~|&Z#i2#$l@kdM0h`6&S~-{LWha_K5eX`! zZ5_p>EuXLC+XDS$i(0x9(rTCjR;Xe=c;tU`dj(BwiGf&$ztBt__6X*`Gx-u#>fsj4 zoqJOkA#SbDP$G+aT6FSuRjN1Y8c>!#HIl(83L=j+@lTLv2!T2jt*%srGPYf=(0D4G zijwDa-&qD5{2}+DrysU)=EpKWS-?8?e~A%y14rosXfLH ziCut?)LpsX4vE5<)E(nOar=Xwhko2Ao!ZKF^P(a=AOU&v!sXY^^}CcUEU=R=;u0GO zwe5BUl;%12hROh>*BctA8PJzh_vK}46a;+jTG3*Sq1v^cf{g76{sG zrns{k(xAT#>9x<#KNWXt83ulO81nvk3x4mj_EG1@UnZr<;BZmTcXkjL)8l63+{iR8 zmo#Y31m#%408CBtn{AW77`}T)n|@bYlbi9{!WRoYGQr1$SCM1SX}I$CVjF|TcTqzR zA`tAprac+GZ3B|(LwwUQv8VH`XQ?#N-Tg%(9*dJ-usj$`lSFA1gNDhUvrA_^Xtf)Z z>rivR0H0UY&x`vQPpNTNXw#V5zMrT^MUo)yS5}{PBA1U57EAMhX6vY5vdui><&!p8 zx19dyG(=n4atd=l3L=ynSpv4^ycBvx=A1N1fr@S7x=j%wsmyxy0Hu zn<*M)S@`UEL_jz>@IiT`X$ri628rYvE#=OKl5|X!SgN>9l8i@?g_&259jCiFI}3QZ zQ}vR?Yre=M%kF)_c{|-9AO@KM50xDuD1k%aa)>EeJ=YVdl~b^2W#_A<) zIPS)C;(at#9wDg}G;HKDTQ*0Q962eWXl=f+OT$fv_s8@=4TeFh4Wt*{lK%L`P{~ob zS+*eA6>K$ad7`@`aS@f2Nsl8&*E0w4?efzM?H#-`O00MHQaV#U3$M`(kOc~XWb%MR z*|gxiyM4p>R-AFs-bk^;k5c8C)90+h`Kbb#GqOzk?ro^cl%X*4@4BX4Qn7?QOj!F2 zXrGkPp<<|395fgtp5Go->`gzU7`(yqCiL;I{YvIOlCtyZ#bYt69Q!M(1-e}XE%Ui7 z7EPh09Q(;z!OW}SjaaI7l_pTBy~)X%LVta{#*I_S71fFf{2E*ZiMKP4Yiz3d?r-1S zr+;MQmLO@`KgT#<4-e=KiMQJ%_jdhOFCv&C0>2F6mCrv~mXQz~sBHR7yQ`CO?^M3o zd^?zSo#&#UrU&{QJ$2PTC|A-niC}31M7mv^djlNCIq&v{tGkSrC@!(Y=wM3K~O!h_uGfb49UztL}$MH(r=Y9F_QKYydNlJF*sM2a+yVMT;)UViMK?8 z+m@pht-n(V8Wlcz8(z!r-M5Y}IF*0TVhrD=7XV#VVY8&2b>eYfinU(V*yKwskt$aS`7;xkoVacpyp)@_}zd|D0ht^FREGkvxw!k;XZ z{+zQvdqms=>-_`?uwNv(T76sc6Bvx$fLmz!h~0gb(b4gN`KiqWoQ#$osagoT0qPCS zkKkmU`8>sHyS79=Tfx>HIM(o5GRhJ%JU`)~B63(vu9xQbkVb?*K(>W#J3)AfU6bI6;LH73r{K3 zZaxEcF`R8;&yNoVXFX#b*C1sfrE2yJcZOdIeNj6>OI~K)8;=pUbSyYD2deIPlJs)M zK&}N26N5!rixo0ndeN>Q-TPK=p{8>ZvuLDVEC_?Ogn^Kol zm{^!>U@mCkNrcaIWO4!M$dh!6ti=uuC8_Uhc-!aWd7CU88ae%3id%GyvCNEO;suS6 z6yLBq)MXkOyVSg1S(GT2Xbxt*6tic^RgF?CXf+wj6FyrSyDm5qUFiy5>u&;Nth>=~ zz(arp;cKvOT;}-U9{iWFhzGH2Aa+`vlNl3gT*mNv#2P)g7QSb>Vf~o4EA*ibwbo=k z+?_nakB?mJTiu3$07AaR44^$PS( zRLEn)KQ#HN$$D|b9l--F!HMfuKGLIV&=X{}==maYs7{`qhkW>PNfN~!m@>PWo<>8L z|8&Ns@tpJWtMae2t>rFzKvq@kQz|XR#mUq$J+Wkg8`)AFsW%%rVf|FT+oBh{;RSK} zu|Uw0HB`KXTXsLl(u9^8Q`>o%%c7EsIdJ|73#ndwZkHA}U1l?{iBp>&IsI@5VV zxKzGIE`|E8+5!)0ZLUdx?d+BS%rG0M_hf4IVW5;~Xms;Lr6xwaF;M1uZGq#lV z0zTt0y`$!UT9GTU4{{l$Xyb}M+YA3^rg zyfgG;`9=`O*|IN-b*{|&oC~wWH1PPc2y7d{eaENi{p<<_Z@_C@T%SP;USvgH-Pf(- z+#kiFCcy1NXDPbWbKFsy>;rd@SljuoV-1K`^3tB&BS&v*bf+``wFkXQ$3ocb*dPU) zWK`W4u;+VGcD|Fn2eX}@5*4LVWtQ-!Ho2%y)GKk!Kz%xmdphYu!;v}X6NA{ z!m#?RxM8hL5uR2I)Z#B?ov|4Vl`@*=`>eaSPiTV16Ekk_z{BNeW^)Riz+Ox zy+*%2Fdj6eQ1+AK|Cpg<%>8VV6iY8Dy;~pC6t7T{GxmVVziXQvO{0jPITgyx#_f!! z`hgB=hv#ijwPy4llqD7o%t(z*plI76&(2VHFL! zFtP?;ri4aaRhW8sneRyfl$NV@u(G>cEzCkBcy7-}aIY$+(CLaIPgLW1O25 zw$e5ID@j;=ME=UE7_fWhSoFpWcK7%p-s9dy4h(HFJhaIwB`8#4ofGdd(ek$#U%%70 zh3|>l<#|neY7<$>*AL0+AxZqMBAhsm2~-Jtvj<6=Y5`3zm>L}|b_$26+{4Z}-{IRS zx1}#2%(mkO#urYA+URp%NV+W?+7${nX0CFPRAi7TuobD6<1Cx19A8_dT9R z{AjZjjo}WoOzLTQ+Ff2RI_tIMyewbssRpLnzWB`Ld#bt~0I8ndxUibHCU))X9n^no zyHMhp#f@yV%;Y{~FXm2Av@@lDpNjXfNvYR;3Lngr1KOvpUiFs!A6?SDC|ieCdHq@U z%Uw*~t}yDKRY>70dEhx2eAPRC7n#V&vCF&_QL^t5oNi+N&R0*es7Pl+XzuiC#F3Md zwK2i1=IR_PF#TRJbte0sFce`$gq#&EKc-|tZxXqOBup~)AvDTqFT46_AunCG6?xYs z{S}G~W}KiQIv!4`;yd75ybr9l293hlGHqvUA>T$T$)(lg5NxQR&Ac#nmCt2W(r?ky zAST1_bXJLVlu|YI-|@29-6yqTkvuMDUN;Cb&JV3hZ&JM~C6lD$jI}t6OkOlIJVDku z-1r=KUBuS_#>u*g8wnb>Ri4FfqH)=*#WMyE>P;l`z@tLKLdKIHjA$;Qy0ewg1TX^{thmDG>avt&?~sIp+;OyBMX|c#w-oOM(DZT0 zql+}6m^Kg&h-dl5Wv*(Lx+A|9iV!{Kah6WaaByP@g7`bRy{(8o)a3ER;jC_kU>tmu zwZn9U#?FE7Gx$Npy?Jz>qMTDcJ8Wss4Z)LNgg0(0dE1u4)0{H>iKdKny$T;HM0|Ej ztNuVFl=zw5)!c8UyXea{eR6So`W;=TgQ0Ri^4P0aH+**l)q5P*new3big{!6c*Cgu zN^22Pwp>9kaZ1p#BVK2zXJxuy>$a2W!^JlzzmC>pWc8YF2waj(poD*y+V7Rxe=9?L zH52XwjGe_C4CJ9ZI*}}p+crtcoQT0WO|Mn5mhvg<HSuJEEQBOqyJZV6Rb#iNcxMb|XTB^eSnCOMz0?051w#;jFPYwUa_AebT=Rra)U zdq9A@4jgYANl5+Xt-}48YK9IzFF6+nK+_`_!vTY3vhVzv!ZuzQHCBF)uK4*g2SX+% z%v4a=z)M}*RTv6G6rv?Glt6ocV%K~0$-Byo@dMPJGAk6&gJaA!(LUAlCgi&;wTY)z z{$VK(S+Dt85_lj5rLzFbfEfQ14mQ7pJ{m^huS^J`fz*_gC7zMq(FoFk#f{obYtW-* zno=fN0-90wTw!g6FEpDP_*mO8wSr?h(2Sav38VVn8g70hbAL=87ICnr&%Nji?JY~E zqp_0d^cefnPOu1wNF=c+X52(FIcA)F*q=kM(LOcirK`qC$Jrob_xHY=$8`5LW-iCi zB#dp2h(C)KmT1lDW9N*|Zd^SI>+gZ*F}Cv%v*O>c*JZ3`lBZ#o_m0*SvkGI{Tw2do zVNKo2n2(@03#w2Yp)i9zdpigL z2U0)2e>swWn3q4Pzqc=c5Ppo(P|)rBFY4dS%)c-GJ1g;D^nWeBcKiE>vG{Y%|1=+e zSc`v*AK{|hc6~n>e;JLxh`+NMe-SzUN&M4z{K@{k_4sq)(omz{uzznv{@VQS&8uIU z&bQ6$-?s9vi2kYhM;r2Y#osl*zZ63-{!{Tct?<7K{(EEeORx#+e+vGL4-NeX1?%>4 NMg{=tewf^V{{!p!J9q#9 literal 0 HcmV?d00001 diff --git a/src/java/org/apache/cassandra/config/Config.java b/src/java/org/apache/cassandra/config/Config.java index f80ea8d838..504ec9e17e 100644 --- a/src/java/org/apache/cassandra/config/Config.java +++ b/src/java/org/apache/cassandra/config/Config.java @@ -57,7 +57,8 @@ public class Config public Integer concurrent_replicates = 32; public Integer memtable_flush_writers = null; // will get set to the length of data dirs in DatabaseDescriptor - + public Integer memtable_total_space_in_mb; + public Integer sliced_buffer_size_in_kb = 64; public Integer storage_port = 7000; diff --git a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java index c75ab524aa..0bc5eb389e 100644 --- a/src/java/org/apache/cassandra/config/DatabaseDescriptor.java +++ b/src/java/org/apache/cassandra/config/DatabaseDescriptor.java @@ -231,6 +231,9 @@ public class DatabaseDescriptor throw new ConfigurationException("conf.concurrent_replicates must be at least 2"); } + if (conf.memtable_total_space_in_mb == null) + conf.memtable_total_space_in_mb = (int) (Runtime.getRuntime().maxMemory() / (3 * 1048576)); + /* Memtable flush writer threads */ if (conf.memtable_flush_writers != null && conf.memtable_flush_writers < 1) { @@ -797,6 +800,8 @@ public class DatabaseDescriptor maxDiskIndex = i; } } + logger.debug("expected data files size is {}; largest free partition has {} bytes free", + expectedCompactedFileSize, maxFreeDisk); // Load factor of 0.9 we do not want to use the entire disk that is too risky. maxFreeDisk = (long)(0.9 * maxFreeDisk); if( expectedCompactedFileSize < maxFreeDisk ) @@ -1057,4 +1062,16 @@ public class DatabaseDescriptor { return conf.memtable_flush_queue_size; } + + public static int getTotalMemtableSpaceInMB() + { + // should only be called if estimatesRealMemtableSize() is true + assert conf.memtable_total_space_in_mb > 0; + return conf.memtable_total_space_in_mb; + } + + public static boolean estimatesRealMemtableSize() + { + return conf.memtable_total_space_in_mb > 0; + } } diff --git a/src/java/org/apache/cassandra/db/BinaryMemtable.java b/src/java/org/apache/cassandra/db/BinaryMemtable.java index 4b4e2ff150..663cc0065a 100644 --- a/src/java/org/apache/cassandra/db/BinaryMemtable.java +++ b/src/java/org/apache/cassandra/db/BinaryMemtable.java @@ -125,7 +125,7 @@ public class BinaryMemtable implements IFlushable private SSTableReader writeSortedContents(List sortedKeys) throws IOException { logger.info("Writing " + this); - SSTableWriter writer = cfs.createFlushWriter(sortedKeys.size()); + SSTableWriter writer = cfs.createFlushWriter(sortedKeys.size(), DatabaseDescriptor.getBMTThreshold()); for (DecoratedKey key : sortedKeys) { diff --git a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java index cb2efae9cd..ea3c2766a8 100644 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@ -24,6 +24,7 @@ import java.nio.ByteBuffer; import java.util.*; import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.concurrent.locks.Lock; import java.util.concurrent.locks.ReentrantLock; @@ -102,6 +103,20 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean public static final ExecutorService postFlushExecutor = new JMXEnabledThreadPoolExecutor("MemtablePostFlusher"); + static + { + if (DatabaseDescriptor.estimatesRealMemtableSize()) + { + logger.info("Global memtable threshold is enabled at {}MB", DatabaseDescriptor.getTotalMemtableSpaceInMB()); + // (can block if flush queue fills up, so don't put on scheduledTasks) + StorageService.tasks.scheduleWithFixedDelay(new MeteredFlusher(), 1000, 1000, TimeUnit.MILLISECONDS); + } + else + { + logger.info("Global memtable threshold is disabled"); + } + } + public final Table table; public final String columnFamily; public final CFMetaData metadata; @@ -143,7 +158,7 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean /** Lock to allow migrations to block all flushing, so we can be sure not to write orphaned data files */ public final Lock flushLock = new ReentrantLock(); - + public static enum CacheType { KEY_CACHE_TYPE("KeyCache"), @@ -166,6 +181,12 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean public final AutoSavingCache, Long> keyCache; public final AutoSavingCache rowCache; + + /** ratio of in-memory memtable size, to serialized size */ + volatile double liveRatio = 1.0; + /** ops count last time we computed liveRatio */ + private final AtomicLong liveRatioComputedAt = new AtomicLong(32); + public void reload() { // metadata object has been mutated directly. make all the members jibe with new settings. @@ -569,12 +590,11 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean * When the sstable object is closed, it will be renamed to a non-temporary * format, so incomplete sstables can be recognized and removed on startup. */ - public String getFlushPath() + public String getFlushPath(long estimatedSize) { - long guessedSize = 2L * memsize.value() * 1024*1024; // 2* adds room for keys, column indexes - String location = DatabaseDescriptor.getDataFileLocationForTable(table.name, guessedSize); + String location = DatabaseDescriptor.getDataFileLocationForTable(table.name, estimatedSize); if (location == null) - throw new RuntimeException("Insufficient disk space to flush"); + throw new RuntimeException("Insufficient disk space to flush " + estimatedSize + " bytes"); return getTempSSTablePath(location); } @@ -733,7 +753,23 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean if (cachedRow != null) cachedRow.addAll(columnFamily); writeStats.addNano(System.nanoTime() - start); - + + if (DatabaseDescriptor.estimatesRealMemtableSize()) + { + while (true) + { + long last = liveRatioComputedAt.get(); + long operations = writeStats.getOpCount(); + if (operations < 2 * last) + break; + if (liveRatioComputedAt.compareAndSet(last, operations)) + { + logger.debug("computing liveRatio of {} at {} ops", this, operations); + mt.updateLiveRatio(); + } + } + } + return flushRequested ? mt : null; } @@ -959,12 +995,20 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean public long getMemtableColumnsCount() { - return getMemtableThreadSafe().getCurrentOperations(); + return getMemtableThreadSafe().getOperations(); } public long getMemtableDataSize() { - return getMemtableThreadSafe().getCurrentThroughput(); + return getMemtableThreadSafe().getLiveSize(); + } + + public long getTotalMemtableLiveSize() + { + long total = 0; + for (ColumnFamilyStore cfs : concatWithIndexes()) + total += cfs.getMemtableThreadSafe().getLiveSize(); + return total; } public int getMemtableSwitchCount() @@ -2022,9 +2066,9 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return intern(name); } - public SSTableWriter createFlushWriter(long estimatedRows) throws IOException + public SSTableWriter createFlushWriter(long estimatedRows, long estimatedSize) throws IOException { - return new SSTableWriter(getFlushPath(), estimatedRows, metadata, partitioner); + return new SSTableWriter(getFlushPath(estimatedSize), estimatedRows, metadata, partitioner); } public SSTableWriter createCompactionWriter(long estimatedRows, String location) throws IOException @@ -2037,4 +2081,8 @@ public class ColumnFamilyStore implements ColumnFamilyStoreMBean return Iterables.concat(indexedColumns.values(), Collections.singleton(this)); } + public Set getMemtablesPendingFlush() + { + return data.getMemtablesPendingFlush(); + } } diff --git a/src/java/org/apache/cassandra/db/Memtable.java b/src/java/org/apache/cassandra/db/Memtable.java index 9cb214149d..25ed4f476f 100644 --- a/src/java/org/apache/cassandra/db/Memtable.java +++ b/src/java/org/apache/cassandra/db/Memtable.java @@ -25,10 +25,7 @@ import java.util.Collection; import java.util.Comparator; import java.util.Iterator; import java.util.Map; -import java.util.concurrent.ConcurrentNavigableMap; -import java.util.concurrent.ConcurrentSkipListMap; -import java.util.concurrent.CountDownLatch; -import java.util.concurrent.ExecutorService; +import java.util.concurrent.*; import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicLong; @@ -38,6 +35,7 @@ import com.google.common.collect.PeekingIterator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import org.apache.cassandra.concurrent.DebuggableThreadPoolExecutor; import org.apache.cassandra.db.columniterator.IColumnIterator; import org.apache.cassandra.db.columniterator.SimpleAbstractColumnIterator; import org.apache.cassandra.db.filter.AbstractColumnIterator; @@ -47,13 +45,32 @@ import org.apache.cassandra.db.marshal.AbstractType; import org.apache.cassandra.io.sstable.SSTableReader; import org.apache.cassandra.io.sstable.SSTableWriter; import org.apache.cassandra.utils.WrappedRunnable; +import org.github.jamm.MemoryMeter; public class Memtable implements Comparable, IFlushable { private static final Logger logger = LoggerFactory.getLogger(Memtable.class); - private volatile boolean isFrozen; + // size in memory can never be less than serialized size + private static final double MIN_SANE_LIVE_RATIO = 1.0; + // max liveratio seen w/ 1-byte columns on a 64-bit jvm was 19. If it gets higher than 64 something is probably broken. + private static final double MAX_SANE_LIVE_RATIO = 64.0; + private static final MemoryMeter meter = new MemoryMeter(); + // we're careful to only allow one count to run at a time because counting is slow + // (can be minutes, for a large memtable and a busy server), so we could keep memtables + // alive after they're flushed and would otherwise be GC'd. + private static final ExecutorService meterExecutor = new ThreadPoolExecutor(1, 1, Integer.MAX_VALUE, TimeUnit.MILLISECONDS, new SynchronousQueue()) + { + @Override + protected void afterExecute(Runnable r, Throwable t) + { + super.afterExecute(r, t); + DebuggableThreadPoolExecutor.logExceptionsAfterExecute(r, t); + } + }; + + private volatile boolean isFrozen; private final AtomicLong currentThroughput = new AtomicLong(0); private final AtomicLong currentOperations = new AtomicLong(0); @@ -63,10 +80,10 @@ public class Memtable implements Comparable, IFlushable private final long THRESHOLD; private final long THRESHOLD_COUNT; + volatile static Memtable activelyMeasuring; public Memtable(ColumnFamilyStore cfs) { - this.cfs = cfs; creationTime = System.currentTimeMillis(); THRESHOLD = cfs.getMemtableThroughputInMB() * 1024L * 1024L; @@ -90,12 +107,18 @@ public class Memtable implements Comparable, IFlushable return 0; } - public long getCurrentThroughput() + public long getLiveSize() + { + // 25% fudge factor + return (long) (currentThroughput.get() * cfs.liveRatio * 1.25); + } + + public long getSerializedSize() { return currentThroughput.get(); } - - public long getCurrentOperations() + + public long getOperations() { return currentOperations.get(); } @@ -126,6 +149,54 @@ public class Memtable implements Comparable, IFlushable resolve(key, columnFamily); } + public void updateLiveRatio() + { + Runnable runnable = new Runnable() + { + public void run() + { + activelyMeasuring = Memtable.this; + + long start = System.currentTimeMillis(); + // ConcurrentSkipListMap has cycles, so measureDeep will have to track a reference to EACH object it visits. + // So to reduce the memory overhead of doing a measurement, we break it up to row-at-a-time. + long deepSize = meter.measure(columnFamilies); + int objects = 0; + for (Map.Entry entry : columnFamilies.entrySet()) + { + deepSize += meter.measureDeep(entry.getKey()) + meter.measureDeep(entry.getValue()); + objects += entry.getValue().getColumnCount(); + } + double newRatio = (double) deepSize / currentThroughput.get(); + + if (newRatio < MIN_SANE_LIVE_RATIO) + { + logger.warn("setting live ratio to minimum of 1.0 instead of {}", newRatio); + newRatio = MIN_SANE_LIVE_RATIO; + } + if (newRatio > MAX_SANE_LIVE_RATIO) + { + logger.warn("setting live ratio to maximum of 64 instead of {}, newRatio"); + newRatio = MAX_SANE_LIVE_RATIO; + } + cfs.liveRatio = Math.max(cfs.liveRatio, newRatio); + + logger.info("{} liveRatio is {} (just-counted was {}). calculation took {}ms for {} columns", + new Object[]{ cfs, cfs.liveRatio, newRatio, System.currentTimeMillis() - start, objects }); + activelyMeasuring = null; + } + }; + + try + { + meterExecutor.submit(runnable); + } + catch (RejectedExecutionException e) + { + logger.debug("Meter thread is busy; skipping liveRatio update for {}", cfs); + } + } + private void resolve(DecoratedKey key, ColumnFamily cf) { currentThroughput.addAndGet(cf.size()); @@ -155,8 +226,10 @@ public class Memtable implements Comparable, IFlushable private SSTableReader writeSortedContents() throws IOException { logger.info("Writing " + this); - SSTableWriter writer = cfs.createFlushWriter(columnFamilies.size()); + SSTableWriter writer = cfs.createFlushWriter(columnFamilies.size(), 2 * getSerializedSize()); // 2* for keys + // (we can't clear out the map as-we-go to free up memory, + // since the memtable is being used for queries in the "pending flush" category) for (Map.Entry entry : columnFamilies.entrySet()) writer.append(entry.getKey(), entry.getValue()); @@ -192,8 +265,8 @@ public class Memtable implements Comparable, IFlushable public String toString() { - return String.format("Memtable-%s@%s(%s bytes, %s operations)", - cfs.getColumnFamilyName(), hashCode(), currentThroughput, currentOperations); + return String.format("Memtable-%s@%s(%s/%s serialized/live bytes, %s ops)", + cfs.getColumnFamilyName(), hashCode(), currentThroughput, getLiveSize(), currentOperations); } /** diff --git a/src/java/org/apache/cassandra/db/MeteredFlusher.java b/src/java/org/apache/cassandra/db/MeteredFlusher.java new file mode 100644 index 0000000000..68230c3c18 --- /dev/null +++ b/src/java/org/apache/cassandra/db/MeteredFlusher.java @@ -0,0 +1,104 @@ +package org.apache.cassandra.db; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.Comparator; +import java.util.List; + +import com.google.common.collect.Iterables; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.apache.cassandra.config.DatabaseDescriptor; + +class MeteredFlusher implements Runnable +{ + private static Logger logger = LoggerFactory.getLogger(MeteredFlusher.class); + + public void run() + { + // first, find how much memory non-active memtables are using + Memtable activelyMeasuring = Memtable.activelyMeasuring; + long flushingBytes = activelyMeasuring == null ? 0 : activelyMeasuring.getLiveSize(); + flushingBytes += countFlushingBytes(); + + // next, flush CFs using more than 1 / (maximum number of memtables it could have in the pipeline) + // of the total size allotted. Then, flush other CFs in order of size if necessary. + long liveBytes = 0; + try + { + for (ColumnFamilyStore cfs : ColumnFamilyStore.all()) + { + long size = cfs.getTotalMemtableLiveSize(); + int maxInFlight = (int) Math.ceil((double) (1 // live memtable + + 1 // potentially a flushed memtable being counted by jamm + + DatabaseDescriptor.getFlushWriters() + + DatabaseDescriptor.getFlushQueueSize()) + / (1 + cfs.getIndexedColumns().size())); + if (size > (DatabaseDescriptor.getTotalMemtableSpaceInMB() * 1048576L - flushingBytes) / maxInFlight) + { + logger.info("flushing high-traffic column family {}", cfs); + cfs.forceFlush(); + } + else + { + liveBytes += size; + } + } + + if (flushingBytes + liveBytes <= DatabaseDescriptor.getTotalMemtableSpaceInMB() * 1048576L) + return; + + logger.info("estimated {} bytes used by all memtables pre-flush", liveBytes); + + // sort memtables by size + List sorted = new ArrayList(); + Iterables.addAll(sorted, ColumnFamilyStore.all()); + Collections.sort(sorted, new Comparator() + { + public int compare(ColumnFamilyStore o1, ColumnFamilyStore o2) + { + long size1 = o1.getTotalMemtableLiveSize(); + long size2 = o2.getTotalMemtableLiveSize(); + if (size1 < size2) + return -1; + if (size1 > size2) + return 1; + return 0; + } + }); + + // flush largest first until we get below our threshold. + // although it looks like liveBytes + flushingBytes will stay a constant, it will not if flushes finish + // while we loop, which is especially likely to happen if the flush queue fills up (so further forceFlush calls block) + while (true) + { + flushingBytes = countFlushingBytes(); + if (liveBytes + flushingBytes <= DatabaseDescriptor.getTotalMemtableSpaceInMB() * 1048576L || sorted.isEmpty()) + break; + + ColumnFamilyStore cfs = sorted.remove(sorted.size() - 1); + long size = cfs.getTotalMemtableLiveSize(); + logger.info("flushing {} to free up {} bytes", cfs, size); + liveBytes -= size; + cfs.forceFlush(); + } + } + finally + { + logger.debug("memtable memory usage is {} bytes with {} live", liveBytes + flushingBytes, liveBytes); + } + } + + private long countFlushingBytes() + { + long flushingBytes = 0; + for (ColumnFamilyStore cfs : ColumnFamilyStore.all()) + { + for (Memtable memtable : cfs.getMemtablesPendingFlush()) + flushingBytes += memtable.getLiveSize(); + } + return flushingBytes; + } +} diff --git a/src/java/org/apache/cassandra/service/StorageService.java b/src/java/org/apache/cassandra/service/StorageService.java index 42fd07ba85..8d05b7e996 100644 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@ -2191,7 +2191,7 @@ public class StorageService implements IEndpointStateChangeSubscriber, StorageSe for (ColumnFamilyStore subordinate : cfs.concatWithIndexes()) { ops += subordinate.getMemtableColumnsCount(); - throughput = subordinate.getMemtableThroughputInMB(); + throughput += subordinate.getMemtableDataSize(); } if (ops > 0 && (largestByOps == null || ops > largestByOps.getMemtableColumnsCount())) diff --git a/src/java/org/apache/cassandra/streaming/StreamIn.java b/src/java/org/apache/cassandra/streaming/StreamIn.java index 5ad87dfde6..fe5d850122 100644 --- a/src/java/org/apache/cassandra/streaming/StreamIn.java +++ b/src/java/org/apache/cassandra/streaming/StreamIn.java @@ -80,7 +80,7 @@ public class StreamIn // new local sstable Table table = Table.open(remotedesc.ksname); ColumnFamilyStore cfStore = table.getColumnFamilyStore(remotedesc.cfname); - Descriptor localdesc = Descriptor.fromFilename(cfStore.getFlushPath()); + Descriptor localdesc = Descriptor.fromFilename(cfStore.getFlushPath(remote.size)); return new PendingFile(localdesc, remote); } diff --git a/test/conf/cassandra.yaml b/test/conf/cassandra.yaml index a3151be2b7..c303136993 100644 --- a/test/conf/cassandra.yaml +++ b/test/conf/cassandra.yaml @@ -33,3 +33,4 @@ encryption_options: truststore: conf/.truststore truststore_password: cassandra incremental_backups: true +flush_largest_memtables_at: 1.0 diff --git a/test/long/org/apache/cassandra/db/MeteredFlusherTest.java b/test/long/org/apache/cassandra/db/MeteredFlusherTest.java new file mode 100644 index 0000000000..4e507d21bd --- /dev/null +++ b/test/long/org/apache/cassandra/db/MeteredFlusherTest.java @@ -0,0 +1,51 @@ +package org.apache.cassandra.db; + +import java.io.IOException; +import java.nio.ByteBuffer; + +import org.junit.Test; + +import org.apache.cassandra.CleanupHelper; +import org.apache.cassandra.config.CFMetaData; +import org.apache.cassandra.config.ConfigurationException; +import org.apache.cassandra.db.marshal.UTF8Type; +import org.apache.cassandra.db.migration.AddColumnFamily; +import org.apache.cassandra.utils.ByteBufferUtil; + +public class MeteredFlusherTest extends CleanupHelper +{ + @Test + public void testManyMemtables() throws IOException, ConfigurationException + { + Table table = Table.open("Keyspace1"); + for (int i = 0; i < 100; i++) + { + CFMetaData metadata = new CFMetaData(table.name, "_CF" + i, ColumnFamilyType.Standard, UTF8Type.instance, null); + new AddColumnFamily(metadata).apply(); + } + + ByteBuffer name = ByteBufferUtil.bytes("c"); + for (int j = 0; j < 200; j++) + { + for (int i = 0; i < 100; i++) + { + RowMutation rm = new RowMutation("Keyspace1", ByteBufferUtil.bytes("key" + j)); + ColumnFamily cf = ColumnFamily.create("Keyspace1", "_CF" + i); + // don't cheat by allocating this outside of the loop; that defeats the purpose of deliberately using lots of memory + ByteBuffer value = ByteBuffer.allocate(100000); + cf.addColumn(new Column(name, value)); + rm.add(cf); + rm.applyUnsafe(); + } + } + + int flushes = 0; + for (ColumnFamilyStore cfs : ColumnFamilyStore.all()) + { + if (cfs.getColumnFamilyName().startsWith("_CF")) + flushes += cfs.getMemtableSwitchCount(); + } + assert flushes > 0; + } +} + diff --git a/test/unit/org/apache/cassandra/db/DefsTest.java b/test/unit/org/apache/cassandra/db/DefsTest.java index 0461a9502b..35da85e30b 100644 --- a/test/unit/org/apache/cassandra/db/DefsTest.java +++ b/test/unit/org/apache/cassandra/db/DefsTest.java @@ -319,7 +319,7 @@ public class DefsTest extends CleanupHelper ColumnFamilyStore store = Table.open(cfm.ksName).getColumnFamilyStore(cfm.cfName); assert store != null; store.forceBlockingFlush(); - store.getFlushPath(); + store.getFlushPath(1024); assert DefsTable.getFiles(cfm.ksName, cfm.cfName).size() > 0; new DropColumnFamily(ks.name, cfm.cfName).apply();