Skip to content

[STORM-855] Add tuple batching - #694

Closed
mjsax wants to merge 3 commits into
apache:masterfrom
mjsax:batching
Closed

[STORM-855] Add tuple batching#694
mjsax wants to merge 3 commits into
apache:masterfrom
mjsax:batching

Conversation

@mjsax

Copy link
Copy Markdown
Member

added new parameters to TopologyBuilder

  • batch_size (int): set the same batch size for all UDF declared output streams
  • batch_sizes (Map<String,Number>): sets the batch sizer for individual streams (including system streams)

added Batch type (as alternative to Tuple)

  • extended Kryo (de)serialization for Batch
  • integrated Tuple and Batch (de)serialization

add tuple batching

  • added internal tuple buffers (one for each output stream)
  • added Spout/Bolt output batching
  • added Bolt input debatching
  • added batching of "acks"

@HeartSaVioR

Copy link
Copy Markdown
Contributor

Could you fix typo in title, to let asfgit matches PR and JIRA issue properly? Thanks!

@mjsaxmjsax changed the title [STROM-855] Add tuple batching[STORM-855] Add tuple batchingAug 20, 2015
@mjsax

Copy link
Copy Markdown
MemberAuthor

I guess I need to close an re-open this PR to get the linkage to JIRA...

@HeartSaVioR

Copy link
Copy Markdown
Contributor

asfbot recognizes your change and links. :)

@mjsax

Copy link
Copy Markdown
MemberAuthor

Great!
Just a heads up: This is my first PR for Storm and I just started to learn Clojure. Please review very carefully. Looking forward to your feedback.

@HeartSaVioR

Copy link
Copy Markdown
Contributor

@mjsax
I'm not familiar with clojure, too. Many committers will take a look, so don't worry. :)

Before taking a look, I think you're encouraged to do benchmark and attach results, since it modifies critical path, especially latency vs throughput.

@HeartSaVioR

Copy link
Copy Markdown
Contributor

FYI, you can use https://github.com/yahoo/storm-perf-test for benchmarking if you don't have your own.
You need to modify pom.xml to let it points recent version of Storm and build.

@mjsax

Copy link
Copy Markdown
MemberAuthor

Sure. Is there any specific approach I should take? How to add/report those result? I did a simple test already, using ExclamationTopology example. Setting batch size to 100 in Spout roughly doubled throughput and increased latency from 1 to 3 seconds. But this is of course not representative.

@HeartSaVioR

Copy link
Copy Markdown
Contributor

Actually there's no specific approach.
FYI you can refer @d2r's approach, #521 (comment)

I think we would be interesting to compare these - not applying this patch, applying this patch & disable batch, applying this patch & some kinds of batch size.

@mjsax

Copy link
Copy Markdown
MemberAuthor

Thanks for your guidance! I will have a look at #521 and add some results. May take a few days...

@HeartSaVioR

Copy link
Copy Markdown
Contributor

Sure, please take your time!

@mjsax

Copy link
Copy Markdown
MemberAuthor

Hi, I did a few performance tests and unfortunately, the impact of my changes is quite high. Not using batching in my branch (compared to master) reduces throughput by 40%. :(

I dug into the problem and it seems that there are two critical point in the code. First, I introduced function emit-msg (in executor.clj). This function uses a Java HashMap (output-batch-buffer) and it branches. Both have a large negative impact. Can it be that the access via HashMap-get() is slow in Clojure? What about branching? To me it seems that there performance impact is ridiculously high (even if I put into account that this code is called a few 100,000 times a second). Especially branch prediction should avoid the branching overhead at all. Because batching was disabled in the test, the else branch is taken every time...

Furthermore, I changed the serialization by overloading KryoTuple(Batch)Serializer.serialize(...) (I renamed this class). It seems that this overloading makes the call to .serialize(Tuple) much more expensive. I was thinking about changing the code, such that even if batching is disabled, I just use batches of size 1 to eliminate the overload. Of course, if comes with the drawback, that a single tuple consumers more space as I need to write the batch size in front of each batch.

Does this observation make sense? Or do you think I oversee something?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Try type-hinting output-batch-buffer here. I haven't profiled it myself, but what you're saying as far as performance goes sounds like it could be the result of reflection. Type-hinting anything you're calling methods on should help.

@knusbaum

Copy link
Copy Markdown
Contributor

What happens when Batch.capacity - 1 tuples are emitted and then there's a long pause in the input?

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Unless I misunderstand, this is going to end up being always true or always false, depending on the result of .get_batch_size for the component at startup. In that case, can we decide this when we're setting up the transfer function instead of each time we're emitting a tuple?

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Your observation is right for the current state of the code. However, I would like to extend it with the possibility to use different batch sizes for different output streams (including mixed-mode batching/non-batching). Nevertheless, we could have 3 functions: non-batching, batching, mixed-mode. Thus, we could choose the correct function at setup time. The branching overhead is avoided for the both main cases (I guess, that in most cases there is only a single output stream).

@mjsax

Copy link
Copy Markdown
MemberAuthor

Type hinting improves by 10% :) But 30% is still a huge gap.
About batches that don't fill up. We need to introduce a timeout (and a flushing thread or other flushing strategy). Otherwise, a batch might starve. I just wanted to get the basic design right before taking care of this problem.

