From 1706746a325b1f9371de0669f2e34ef443a1da2f Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Fri, 22 Jun 2018 09:54:24 +0000 Subject: [PATCH 1/3] Faster report merging through mutating objects When we know we have the only reference to a Report or Node object we can avoid copying the data to change it. Add "Unsafe" variants of various Merge operations which mutate the receiver, and a new Merger which takes advantage of them. --- app/benchmark_internal_test.go | 9 +++++++++ app/merger.go | 18 ++++++++++++++++++ app/merger_test.go | 4 ++++ report/report.go | 19 ++++++++++++------- report/topology.go | 28 ++++++++++++++++++++++++---- 5 files changed, 67 insertions(+), 11 deletions(-) diff --git a/app/benchmark_internal_test.go b/app/benchmark_internal_test.go index afc61e63f..fa5944654 100644 --- a/app/benchmark_internal_test.go +++ b/app/benchmark_internal_test.go @@ -72,6 +72,15 @@ func BenchmarkReportMerge(b *testing.B) { } } +func BenchmarkReportFastMerge(b *testing.B) { + reports := upgradeReports(readReportFiles(b, *benchReportPath)) + merger := NewFastMerger() + b.ResetTimer() + for i := 0; i < b.N; i++ { + merger.Merge(reports) + } +} + func getReport(b *testing.B) report.Report { r := fixture.Report if *benchReportPath != "" { diff --git a/app/merger.go b/app/merger.go index 94231baf8..ac5c936f3 100644 --- a/app/merger.go +++ b/app/merger.go @@ -60,3 +60,21 @@ func (smartMerger) Merge(reports []report.Report) report.Report { } return <-c } + +type fastMerger struct{} + +// NewFastMerger makes a Merger which merges together reports, mutating the one we are building up +func NewFastMerger() Merger { + return fastMerger{} +} + +func (fastMerger) Merge(reports []report.Report) report.Report { + rpt := report.MakeReport() + id := murmur3.New64() + for _, r := range reports { + rpt.UnsafeMerge(r) + id.Write([]byte(r.ID)) + } + rpt.ID = fmt.Sprintf("%x", id.Sum64()) + return rpt +} diff --git a/app/merger_test.go b/app/merger_test.go index 5ec36c307..ec114eedb 100644 --- a/app/merger_test.go +++ b/app/merger_test.go @@ -52,6 +52,10 @@ func BenchmarkDumbMerger(b *testing.B) { benchmarkMerger(b, app.MakeDumbMerger()) } +func BenchmarkFastMerger(b *testing.B) { + benchmarkMerger(b, app.NewFastMerger()) +} + const numHosts = 15 func benchmarkMerger(b *testing.B, merger app.Merger) { diff --git a/report/report.go b/report/report.go index a217550e9..d21551c50 100644 --- a/report/report.go +++ b/report/report.go @@ -305,16 +305,21 @@ func (r Report) Copy() Report { // original is not modified. func (r Report) Merge(other Report) Report { newReport := r.Copy() - newReport.DNS = newReport.DNS.Merge(other.DNS) - newReport.Sampling = newReport.Sampling.Merge(other.Sampling) - newReport.Window = newReport.Window + other.Window - newReport.Plugins = newReport.Plugins.Merge(other.Plugins) - newReport.WalkPairedTopologies(&other, func(ourTopology, theirTopology *Topology) { - *ourTopology = ourTopology.Merge(*theirTopology) - }) + newReport.UnsafeMerge(other) return newReport } +// UnsafeMerge merges another Report into the receiver. The original is modified. +func (r *Report) UnsafeMerge(other Report) { + r.DNS = r.DNS.Merge(other.DNS) + r.Sampling = r.Sampling.Merge(other.Sampling) + r.Window = r.Window + other.Window + r.Plugins = r.Plugins.Merge(other.Plugins) + r.WalkPairedTopologies(&other, func(ourTopology, theirTopology *Topology) { + ourTopology.UnsafeMerge(*theirTopology) + }) +} + // WalkTopologies iterates through the Topologies of the report, // potentially modifying them func (r *Report) WalkTopologies(f func(*Topology)) { diff --git a/report/topology.go b/report/topology.go index 83baa0f94..49c01678c 100644 --- a/report/topology.go +++ b/report/topology.go @@ -163,6 +163,21 @@ func (t Topology) Merge(other Topology) Topology { } } +// UnsafeMerge merges the other object into this one, modifying the original. +func (t *Topology) UnsafeMerge(other Topology) { + if t.Shape == "" { + t.Shape = other.Shape + } + if t.Label == "" { + t.Label, t.LabelPlural = other.Label, other.LabelPlural + } + t.Nodes.UnsafeMerge(other.Nodes) + t.Controls = t.Controls.Merge(other.Controls) + t.MetadataTemplates = t.MetadataTemplates.Merge(other.MetadataTemplates) + t.MetricTemplates = t.MetricTemplates.Merge(other.MetricTemplates) + t.TableTemplates = t.TableTemplates.Merge(other.TableTemplates) +} + // Nodes is a collection of nodes in a topology. Keys are node IDs. // TODO(pb): type Topology map[string]Node type Nodes map[string]Node @@ -186,14 +201,19 @@ func (n Nodes) Merge(other Nodes) Nodes { return n } cp := n.Copy() + cp.UnsafeMerge(other) + return cp +} + +// UnsafeMerge merges the other object into this one, modifying the original. +func (n *Nodes) UnsafeMerge(other Nodes) { for k, v := range other { - if n, ok := cp[k]; ok { // don't overwrite - cp[k] = v.Merge(n) + if existing, ok := (*n)[k]; ok { // don't overwrite + (*n)[k] = v.Merge(existing) } else { - cp[k] = v + (*n)[k] = v } } - return cp } // Validate checks the topology for various inconsistencies. From 126a171f62f93b01396f8abb93e12c0b50140cfd Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Fri, 22 Jun 2018 10:03:14 +0000 Subject: [PATCH 2/3] Make 'fast' merger the default --- app/benchmark_internal_test.go | 9 --------- app/collector.go | 4 ++-- app/multitenant/aws_collector.go | 2 +- 3 files changed, 3 insertions(+), 12 deletions(-) diff --git a/app/benchmark_internal_test.go b/app/benchmark_internal_test.go index fa5944654..ec7814e5a 100644 --- a/app/benchmark_internal_test.go +++ b/app/benchmark_internal_test.go @@ -64,15 +64,6 @@ func BenchmarkReportUpgrade(b *testing.B) { } func BenchmarkReportMerge(b *testing.B) { - reports := upgradeReports(readReportFiles(b, *benchReportPath)) - merger := NewSmartMerger() - b.ResetTimer() - for i := 0; i < b.N; i++ { - merger.Merge(reports) - } -} - -func BenchmarkReportFastMerge(b *testing.B) { reports := upgradeReports(readReportFiles(b, *benchReportPath)) merger := NewFastMerger() b.ResetTimer() diff --git a/app/collector.go b/app/collector.go index 93e71ab2b..e0906b11d 100644 --- a/app/collector.go +++ b/app/collector.go @@ -107,7 +107,7 @@ func NewCollector(window time.Duration) Collector { waitableCondition: waitableCondition{ waiters: map[chan struct{}]struct{}{}, }, - merger: NewSmartMerger(), + merger: NewFastMerger(), } } @@ -292,7 +292,7 @@ func NewFileCollector(path string, window time.Duration) (Collector, error) { go replay(collector, timestamps, reports) return collector, nil } - return StaticCollector(NewSmartMerger().Merge(reports).Upgrade()), nil + return StaticCollector(NewFastMerger().Merge(reports).Upgrade()), nil } func timestampFromFilepath(path string) (time.Time, error) { diff --git a/app/multitenant/aws_collector.go b/app/multitenant/aws_collector.go index d0a75dbf3..e9976d958 100644 --- a/app/multitenant/aws_collector.go +++ b/app/multitenant/aws_collector.go @@ -154,7 +154,7 @@ func NewAWSCollector(config AWSCollectorConfig) (AWSCollector, error) { s3: config.S3Store, userIDer: config.UserIDer, tableName: config.DynamoTable, - merger: app.NewSmartMerger(), + merger: app.NewFastMerger(), inProcess: newInProcessStore(reportCacheSize, config.Window), memcache: config.MemcacheClient, window: config.Window, From 3309d09ad898baf89fcf3ea657465222f9a45024 Mon Sep 17 00:00:00 2001 From: Bryan Boreham Date: Fri, 22 Jun 2018 10:47:00 +0000 Subject: [PATCH 3/3] Remove slower mergers --- app/benchmark_internal_test.go | 2 +- app/merger.go | 48 ---------------------------------- app/merger_test.go | 10 +------ 3 files changed, 2 insertions(+), 58 deletions(-) diff --git a/app/benchmark_internal_test.go b/app/benchmark_internal_test.go index ec7814e5a..1c9e06436 100644 --- a/app/benchmark_internal_test.go +++ b/app/benchmark_internal_test.go @@ -75,7 +75,7 @@ func BenchmarkReportMerge(b *testing.B) { func getReport(b *testing.B) report.Report { r := fixture.Report if *benchReportPath != "" { - r = NewSmartMerger().Merge(upgradeReports(readReportFiles(b, *benchReportPath))) + r = NewFastMerger().Merge(upgradeReports(readReportFiles(b, *benchReportPath))) } return r } diff --git a/app/merger.go b/app/merger.go index ac5c936f3..17397d16a 100644 --- a/app/merger.go +++ b/app/merger.go @@ -13,54 +13,6 @@ type Merger interface { Merge([]report.Report) report.Report } -type dumbMerger struct{} - -// MakeDumbMerger makes a Merger which merges together reports in the simplest possible way. -func MakeDumbMerger() Merger { - return dumbMerger{} -} - -func (dumbMerger) Merge(reports []report.Report) report.Report { - rpt := report.MakeReport() - id := murmur3.New64() - for _, r := range reports { - rpt = rpt.Merge(r) - id.Write([]byte(r.ID)) - } - rpt.ID = fmt.Sprintf("%x", id.Sum64()) - return rpt -} - -type smartMerger struct{} - -// NewSmartMerger makes a Merger which merges reports in -// parallel. Speed up comes from the fact that a) most merges are -// between small reports, and b) we take advantage of available cores. -func NewSmartMerger() Merger { - return smartMerger{} -} - -func (smartMerger) Merge(reports []report.Report) report.Report { - l := len(reports) - switch l { - case 0: - return report.MakeReport() - case 1: - return reports[0] - } - c := make(chan report.Report, l) - for _, r := range reports { - c <- r - } - for ; l > 1; l-- { - left, right := <-c, <-c - go func() { - c <- left.Merge(right) - }() - } - return <-c -} - type fastMerger struct{} // NewFastMerger makes a Merger which merges together reports, mutating the one we are building up diff --git a/app/merger_test.go b/app/merger_test.go index ec114eedb..dae040e95 100644 --- a/app/merger_test.go +++ b/app/merger_test.go @@ -27,7 +27,7 @@ func TestMerger(t *testing.T) { want.Endpoint.AddNode(report.MakeNode("bar")) want.Endpoint.AddNode(report.MakeNode("baz")) - for _, merger := range []app.Merger{app.MakeDumbMerger(), app.NewSmartMerger()} { + for _, merger := range []app.Merger{app.NewFastMerger()} { // Test the empty list case if have := merger.Merge([]report.Report{}); !reflect.DeepEqual(have, report.MakeReport()) { t.Errorf("Bad merge: %s", test.Diff(have, want)) @@ -44,14 +44,6 @@ func TestMerger(t *testing.T) { } } -func BenchmarkSmartMerger(b *testing.B) { - benchmarkMerger(b, app.NewSmartMerger()) -} - -func BenchmarkDumbMerger(b *testing.B) { - benchmarkMerger(b, app.MakeDumbMerger()) -} - func BenchmarkFastMerger(b *testing.B) { benchmarkMerger(b, app.NewFastMerger()) }