Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
65 changes: 37 additions & 28 deletions asap-precompute-go/window.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ import (
// owns one Sketch instance and the labels needed to reconstruct
// the SketchEnvelope at flush time.
type seriesEntry struct {
mu sync.Mutex
seriesKey string
heapIndex int
// Sketch is the running sketch for this series. Owned here;
Expand Down Expand Up @@ -120,6 +121,7 @@ type windowState struct {
mu sync.RWMutex
series map[string]*seriesEntry
eviction seriesEvictionHeap
evictionMu sync.Mutex
snapshotCache *SnapshotCache
sketchSink *atomic.Pointer[SketchSink]
activeStartMs uint64
Expand Down Expand Up @@ -303,18 +305,6 @@ func (w *windowState) observe(
observer SketchObserver,
stats *PrecomputeStats,
) error {
w.mu.Lock()
defer w.mu.Unlock()

if w.sketchFactory == nil {
w.sketchFactory = sketchFactory
}
w.initWindow(obs.TimestampMs, cfg)

if err := w.validateTimestampLocked(obs.TimestampMs, cfg); err != nil {
return err
}

// Build the lookup key into a pooled byte buffer so the common
// case (an already-admitted series) costs no allocation: the
// `w.series[string(sc.buf)]` index is the compiler's zero-alloc
Expand All @@ -323,14 +313,7 @@ func (w *windowState) observe(
sc := getSeriesKeyScratch()
defer putSeriesKeyScratch(sc)
cfg.buildSeriesKey(sc, obs)
entry, ok := w.series[string(sc.buf)]
if !ok {
var err error
if entry, err = w.admitSeriesLocked(string(sc.buf), obs, cfg, sketchFactory, stats); err != nil {
return err
}
}
return w.recordLocked(entry, obs, observer)
return w.observeWithKey(string(sc.buf), obs, cfg, sketchFactory, observer, stats)
}

// observeKeyed is the shared-key entry point for the fused asap_edge
Expand All @@ -346,26 +329,48 @@ func (w *windowState) observeKeyed(
observer SketchObserver,
stats *PrecomputeStats,
) error {
return w.observeWithKey(key, obs, cfg, sketchFactory, observer, stats)
}

func (w *windowState) observeWithKey(key string, obs *Observation, cfg *PrecomputeConfig, sketchFactory SketchFactory, observer SketchObserver, stats *PrecomputeStats) error {
// The generation read lock prevents rotation while independent entries use
// their own locks. Established series therefore update concurrently.
w.mu.RLock()
if w.initialized {
if err := w.validateTimestampLocked(obs.TimestampMs, cfg); err != nil {
w.mu.RUnlock()
return err
}
if entry := w.series[key]; entry != nil {
entry.mu.Lock()
err := w.recordLocked(entry, obs, observer, cfg.MaxSeries > 0 && cfg.OnOverflow == OnOverflowEvictOldest)
entry.mu.Unlock()
w.mu.RUnlock()
return err
}
}
w.mu.RUnlock()

// New-series admission mutates the map/cardinality index and is rare after
// warm-up. Recheck after upgrading because another goroutine may admit it.
w.mu.Lock()
defer w.mu.Unlock()

if w.sketchFactory == nil {
w.sketchFactory = sketchFactory
}
w.initWindow(obs.TimestampMs, cfg)

if err := w.validateTimestampLocked(obs.TimestampMs, cfg); err != nil {
return err
}

entry, ok := w.series[key]
if !ok {
entry := w.series[key]
if entry == nil {
var err error
if entry, err = w.admitSeriesLocked(key, obs, cfg, sketchFactory, stats); err != nil {
entry, err = w.admitSeriesLocked(key, obs, cfg, sketchFactory, stats)
if err != nil {
return err
}
}
return w.recordLocked(entry, obs, observer)
return w.recordLocked(entry, obs, observer, cfg.MaxSeries > 0 && cfg.OnOverflow == OnOverflowEvictOldest)
}

// validateTimestampLocked prevents host scheduling jitter from changing
Expand Down Expand Up @@ -484,8 +489,12 @@ func (w *windowState) evictOldestLocked() (string, *seriesEntry) {

// recordLocked feeds one observation into a series' sketch and advances its
// bookkeeping. Caller holds w.mu. Shared by observe and observeKeyed.
func (w *windowState) recordLocked(entry *seriesEntry, obs *Observation, observer SketchObserver) error {
func (w *windowState) recordLocked(entry *seriesEntry, obs *Observation, observer SketchObserver, trackEviction bool) error {
if obs.TimestampMs > entry.LastSeenMs {
if trackEviction {
w.evictionMu.Lock()
defer w.evictionMu.Unlock()
}
entry.LastSeenMs = obs.TimestampMs
if entry.heapIndex >= 0 {
heap.Fix(&w.eviction, entry.heapIndex)
Expand Down
37 changes: 37 additions & 0 deletions asap-precompute-go/window_parallel_bench_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,37 @@
package precompute

import (
"strconv"
"sync/atomic"
"testing"
"time"
)

func BenchmarkParallelExistingSeries(b *testing.B) {
const series = 1024
cfg := &PrecomputeConfig{AggID: 1, SketchType: SketchTypeDDSketch,
Mode: Tumbling, Window: WindowSpec{Size: time.Hour}}
w := newWindowState()
factory, observer, stats := newFakeFactory(), &fakeObserver{}, NewStats()
observations := make([]*Observation, series)
keys := make([]string, series)
for i := range observations {
observations[i] = &Observation{TimestampMs: 1, Labels: []KeyValue{{Key: "k", Value: strconv.Itoa(i)}}, Value: FloatValue(1)}
keys[i] = cfg.SeriesKeyFor(observations[i])
if err := w.observeKeyed(keys[i], observations[i], cfg, factory, observer, stats); err != nil {
b.Fatal(err)
}
}
var next atomic.Uint64
b.ReportAllocs()
b.ResetTimer()
b.RunParallel(func(pb *testing.PB) {
for pb.Next() {
i := int(next.Add(1)) & (series - 1)
if err := w.observeKeyed(keys[i], observations[i], cfg, factory, observer, stats); err != nil {
b.Error(err)
return
}
}
})
}
Loading