@knusbaum

Copy link
Copy Markdown
Contributor

Keep an eye out for more missing types. Those are often a source of unexpectedly bad performance.
I'm hesitant to add another thread, especially since we're operating on the critical path. We need to figure out some other way to flush.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This line, for instance, is eating 9% of the CPU time.
screen shot 2015-08-24 at 4 26 07 pm

Copy link
Copy Markdown
MemberAuthor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks. However, this is not related to the decreased throughput I observed, because the spout output rate limits the throughput in my experiment.

@knusbaum

Copy link
Copy Markdown
Contributor

@mjsax, I'm not sure it's not related. Here is the output from my benchmarks:

Pre-batching (apache master)

status time time-diff ms transferred throughput (MB/s)
WAITING 1440456969132 0 0 0.0
WAITING 1440456999132 30000 99560 0.031649271647135414
WAITING 1440457029131 29999 2073460 0.6591571546037212
RUNNING 1440457059131 30000 2043680 0.6496683756510416
RUNNING 1440457089132 30001 2000600 0.6359524140536461
RUNNING 1440457119133 30001 1921980 0.6109606221947549
RUNNING 1440457149131 29998 1846040 0.5868794369819967
RUNNING 1440457179132 30001 1847640 0.5873293603429365
RUNNING 1440457209133 30001 1688440 0.5367227301733171
RUNNING 1440457239131 29998 1733740 0.5511778482986105
RUNNING 1440457269131 30000 1751700 0.5568504333496094
RUNNING 1440457299133 30002 1748280 0.5557261962158252
RUNNING 1440457329132 29999 1706500 0.54249982364321
RUNNING 1440457359133 30001 1787720 0.568281940243919
RUNNING 1440457389131 29998 1717900 0.5461421121922451
RUNNING 1440457419132 30001 1779180 0.5655672378466291
RUNNING 1440457449132 30000 1639820 0.5212847391764323
RUNNING 1440457479131 29999 1743720 0.5543321374058823
RUNNING 1440457509134 30003 1681140 0.5343665767700574
RUNNING 1440457539133 29999 1716760 0.5457614985278155
RUNNING 1440457569131 29998 1700800 0.5407058061683279
RUNNING 1440457599132 30001 1754100 0.5575947863098574
RUNNING 1440457629134 30002 1669220 0.53059537445225
RUNNING 1440457659132 29998 1757240 0.5586487951735845
RUNNING 1440457689131 29999 1781520 0.5663488343491658
RUNNING 1440457719131 30000 1772920 0.5635960896809896
RUNNING 1440457749134 30003 1656020 0.5263819422907969
RUNNING 1440457779135 30001 1557820 0.49520113449017844
RUNNING 1440457809134 29999 1750760 0.5565701677360599
RUNNING 1440457839131 29997 1773740 0.5639131519760049
RUNNING 1440457869132 30001 1656360 0.5265251127371275
RUNNING 1440457899131 29999 1770140 0.5627311091847593
RUNNING 1440457929131 30000 1767480 0.5618667602539062

With batching, no typehints. (This pull request as of now)

status time time-diff ms transferred throughput (MB/s)
WAITING 1440455622424 0 0 0.0
WAITING 1440455652425 30001 87340 0.027763712807880363
WAITING 1440455682427 30002 652540 0.2074230512724933
RUNNING 1440455712428 30001 678920 0.21581566177611788
RUNNING 1440455742425 29997 729960 0.23207124179214797
RUNNING 1440455772425 30000 667320 0.21213531494140625
RUNNING 1440455802427 30002 670380 0.2130938564870415
RUNNING 1440455832423 29996 616080 0.19587267397371733
RUNNING 1440455862423 30000 645540 0.20521163940429688
RUNNING 1440455892428 30005 635080 0.2018528528122917
RUNNING 1440455922427 29999 622760 0.1979766716507738
RUNNING 1440455952428 30001 610080 0.19393274455955634
RUNNING 1440455982424 29996 633800 0.20150646144095255
RUNNING 1440456012423 29999 573820 0.18241854603161253
RUNNING 1440456042427 30004 596680 0.18965417648089627
RUNNING 1440456072423 29996 592880 0.18849660911819494
RUNNING 1440456102425 30002 558000 0.17737159807835728
RUNNING 1440456132424 29999 614140 0.19523635610444828
RUNNING 1440456162423 29999 580520 0.18454848966970774
RUNNING 1440456192423 30000 585380 0.18608729044596353
RUNNING 1440456222425 30002 568560 0.18072830789145308
RUNNING 1440456252423 29998 598440 0.19025163607912401
RUNNING 1440456282426 30003 560560 0.1781794069941964
RUNNING 1440456312423 29997 579480 0.1842301539724285
RUNNING 1440456342424 30001 589860 0.18750519391866624
RUNNING 1440456372428 30004 569560 0.18103411000278088
RUNNING 1440456402423 29995 536220 0.17048816200812114
RUNNING 1440456432427 30004 607340 0.19304244744906407
RUNNING 1440456462425 29998 563460 0.17913105217756703
RUNNING 1440456492427 30002 588020 0.18691406290687393
RUNNING 1440456522423 29996 594640 0.18905617265895872
RUNNING 1440456552423 30000 594260 0.18891016642252603
RUNNING 1440456582426 30003 582340 0.18510239023298192
RUNNING 1440456612425 29999 554180 0.17617495005367367

