From 1d03fd5cf6e22b9c2231aac6206c580d93c721e1 Mon Sep 17 00:00:00 2001 From: songlq <31207055+songlonqi-java@users.noreply.github.com> Date: Mon, 13 Jul 2026 15:04:21 +0800 Subject: [PATCH 1/2] tail-sampling --- aggregate/tail_sampling_processor.go | 34 +++++++++++++++++----------- 1 file changed, 21 insertions(+), 13 deletions(-) diff --git a/aggregate/tail_sampling_processor.go b/aggregate/tail_sampling_processor.go index 72a52fcb..ef366e65 100644 --- a/aggregate/tail_sampling_processor.go +++ b/aggregate/tail_sampling_processor.go @@ -117,6 +117,24 @@ func (r *TailSamplingProcessor) AdvanceTime() map[uint64]*DataGroup { } func (r *TailSamplingProcessor) TailSamplingData(dataGroups map[uint64]*DataGroup) map[uint64]*DataPacket { + outcomes := r.TailSamplingOutcomes(dataGroups) + if outcomes == nil { + return nil + } + + keptPackets := make(map[uint64]*DataPacket) + + for key, outcome := range outcomes { + if outcome == nil || outcome.Packet == nil { + continue + } + keptPackets[key] = outcome.Packet + } + + return keptPackets +} + +func (r *TailSamplingProcessor) TailSamplingOutcomes(dataGroups map[uint64]*DataGroup) map[uint64]*TailSamplingOutcome { if r == nil || r.sampler == nil { return nil } @@ -131,19 +149,12 @@ func (r *TailSamplingProcessor) TailSamplingData(dataGroups map[uint64]*DataGrou } outcomes := r.sampler.TailSamplingOutcomes(dataGroups) - keptPackets := make(map[uint64]*DataPacket) if r.collector == nil || len(r.metrics) == 0 { - for key, outcome := range outcomes { - if outcome == nil || outcome.Packet == nil { - continue - } - keptPackets[key] = outcome.Packet - } - return keptPackets + return outcomes } - for key, outcome := range outcomes { + for _, outcome := range outcomes { if outcome == nil { continue } @@ -153,12 +164,9 @@ func (r *TailSamplingProcessor) TailSamplingData(dataGroups map[uint64]*DataGrou r.collector.Add(r.filterBuiltinRecords(packetForMetrics, r.metrics.OnDecision(packetForMetrics, outcome.Decision))) } - if outcome.Packet != nil { - keptPackets[key] = outcome.Packet - } } - return keptPackets + return outcomes } func (r *TailSamplingProcessor) RecordDecision(packet *DataPacket, decision DerivedMetricDecision) { From 1312f0f7ed603b47d46c7219d0cc7d90ed381f2e Mon Sep 17 00:00:00 2001 From: songlq <31207055+songlonqi-java@users.noreply.github.com> Date: Mon, 13 Jul 2026 15:34:48 +0800 Subject: [PATCH 2/2] tail-sampling --- aggregate/tail_sampling_processor.go | 1 - 1 file changed, 1 deletion(-) diff --git a/aggregate/tail_sampling_processor.go b/aggregate/tail_sampling_processor.go index ef366e65..49569e5e 100644 --- a/aggregate/tail_sampling_processor.go +++ b/aggregate/tail_sampling_processor.go @@ -163,7 +163,6 @@ func (r *TailSamplingProcessor) TailSamplingOutcomes(dataGroups map[uint64]*Data if packetForMetrics != nil { r.collector.Add(r.filterBuiltinRecords(packetForMetrics, r.metrics.OnDecision(packetForMetrics, outcome.Decision))) } - } return outcomes