mirror of
https://github.com/prometheus/prometheus.git
synced 2024-11-09 23:24:05 -08:00
Add ref counting to string interning so we can remove
a string when there are no longer any refs. Add tests for interning. Co-authored-by: Tom Wilkie <tom.wilkie@gmail.com> Signed-off-by: Callum Styan <callumstyan@gmail.com>
This commit is contained in:
parent
cbf5f13285
commit
1a7923dde3
|
@ -18,18 +18,26 @@
|
||||||
|
|
||||||
package remote
|
package remote
|
||||||
|
|
||||||
import "sync"
|
import (
|
||||||
|
"sync"
|
||||||
|
"sync/atomic"
|
||||||
|
)
|
||||||
|
|
||||||
var interner = newPool()
|
var interner = newPool()
|
||||||
|
|
||||||
type pool struct {
|
type pool struct {
|
||||||
mtx sync.RWMutex
|
mtx sync.RWMutex
|
||||||
pool map[string]string
|
pool map[string]*entry
|
||||||
|
}
|
||||||
|
|
||||||
|
type entry struct {
|
||||||
|
s string
|
||||||
|
refs int64
|
||||||
}
|
}
|
||||||
|
|
||||||
func newPool() *pool {
|
func newPool() *pool {
|
||||||
return &pool{
|
return &pool{
|
||||||
pool: map[string]string{},
|
pool: map[string]*entry{},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -42,18 +50,44 @@ func (p *pool) intern(s string) string {
|
||||||
interned, ok := p.pool[s]
|
interned, ok := p.pool[s]
|
||||||
p.mtx.RUnlock()
|
p.mtx.RUnlock()
|
||||||
if ok {
|
if ok {
|
||||||
return interned
|
atomic.AddInt64(&interned.refs, 1)
|
||||||
|
return interned.s
|
||||||
|
}
|
||||||
|
p.mtx.Lock()
|
||||||
|
defer p.mtx.Unlock()
|
||||||
|
if interned, ok := p.pool[s]; ok {
|
||||||
|
atomic.AddInt64(&interned.refs, 1)
|
||||||
|
return interned.s
|
||||||
|
}
|
||||||
|
|
||||||
|
s = pack(s)
|
||||||
|
p.pool[s] = &entry{
|
||||||
|
s: s,
|
||||||
|
refs: 1,
|
||||||
|
}
|
||||||
|
return s
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *pool) release(s string) {
|
||||||
|
p.mtx.RLock()
|
||||||
|
interned, ok := p.pool[s]
|
||||||
|
p.mtx.RUnlock()
|
||||||
|
|
||||||
|
if !ok {
|
||||||
|
panic("released unknown string")
|
||||||
|
}
|
||||||
|
|
||||||
|
refs := atomic.AddInt64(&interned.refs, -1)
|
||||||
|
if refs > 0 {
|
||||||
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
p.mtx.Lock()
|
p.mtx.Lock()
|
||||||
defer p.mtx.Unlock()
|
defer p.mtx.Unlock()
|
||||||
if interned, ok := p.pool[s]; ok {
|
if atomic.LoadInt64(&interned.refs) != 0 {
|
||||||
return interned
|
return
|
||||||
}
|
}
|
||||||
|
delete(p.pool, s)
|
||||||
s = pack(s)
|
|
||||||
p.pool[s] = s
|
|
||||||
return s
|
|
||||||
}
|
}
|
||||||
|
|
||||||
// StrPack returns a new instance of s which is tightly packed in memory.
|
// StrPack returns a new instance of s which is tightly packed in memory.
|
||||||
|
|
85
storage/remote/intern_test.go
Normal file
85
storage/remote/intern_test.go
Normal file
|
@ -0,0 +1,85 @@
|
||||||
|
// Copyright 2019 The Prometheus Authors
|
||||||
|
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||||
|
// you may not use this file except in compliance with the License.
|
||||||
|
// You may obtain a copy of the License at
|
||||||
|
//
|
||||||
|
// http://www.apache.org/licenses/LICENSE-2.0
|
||||||
|
//
|
||||||
|
// Unless required by applicable law or agreed to in writing, software
|
||||||
|
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||||
|
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||||
|
// See the License for the specific language governing permissions and
|
||||||
|
// limitations under the License.
|
||||||
|
//
|
||||||
|
// Inspired / copied / modified from https://gitlab.com/cznic/strutil/blob/master/strutil.go,
|
||||||
|
// which is MIT licensed, so:
|
||||||
|
//
|
||||||
|
// Copyright (c) 2014 The strutil Authors. All rights reserved.
|
||||||
|
|
||||||
|
package remote
|
||||||
|
|
||||||
|
import (
|
||||||
|
"fmt"
|
||||||
|
"testing"
|
||||||
|
"time"
|
||||||
|
|
||||||
|
"github.com/prometheus/prometheus/util/testutil"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestIntern(t *testing.T) {
|
||||||
|
testString := "TestIntern"
|
||||||
|
interner.intern(testString)
|
||||||
|
interned, ok := interner.pool[testString]
|
||||||
|
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 1, fmt.Sprintf("expected refs to be 1 but it was %d", interned.refs))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIntern_MultiRef(t *testing.T) {
|
||||||
|
testString := "TestIntern_MultiRef"
|
||||||
|
|
||||||
|
interner.intern(testString)
|
||||||
|
interned, ok := interner.pool[testString]
|
||||||
|
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 1, fmt.Sprintf("expected refs to be 1 but it was %d", interned.refs))
|
||||||
|
|
||||||
|
interner.intern(testString)
|
||||||
|
interned, ok = interner.pool[testString]
|
||||||
|
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 2, fmt.Sprintf("expected refs to be 2 but it was %d", interned.refs))
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIntern_DeleteRef(t *testing.T) {
|
||||||
|
testString := "TestIntern_DeleteRef"
|
||||||
|
|
||||||
|
interner.intern(testString)
|
||||||
|
interned, ok := interner.pool[testString]
|
||||||
|
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 1, fmt.Sprintf("expected refs to be 1 but it was %d", interned.refs))
|
||||||
|
|
||||||
|
interner.release(testString)
|
||||||
|
_, ok = interner.pool[testString]
|
||||||
|
testutil.Equals(t, ok, false)
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestIntern_MultiRef_Concurrent(t *testing.T) {
|
||||||
|
testString := "TestIntern_MultiRef_Concurrent"
|
||||||
|
|
||||||
|
interner.intern(testString)
|
||||||
|
interned, ok := interner.pool[testString]
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 1, fmt.Sprintf("expected refs to be 1 but it was %d", interned.refs))
|
||||||
|
|
||||||
|
go interner.release(testString)
|
||||||
|
|
||||||
|
interner.intern(testString)
|
||||||
|
|
||||||
|
time.Sleep(time.Millisecond)
|
||||||
|
|
||||||
|
interned, ok = interner.pool[testString]
|
||||||
|
testutil.Equals(t, ok, true)
|
||||||
|
testutil.Assert(t, interned.refs == 1, fmt.Sprintf("expected refs to be 1 but it was %d", interned.refs))
|
||||||
|
}
|
|
@ -343,8 +343,12 @@ func (t *QueueManager) StoreSeries(series []tsdb.RefSeries, index int) {
|
||||||
t.seriesMtx.Lock()
|
t.seriesMtx.Lock()
|
||||||
defer t.seriesMtx.Unlock()
|
defer t.seriesMtx.Unlock()
|
||||||
for ref, labels := range temp {
|
for ref, labels := range temp {
|
||||||
t.seriesLabels[ref] = labels
|
|
||||||
t.seriesSegmentIndexes[ref] = index
|
t.seriesSegmentIndexes[ref] = index
|
||||||
|
|
||||||
|
if orig, ok := t.seriesLabels[ref]; ok {
|
||||||
|
release(orig)
|
||||||
|
}
|
||||||
|
t.seriesLabels[ref] = labels
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -359,12 +363,20 @@ func (t *QueueManager) SeriesReset(index int) {
|
||||||
// that were not also present in the checkpoint.
|
// that were not also present in the checkpoint.
|
||||||
for k, v := range t.seriesSegmentIndexes {
|
for k, v := range t.seriesSegmentIndexes {
|
||||||
if v < index {
|
if v < index {
|
||||||
delete(t.seriesLabels, k)
|
|
||||||
delete(t.seriesSegmentIndexes, k)
|
delete(t.seriesSegmentIndexes, k)
|
||||||
|
release(t.seriesLabels[k])
|
||||||
|
delete(t.seriesLabels, k)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func release(ls []prompb.Label) {
|
||||||
|
for _, l := range ls {
|
||||||
|
interner.release(l.Name)
|
||||||
|
interner.release(l.Value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// processExternalLabels merges externalLabels into ls. If ls contains
|
// processExternalLabels merges externalLabels into ls. If ls contains
|
||||||
// a label in externalLabels, the value in ls wins.
|
// a label in externalLabels, the value in ls wins.
|
||||||
func processExternalLabels(ls tsdbLabels.Labels, externalLabels labels.Labels) labels.Labels {
|
func processExternalLabels(ls tsdbLabels.Labels, externalLabels labels.Labels) labels.Labels {
|
||||||
|
|
Loading…
Reference in a new issue