With batching, typehints added:

status time time-diff ms transferred throughput (MB/s)
WAITING 1440453725066 0 0 0.0
WAITING 1440453755065 29999 78920 0.025088828644548568
WAITING 1440453785067 30002 2003120 0.6367322500765574
RUNNING 1440453815066 29999 1940420 0.6168634677959318
RUNNING 1440453845065 29999 2073440 0.6591507965630105
RUNNING 1440453875068 30003 1952840 0.6207290444458158
RUNNING 1440453905067 29999 1751360 0.556760908957382
RUNNING 1440453935068 30001 1823440 0.5796366439478059
RUNNING 1440453965065 29997 1835720 0.5836180338411445
RUNNING 1440453995070 30005 1776660 0.5646908885140237
RUNNING 1440454025066 29996 1599600 0.5085669544350705
RUNNING 1440454055065 29999 1566060 0.4978536617724497
RUNNING 1440454085068 30003 1632100 0.5187787393949407
RUNNING 1440454115065 29997 1619780 0.5149656913119698
RUNNING 1440454145066 30001 1665920 0.5295640535940468
RUNNING 1440454175068 30002 1697340 0.5395338858106075
RUNNING 1440454205065 29997 1705360 0.542173561431664
RUNNING 1440454235067 30002 1628760 0.5177343442403319
RUNNING 1440454265066 29999 1721040 0.5471221192399122
RUNNING 1440454295070 30004 1731960 0.5505018561001762
RUNNING 1440454325067 29997 1648780 0.5241854650146004
RUNNING 1440454355065 29998 1665100 0.529356325171027
RUNNING 1440454385065 30000 1573460 0.500189463297526
RUNNING 1440454415068 30003 1700960 0.5406665550892827
RUNNING 1440454445066 29998 1744260 0.554522289197535
RUNNING 1440454475065 29999 1580540 0.5024568832470198
RUNNING 1440454505068 30003 1695720 0.5390009705084179
RUNNING 1440454535066 29998 1746140 0.5551199649475329
RUNNING 1440454565065 29999 1685440 0.5358048067748091
RUNNING 1440454595066 30001 1738360 0.5525913418445947
RUNNING 1440454625069 30003 1431160 0.4549080207539142
RUNNING 1440454655065 29996 1421760 0.45202560211153153
RUNNING 1440454685070 30005 1403440 0.4460672163363398

With the type-hints in place, the throughput goes up to very near apache-master levels.
These were gathered with the command:
storm jar storm_perf_test-1.0.0-SNAPSHOT-jar-with-dependencies.jar com.yahoo.storm.perftest.Main --ack --bolt 4 --name test -l 1 -n 1 --workers 4 --spout 3 --testTimeSec 900 -c topology.max.spout.pending=1092 --messageSize 10

@mjsax

Copy link
Copy Markdown
MemberAuthor

I just pushed some changes (I added a new commit, so you can better see what I changed):

  • added type hints
  • split tuple and batch serialization in separate classes
  • assemble different "emit function" in Clojure for single tuple and batch case (to add more type hints)

I get the following result running on a 4 node cluster with parameters: storm jar storm_perf_test-1.0.0-SNAPSHOT-jar-with-dependencies.jar com.yahoo.storm.perftest.Main --bolt 3 --name test -l 1 -n 1 --messageSize 4 --workers 4 --spout 1 --testTimeSec 300

Master Branch:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 48 0 4 0 1440466170638 0 0 0.0
WAITING 1 48 4 4 4 1440466200638 30000 6122840 0.7785593668619791
WAITING 1 48 4 4 4 1440466230638 30000 11565400 1.4706166585286458
RUNNING 1 48 4 4 4 1440466260638 30000 11394040 1.4488271077473958
RUNNING 1 48 4 4 4 1440466290638 30000 11718240 1.49005126953125
RUNNING 1 48 4 4 4 1440466320638 30000 11615920 1.4770406087239583
RUNNING 1 48 4 4 4 1440466350638 30000 11557380 1.4695968627929688
RUNNING 1 48 4 4 4 1440466380638 30000 11581080 1.4726104736328125
RUNNING 1 48 4 4 4 1440466410638 30000 11492600 1.4613596598307292
RUNNING 1 48 4 4 4 1440466440638 30000 11413760 1.4513346354166667
RUNNING 1 48 4 4 4 1440466470638 30000 11300580 1.4369430541992188
RUNNING 1 48 4 4 4 1440466500638 30000 11368760 1.4456125895182292
RUNNING 1 48 4 4 4 1440466530638 30000 11509820 1.463549296061198

Batching branch with batching disabled:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 48 0 4 0 1440467016767 0 0 0.0
WAITING 1 48 4 4 4 1440467046767 30000 7095940 0.9022954305013021
WAITING 1 48 4 4 4 1440467076767 30000 11136640 1.4160970052083333
RUNNING 1 48 4 4 4 1440467106767 30000 11159220 1.4189682006835938
RUNNING 1 48 4 4 4 1440467136767 30000 7757660 0.9864374796549479
RUNNING 1 48 4 4 4 1440467166767 30000 11375580 1.4464797973632812
RUNNING 1 48 4 4 4 1440467196767 30000 11669980 1.4839146931966145
RUNNING 1 48 4 4 4 1440467226767 30000 11344380 1.4425125122070312
RUNNING 1 48 4 4 4 1440467256767 30000 11521460 1.4650293986002605
RUNNING 1 48 4 4 4 1440467286767 30000 11401040 1.4497172037760417
RUNNING 1 48 4 4 4 1440467316767 30000 11493700 1.461499532063802
RUNNING 1 48 4 4 4 1440467346767 30000 11452680 1.4562835693359375
RUNNING 1 48 4 4 4 1440467376767 30000 11148300 1.4175796508789062

Batching branch with batch size of 100 tuples:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 48 1 4 0 1440467461710 0 0 0.0
WAITING 1 48 4 4 4 1440467491710 30000 11686000 1.4859517415364583
WAITING 1 48 4 4 4 1440467521710 30000 18026640 2.292205810546875
RUNNING 1 48 4 4 4 1440467551710 30000 17936300 2.2807184855143228
RUNNING 1 48 4 4 4 1440467581710 30000 18969300 2.4120712280273438
RUNNING 1 48 4 4 4 1440467611710 30000 18581620 2.3627751668294272
RUNNING 1 48 4 4 4 1440467641711 30001 18963120 2.4112050268897285
RUNNING 1 48 4 4 4 1440467671710 29999 18607200 2.3661067022546587
RUNNING 1 48 4 4 4 1440467701710 30000 19333620 2.4583969116210938
RUNNING 1 48 4 4 4 1440467731710 30000 18629100 2.3688125610351562
RUNNING 1 48 4 4 4 1440467761711 30001 18847820 2.3965443624209923
RUNNING 1 48 4 4 4 1440467791710 29999 18021400 2.291615897287722
RUNNING 1 48 4 4 4 1440467821710 30000 18143360 2.3070475260416665

The negative impact is gone and batching increases output rate by about 50%. Need to do more tests. Also need to investigate the performance impact of input debachting. Furthermore, need to test with acking enabled.

Some more question:

  • What about assert-can-serialize? Is it performance critical? Did not test it, but it seems that a generic approach for tuple and batch should be good enough.
  • What about batching acks? Would it make sense? I don't understand the acking code path good enough right now to judge. As acking is quite expensive, it might be a good idea.

@knusbaum

Copy link
Copy Markdown
Contributor

Excellent work. This looks very promising.
I'll address your questions and review the new code tomorrow.
I am very interested, though, in how the size of the messages affects the throughput. Have you tried with larger --messageSize?

@mjsax

Copy link
Copy Markdown
MemberAuthor

I just checked some older benchmark result doing batching in user land, ie, on top of Storm (=> Aeolus). For this case, a batch size of 100 increased the spout output rate by a factor of 6 (instead of 1.5 as the benchmark above shows). The benchmark should yield more than 70M tuples per 30 seconds... (and not about 19M).

Of course, batching is done a little different now. In Aeolus, a fat-tuple is used as batch. Thus, the system sees only a single batch-tuple. Now, the system sees all tuples, but emitting is delayed until the batch is full (this still saved the overhead of going through the disruptor for each tuple). However, we generate a tuple-ID for each tuple in the batch, instead of a single ID per batch. Not sure how expensive this is. Because acking was not enabled, it should not be too expensive, because the IDs have not to be "registered" at the ackers (right?).

As a further optimization, it might be a good idea not to batch whole tuples, but only Values and tuple-id. The worker-context, task-id, and outstream-id is the same for all tuples within a batch. I will try this out, and push a new version the next days if it works.

@mjsax

Copy link
Copy Markdown
MemberAuthor

Here are some additional benchmark results with larger --messageSize (ie, 100 and 250). Those benchmarks are run in a 12 node cluster (with nimbus.thrift.max_buffer_size: 33554432 and worker.childopts: "-Xmx2g") as follows:
storm jar target/storm_perf_test-1.0.0-SNAPSHOT-jar-with-dependencies.jar com.yahoo.storm.perftest.Main --name test -l 1 -n 1 --messageSize 100 --workers 24 --spout 1 --bolt 10 --testTimeSec 300

As you can observe, the output rate is slightly reduced for larger tuple size in all cases (what is expected due to larger serialization costs). That batching case with batch size 100 and--messageSize 250 does not improve output rate as significant as before, because it hits the underlaying 1Gbit Ethernet limit. In the beginning, it starts promising. After about 2 minutes performance does down. Because tuples cannot be transfered over the network fast enough, they get buffered in main memory which reaches it's limit (assumption) and thus Storm slows down the spout. Not sure, why performances goes down in each report. Any ideas?

--messageSize 100
master branch

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440523358523 0 0 0.0
WAITING 1 144 11 11 11 1440523388523 30000 6336040 20.14172871907552
WAITING 1 144 11 11 11 1440523418523 30000 10554200 33.55089823404948
RUNNING 1 144 11 11 11 1440523448523 30000 10005840 31.807708740234375
RUNNING 1 144 11 11 11 1440523478523 30000 10514120 33.4234873453776
RUNNING 1 144 11 11 11 1440523508523 30000 10430700 33.158302307128906
RUNNING 1 144 11 11 11 1440523538523 30000 7608960 24.188232421875
RUNNING 1 144 11 11 11 1440523568523 30000 10279460 32.67752329508463
RUNNING 1 144 11 11 11 1440523598524 30001 10496260 33.36559974774929
RUNNING 1 144 11 11 11 1440523628523 29999 10382860 33.00732328691556
RUNNING 1 144 11 11 11 1440523658523 30000 10138280 32.22872416178385
RUNNING 1 144 11 11 11 1440523688523 30000 10072940 32.02101389567057
RUNNING 1 144 11 11 11 1440523718523 30000 10095820 32.09374745686849

batching branch (no batching)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440520763917 0 0 0.0
WAITING 1 144 11 11 11 1440520793917 30000 4467900 14.203071594238281
WAITING 1 144 11 11 11 1440520823917 30000 10616160 33.74786376953125
RUNNING 1 144 11 11 11 1440520853917 30000 10473700 33.294995625813804
RUNNING 1 144 11 11 11 1440520883917 30000 10556860 33.55935414632162
RUNNING 1 144 11 11 11 1440520913917 30000 10580760 33.63533020019531
RUNNING 1 144 11 11 11 1440520943917 30000 10367580 32.95764923095703
RUNNING 1 144 11 11 11 1440520973917 30000 10646760 33.84513854980469
RUNNING 1 144 11 11 11 1440521003917 30000 10750300 34.17428334554037
RUNNING 1 144 11 11 11 1440521033917 30000 10607220 33.719444274902344
RUNNING 1 144 11 11 11 1440521063917 30000 10456920 33.24165344238281
RUNNING 1 144 11 11 11 1440521093917 30000 10108000 32.132466634114586
RUNNING 1 144 11 11 11 1440521123917 30000 10576120 33.6205800374349

batching branch (batch size: 100)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440521937049 0 0 0.0
WAITING 1 144 11 11 11 1440521967049 30000 11346480 36.069488525390625
WAITING 1 144 11 11 11 1440521997050 30001 17333500 55.0998758822297
RUNNING 1 144 11 11 11 1440522027049 29999 17815260 56.635074176137906
RUNNING 1 144 11 11 11 1440522057049 30000 17993660 57.200304667154946
RUNNING 1 144 11 11 11 1440522087049 30000 17720880 56.333160400390625
RUNNING 1 144 11 11 11 1440522117050 30001 17957200 57.08249869861108
RUNNING 1 144 11 11 11 1440522147049 29999 18286500 58.13315572840058
RUNNING 1 144 11 11 11 1440522177050 30001 18027820 57.306986149778076
RUNNING 1 144 11 11 11 1440522207049 29999 12470520 39.64403692199896
RUNNING 1 144 11 11 11 1440522237049 30000 17256760 54.8577626546224
RUNNING 1 144 11 11 11 1440522267049 30000 17288300 54.95802561442057
RUNNING 1 144 11 11 11 1440522297049 30000 17077300 54.28727467854818

--messageSize 250
master branch

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440523772498 0 0 0.0
WAITING 1 144 11 11 11 1440523802498 30000 6047420 48.06057612101237
WAITING 1 144 11 11 11 1440523832498 30000 9909100 78.7504514058431
RUNNING 1 144 11 11 11 1440523862498 30000 9846400 78.25215657552083
RUNNING 1 144 11 11 11 1440523892498 30000 7100940 56.43320083618164
RUNNING 1 144 11 11 11 1440523922498 30000 9870500 78.4436861673991
RUNNING 1 144 11 11 11 1440523952498 30000 9856460 78.33210627237956
RUNNING 1 144 11 11 11 1440523982498 30000 9662740 76.79255803426106
RUNNING 1 144 11 11 11 1440524012498 30000 10060860 79.9565315246582
RUNNING 1 144 11 11 11 1440524042499 30001 10049660 79.86485975980162
RUNNING 1 144 11 11 11 1440524072498 29999 9937680 78.98021751278428
RUNNING 1 144 11 11 11 1440524102498 30000 10061600 79.96241251627605
RUNNING 1 144 11 11 11 1440524132498 30000 9804960 77.92282104492188

batching branch (no batching)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440521276880 0 0 0.0
WAITING 1 144 11 11 11 1440521306880 30000 4983420 39.60466384887695
WAITING 1 144 11 11 11 1440521336880 30000 9760600 77.57027943929036
RUNNING 1 144 11 11 11 1440521366881 30001 9344540 74.26125626338171
RUNNING 1 144 11 11 11 1440521396880 29999 9323920 74.10232867951068
RUNNING 1 144 11 11 11 1440521426880 30000 9354460 74.3425687154134
RUNNING 1 144 11 11 11 1440521456880 30000 9601440 76.30538940429688
RUNNING 1 144 11 11 11 1440521486880 30000 9476900 75.3156344095866
RUNNING 1 144 11 11 11 1440521516880 30000 9581560 76.14739735921223
RUNNING 1 144 11 11 11 1440521546881 30001 9494500 75.4529915429414
RUNNING 1 144 11 11 11 1440521576880 29999 9260020 73.59448017774095
RUNNING 1 144 11 11 11 1440521606880 30000 9054000 71.95472717285156
RUNNING 1 144 11 11 11 1440521636880 30000 9278340 73.73762130737305
RUNNING 1 144 11 11 11 1440521666880 30000 9153180 72.74293899536133

batching branch (batch size: 100)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 144 0 11 0 1440522380492 0 0 0.0
WAITING 1 144 11 11 11 1440522410492 30000 10616760 84.37442779541016
WAITING 1 144 11 11 11 1440522440492 30000 15851220 125.97417831420898
RUNNING 1 144 11 11 11 1440522470492 30000 12884860 102.39966710408528
RUNNING 1 144 11 11 11 1440522500492 30000 8743860 69.48995590209961
RUNNING 1 144 11 11 11 1440522530492 30000 6445660 51.225503285725914
RUNNING 1 144 11 11 11 1440522560492 30000 5974480 47.48090108235677
RUNNING 1 144 11 11 11 1440522590493 30001 5198640 41.3137016119645
RUNNING 1 144 11 11 11 1440522620492 29999 5209800 41.405150618464624
RUNNING 1 144 11 11 11 1440522650492 30000 4585960 36.445935567220054
RUNNING 1 144 11 11 11 1440522680492 30000 3870140 30.75710932413737
RUNNING 1 144 11 11 11 1440522710492 30000 3751040 29.810587565104168
RUNNING 1 144 11 11 11 1440522740493 30001 3450600 27.421990901898322

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

It would be nice to make this a vector, for consistency.

@revans2

Copy link
Copy Markdown
Contributor

acking does seem to be having some issues. The ack tuples don't seem to be batched, so the serialization code gets confused.

@harshach

Copy link
Copy Markdown
Contributor

@mjsax did you get a chance to look at the ack tuples issue. This is going to be great perf improvement would like to see this merged in.

@mjsax

Copy link
Copy Markdown
MemberAuthor

I did not have time yet. However, I am still on it... I did already some more work for this PR that is not pushed yet (for example to support different output batches for different output streams and some "clean up" of my changes). I think I have an overall cleaner design now of my code. As a next step, I actually want to investigate why this PR only gives factor two throughput improvement (because Aeolus give factor five). But if you prefer, I can also work on the ack problem first to get a fully running PR for now. (I just disabled acking in my current tests ;))

@harshach

Copy link
Copy Markdown
Contributor

@mjsax Thanks for the quick reply. It would be great to have ack issue fixed than others can also run some tests.

@mjsax

Copy link
Copy Markdown
MemberAuthor

Ok. Hope to get it done over the weekend...

@mjsax

Copy link
Copy Markdown
MemberAuthor

I just realized, that some commits from other people got added to this PR. This confuses me. Can you guide me through the process Storm development is following here? I am not used to that...

@revans2

Copy link
Copy Markdown
Contributor

@mjsax the issue is with how you upmerge. If you just do a git merge or a git pull github can become confused because it thinks you are still based off of the original commit, and will include the commits from the upmerge.

Alternatively you could do a git rebase and rebase all of your changes on a new version, instead of merging in the new version. This will clean it up, but before you do this please make sure you have a backup of your changes, and it does a destructive write that can, and I have had this personally happen to me, delete all of your code. When you are pushing to the github repo you will have to include a -f because git by default does not like destructive writes and you have to force it to do them.

@mjsax

Copy link
Copy Markdown
MemberAuthor

I see. It's a github issue... Usually I rebase before updating a PR. This time I did not... Thanks for the quick response.

@mjsax

Copy link
Copy Markdown
MemberAuthor

I just pushed a new version (rebased to current master and squahed commits etc). The PR is in much cleaner state now. Looking forward to your feedback.

Batching does now work with acks and each output stream can have a different batch size. The ack streams (__ack_init, __ack_ack, and __ack_fail) are treated as regular streams from a batching point of view. Metric and Eventlogger should work too (but I did not test is).

If you want to test this, be aware that hybrid batching is not supported yet (for performance reasons -- at least for the non-acking case). Thus, right now you set batch size via and int only if acking is disabled or via HashMap and the map must contain an entry for each output stream (including ack streams) and the batch size must be at least 1. (After we decided how to proceed -- see next paragraph --, this can be cleaned up.)

The problem with hybrid batching is the serialization path in executor.clj at mk-transfer-fn and start-batch-transfer->worker-handler!. I wrote the code for hybrid serialization already but disabled it, ie, put it into comments. Because, I am not able to set different serializers for different output stream, only a hybrid serializer could be used. However, the runtime binding to the correct method for TupleImpl or Batch reduced the throughput (see numbers below). Not sure if/how this could be resolved. On the other hand, a batch size of one does not have a big performance penalty -- maybe it would be worth to enable batching all the time (ie, even for the non-batching case, just use a batch size of one) to avoid the hybrid setup.

I also did some network monitoring using nmon. It shows that batching does reduce the actually transfered number of byte over the network. The perf tool does not measure but compute the number of transfered bytes (what is not quite accurate). Right now, I don't have the numbers, but if you wish I could rerun those experiments and post here.

I collected the following number for non-acking and acking:

NO ACKING

--name test -l 1 -n 1 --messageSize 100 --workers 24 --spout 1 --bolt 10 --testTimeSec 40

no batching:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443539268669 0 0 0.0
WAITING 1 96 11 11 11 1443539298669 30000 6688960 21.263631184895832
WAITING 1 96 11 11 11 1443539328669 30000 7518460 23.900540669759113
RUNNING 1 96 11 11 11 1443539358669 30000 10428980 33.15283457438151
RUNNING 1 96 11 11 11 1443539388669 30000 10395200 33.045450846354164

batch size = 1 (to measure overhead; about 5%):

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443539471193 0 0 0.0
WAITING 1 96 11 11 11 1443539501193 30000 3089120 9.820048014322916
WAITING 1 96 11 11 11 1443539531193 30000 9134740 29.038556416829426
RUNNING 1 96 11 11 11 1443539561193 30000 9502680 30.208206176757812
RUNNING 1 96 11 11 11 1443539591193 30000 9672300 30.747413635253906

batch size = 100 (throughput improvement by about 85%)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443539658994 0 0 0.0
WAITING 1 96 11 11 11 1443539688994 30000 8345560 26.529820760091145
WAITING 1 96 11 11 11 1443539718994 30000 19876460 63.18556467692057
RUNNING 1 96 11 11 11 1443539748994 30000 18229880 57.95122782389323
RUNNING 1 96 11 11 11 1443539778994 30000 18294660 58.15715789794922

ACKING

--name test -l 1 -n 1 --messageSize 100 --workers 24 --spout 1 --bolt 10 --testTimeSec 40 --ack --ackers 10

no batching:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 21 0 1443539868024 0 0 0.0
WAITING 1 96 21 21 21 1443539898024 30000 864800 2.7491251627604165
WAITING 1 96 21 21 21 1443539928024 30000 1768760 5.6227366129557295
RUNNING 1 96 21 21 21 1443539958024 30000 1910340 6.072807312011719
RUNNING 1 96 21 21 21 1443539988025 30001 1888740 6.003942629809475

Complete Latency (from WebUI): 6.256

all batch sizes = 1 (to measure overhead; acking dominates; no overhead measurable):

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 2 21 0 1443540043224 0 0 0.0
WAITING 1 96 21 21 21 1443540073225 30001 803060 2.552773895980811
WAITING 1 96 21 21 21 1443540103225 30000 2001520 6.362660725911458
RUNNING 1 96 21 21 21 1443540133224 29999 1789860 5.690001373255411
RUNNING 1 96 21 21 21 1443540163224 30000 1925420 6.120745340983073

Complete Latency (from WebUI): 9.686 (no impact -- 3ms difference is just too small to be a reliable number)

default batch size = 100 (almost no throughput improvement; about 15%; acking dominates as acks are not batched)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 21 0 1443540446423 0 0 0.0
WAITING 1 96 21 21 21 1443540476423 30000 866980 2.7560551961263022
WAITING 1 96 21 21 21 1443540506423 30000 2076100 6.599744160970052
RUNNING 1 96 21 21 21 1443540536424 30001 2100200 6.676133459939356
RUNNING 1 96 21 21 21 1443540566424 30000 2191620 6.966972351074219

Complete Latency (from WebUI): 11.721 (compared to 6ms, almost doubles, but this is still tiny)

all batch sizes = 100 (acks are batched too; throughput improvement increases to more than 50%)

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 21 0 1443540631845 0 0 0.0
WAITING 1 96 21 21 21 1443540661846 30001 1728100 5.493298843977336
WAITING 1 96 21 21 21 1443540691845 29999 2814100 8.946081182035496
RUNNING 1 96 21 21 21 1443540721845 30000 2084440 6.6262563069661455
RUNNING 1 96 21 21 21 1443540751845 30000 2949700 9.376843770345053

Complete Latency (from WebUI): 65.225 (latency increased by factor of 10 due to ack batching)

===== Hybrid =====
I repeated the same experiment using hybrid serializer (KryoTupleBatchSerializer): It is a similar, however with reduced throughput by 25% for the non-acking case. In the acking case, the hybrid approach has no influence as acking dominates.

NO ACKING (hybrid)

--name test -l 1 -n 1 --messageSize 100 --workers 24 --spout 1 --bolt 10 --testTimeSec 40

no batching:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443536488016 0 0 0.0
WAITING 1 96 11 11 11 1443536518016 30000 5118760 16.27209981282552
WAITING 1 96 11 11 11 1443536548017 30001 8131680 25.84905291568406
RUNNING 1 96 11 11 11 1443536578016 29999 7960300 25.305955734820067
RUNNING 1 96 11 11 11 1443536608016 30000 7711800 24.515151977539062

batch size = 1 (to measure overhead):

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443536668726 0 0 0.0
WAITING 1 96 11 11 11 1443536698726 30000 4799100 15.255928039550781
WAITING 1 96 11 11 11 1443536728726 30000 7695580 24.463589986165363
RUNNING 1 96 11 11 11 1443536758727 30001 7828800 24.886255419090197
RUNNING 1 96 11 11 11 1443536788726 29999 5613780 17.84632089054661

batch size = 100:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 11 0 1443536914666 0 0 0.0
WAITING 1 96 11 11 11 1443536944666 30000 9114920 28.975550333658855
WAITING 1 96 11 11 11 1443536974666 30000 16459180 52.32232411702474
RUNNING 1 96 11 11 11 1443537004666 30000 16347100 51.96603139241537
RUNNING 1 96 11 11 11 1443537034666 30000 16594400 52.752176920572914

ACKING (hybrid)

--name test -l 1 -n 1 --messageSize 100 --workers 24 --spout 1 --bolt 10 --testTimeSec 40 --ack --ackers 10

no batching:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 21 0 1443538020486 0 0 0.0
WAITING 1 96 21 21 21 1443538050486 30000 934640 2.9711405436197915
WAITING 1 96 21 21 21 1443538080486 30000 1774140 5.639839172363281
RUNNING 1 96 21 21 21 1443538110486 30000 1875420 5.961799621582031
RUNNING 1 96 21 21 21 1443538140486 30000 1832160 5.82427978515625

Complete Latency (from WebUI): 43.945

all batch sizes = 1 (to measure overhead):

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 3 21 0 1443538188807 0 0 0.0
WAITING 1 96 21 21 21 1443538218807 30000 1192380 3.7904739379882812
WAITING 1 96 21 21 21 1443538248807 30000 1915040 6.087748209635417
RUNNING 1 96 21 21 21 1443538278808 30001 1387960 4.412058945365883
RUNNING 1 96 21 21 21 1443538308807 29999 2106580 6.696860700206934

Complete Latency (from WebUI): 9.132

default batch size = 100:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 5 21 0 1443538350612 0 0 0.0
WAITING 1 96 21 21 21 1443538380612 30000 1449080 4.6065012613932295
WAITING 1 96 21 21 21 1443538410613 30001 2369020 7.5306607414843985
RUNNING 1 96 21 21 21 1443538440612 29999 2323580 7.386708117321359
RUNNING 1 96 21 21 21 1443538470612 30000 2414040 7.6740264892578125

Complete Latency (from WebUI): 11.063

all batch sizes = 100:

status topologies totalSlots slotsUsed totalExecutors executorsWithMetrics time time-diff ms transferred throughput (MB/s)
WAITING 1 96 0 21 0 1443538562562 0 0 0.0
WAITING 1 96 21 21 21 1443538592562 30000 1278060 4.062843322753906
WAITING 1 96 21 21 21 1443538622562 30000 2170020 6.898307800292969
RUNNING 1 96 21 21 21 1443538652563 30001 2209060 7.022178545383122
RUNNING 1 96 21 21 21 1443538682562 29999 2155180 6.851361089477723

Complete Latency (from WebUI): 122.958

mjsax added 3 commits September 30, 2015 00:37
 - batch_size (int): set the same batch size for all UDF declared output streams
- batch_sizes (Map<String,Number>): sets the batch sizer for individual streams (including system streams)
 - extended Kryo (de)serialization for Batch
- integrated Tuple and Batch (de)serialization
 - added internal tuple buffers (one for each output stream)
- added Spout/Bolt output batching
- added Bolt input debatching
- added batching of "acks"
@mjsax

Copy link
Copy Markdown
MemberAuthor

One question: There are a lot of changes in storm-core/src/jvm/backtype/storm/generated/* resulting from rebuild those files with genthrift.sh. However, it seems to me that only the changes to ComponentCommon.java (and I guess to storm-core/src/py/storm/ttypes.py) are relevant. For all other classes it this package it seems they include variable renaming only. The files in which only the generation date was changed are not included already. Can I safely revert those other files, too? Or might I break something?

@mjsax

mjsax commented Oct 5, 2015

Copy link
Copy Markdown
MemberAuthor

Hi, what is the next step?

@revans2

Copy link
Copy Markdown
Contributor

I still need to run some tests. I am way behind on my open source commitments. I will try really hard this week to play around with this and let you know the results.

@mjsax

mjsax commented Oct 5, 2015

Copy link
Copy Markdown
MemberAuthor

Thx. Take your time; there is actually no rush. I was just curious :)

@mjsax

Copy link
Copy Markdown
MemberAuthor

Closing this (also closing the JIRA) as Bobby's work (#765) got merged and I don't have time to work on this right now.

@mjsaxmjsax closed this Mar 31, 2016
Sign up for freeto join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants

@mjsax@HeartSaVioR@knusbaum@revans2@harshach