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
2 changes: 1 addition & 1 deletion internal/component/common/loki/client/consumer_wal.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ func NewWALConsumer(logger *slog.Logger, reg prometheus.Registerer, walCfg wal.C
return nil, fmt.Errorf("at least one endpoint config must be provided")
}

writer, err := wal.NewWriter(walCfg, logger, reg)
writer, err := wal.NewWriter(walCfg, logger, reg, wal.NewWriterMetrics(reg))
if err != nil {
return nil, fmt.Errorf("error creating wal writer: %w", err)
}
Expand Down
39 changes: 6 additions & 33 deletions internal/component/common/loki/wal/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,15 +55,13 @@ type Writer struct {
writeSubscribersLock sync.RWMutex
writeSubscribers []WriteEventSubscriber

reclaimedOldSegmentsSpaceCounter *prometheus.CounterVec
lastReclaimedSegment *prometheus.GaugeVec
lastWrittenTimestamp *prometheus.GaugeVec
metrics *WriterMetrics

closeCleaner chan struct{}
}

// NewWriter creates a new Writer.
func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer) (*Writer, error) {
func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer, metrics *WriterMetrics) (*Writer, error) {
// Start WAL
wl, err := New(Config{
Dir: walCfg.Dir,
Expand All @@ -80,32 +78,7 @@ func NewWriter(walCfg Config, logger *slog.Logger, reg prometheus.Registerer) (*
wal: wl,
entryWriter: newEntryWriter(),
closeCleaner: make(chan struct{}, 1),
}

wrt.reclaimedOldSegmentsSpaceCounter = prometheus.NewCounterVec(prometheus.CounterOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "reclaimed_space",
Help: "Number of bytes reclaimed from storage.",
}, []string{})

wrt.lastReclaimedSegment = prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "last_reclaimed_segment",
Help: "Last reclaimed segment number",
}, []string{})
wrt.lastWrittenTimestamp = prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "last_written_timestamp",
Help: "Latest timestamp that was written to the WAL",
}, []string{})

if reg != nil {
_ = reg.Register(wrt.reclaimedOldSegmentsSpaceCounter)
_ = reg.Register(wrt.lastReclaimedSegment)
_ = reg.Register(wrt.lastWrittenTimestamp)
metrics: metrics,
}

return wrt, nil
Expand All @@ -122,7 +95,7 @@ func (wrt *Writer) Start(maxSegmentAge time.Duration) {
}

// emit metric with latest written timestamp, to be able to track delay from writer to watcher
wrt.lastWrittenTimestamp.WithLabelValues().Set(float64(e.Timestamp.Unix()))
wrt.metrics.lastWrittenTimestamp.WithLabelValues().Set(float64(e.Timestamp.Unix()))

wrt.writeSubscribersLock.RLock()
for _, s := range wrt.writeSubscribers {
Expand Down Expand Up @@ -202,7 +175,7 @@ func (wrt *Writer) cleanSegments(maxAge time.Duration) error {
wrt.logger.Error("Error old wal segment", "err", err, "segmentNum", segment.number)
}
wrt.logger.Debug("Deleted old wal segment", "segmentNum", segment.number)
wrt.reclaimedOldSegmentsSpaceCounter.WithLabelValues().Add(float64(segment.size))
wrt.metrics.reclaimedOldSegmentsSpaceCounter.WithLabelValues().Add(float64(segment.size))
// keep track of the largest segment number reclaimed
if segment.number > maxReclaimed {
maxReclaimed = segment.number
Expand All @@ -216,7 +189,7 @@ func (wrt *Writer) cleanSegments(maxAge time.Duration) error {
for _, subscriber := range wrt.cleanupSubscribers {
subscriber.SeriesReset(maxReclaimed)
}
wrt.lastReclaimedSegment.WithLabelValues().Set(float64(maxReclaimed))
wrt.metrics.lastReclaimedSegment.WithLabelValues().Set(float64(maxReclaimed))
}
return nil
}
Expand Down
43 changes: 43 additions & 0 deletions internal/component/common/loki/wal/writer_metrics.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
package wal

import (
"github.com/grafana/alloy/internal/util"
"github.com/prometheus/client_golang/prometheus"
)

type WriterMetrics struct {
lastReclaimedSegment *prometheus.GaugeVec
lastWrittenTimestamp *prometheus.GaugeVec
reclaimedOldSegmentsSpaceCounter *prometheus.CounterVec
}

func NewWriterMetrics(reg prometheus.Registerer) *WriterMetrics {
m := &WriterMetrics{
lastReclaimedSegment: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "last_reclaimed_segment",
Help: "Last reclaimed segment number",
}, []string{}),
lastWrittenTimestamp: prometheus.NewGaugeVec(prometheus.GaugeOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "last_written_timestamp",
Help: "Latest timestamp that was written to the WAL",
}, []string{}),
reclaimedOldSegmentsSpaceCounter: prometheus.NewCounterVec(prometheus.CounterOpts{
Namespace: "loki_write",
Subsystem: "wal_writer",
Name: "reclaimed_space",
Help: "Number of bytes reclaimed from storage.",
}, []string{}),
}

if reg != nil {
m.lastReclaimedSegment = util.MustRegisterOrGet(reg, m.lastReclaimedSegment).(*prometheus.GaugeVec)
m.lastWrittenTimestamp = util.MustRegisterOrGet(reg, m.lastWrittenTimestamp).(*prometheus.GaugeVec)
m.reclaimedOldSegmentsSpaceCounter = util.MustRegisterOrGet(reg, m.reclaimedOldSegmentsSpaceCounter).(*prometheus.CounterVec)
}

return m
}
97 changes: 80 additions & 17 deletions internal/component/common/loki/wal/writer_test.go

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.

I think the tests are not able still to discover this as a bug. Would it be possible to add a check for this?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done

Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,16 @@ import (
"fmt"
"os"
"path/filepath"
"strings"
"sync"
"testing"
"time"

"github.com/grafana/loki/pkg/push"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/testutil"
"github.com/prometheus/common/model"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"

"github.com/grafana/alloy/internal/component/common/loki"
Expand All @@ -20,13 +23,16 @@ import (
)

func TestWriter_EntriesAreWrittenToWAL(t *testing.T) {
dir := t.TempDir()
var (
dir = t.TempDir()
reg = prometheus.NewRegistry()
)

writer, err := NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: time.Minute,
}, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry())
}, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg))
require.NoError(t, err)
defer func() {
writer.Stop()
Expand Down Expand Up @@ -62,6 +68,58 @@ func TestWriter_EntriesAreWrittenToWAL(t *testing.T) {
require.Equal(t, testLabels, readEntries[0].Labels)
}

func TestWriter_MetricsWorkAfterRecreation(t *testing.T) {
var (
dir = t.TempDir()
reg = prometheus.NewRegistry()
)

writer, err := NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: time.Minute,
}, logging.NewSlogNop(), reg, NewWriterMetrics(reg))
require.NoError(t, err)

writer.Start(time.Minute)
entry := loki.NewEntry(model.LabelSet{"foo": "bar"}, push.Entry{Timestamp: time.Now(), Line: "line"})
writer.Chan() <- entry

expected := fmt.Sprintf(`
# HELP loki_write_wal_writer_last_written_timestamp Latest timestamp that was written to the WAL
# TYPE loki_write_wal_writer_last_written_timestamp gauge
loki_write_wal_writer_last_written_timestamp %d
`, entry.Timestamp.Unix())

require.EventuallyWithT(t, func(c *assert.CollectT) {
require.NoError(c, testutil.GatherAndCompare(reg, strings.NewReader(expected), "loki_write_wal_writer_last_written_timestamp"))
}, 2*time.Second, 100*time.Millisecond)

writer.Stop()

writer, err = NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: time.Minute,
}, logging.NewSlogNop(), reg, NewWriterMetrics(reg))
require.NoError(t, err)
writer.Start(time.Minute)
defer writer.Stop()

newEntry := loki.NewEntry(model.LabelSet{"foo": "bar"}, push.Entry{Timestamp: time.Now().Add(1 * time.Second), Line: "line"})
writer.Chan() <- newEntry

expected = fmt.Sprintf(`
# HELP loki_write_wal_writer_last_written_timestamp Latest timestamp that was written to the WAL
# TYPE loki_write_wal_writer_last_written_timestamp gauge
loki_write_wal_writer_last_written_timestamp %d
`, newEntry.Timestamp.Unix())

require.EventuallyWithT(t, func(c *assert.CollectT) {
require.NoError(c, testutil.GatherAndCompare(reg, strings.NewReader(expected), "loki_write_wal_writer_last_written_timestamp"))
}, 2*time.Second, 100*time.Millisecond)
}

type notifySegmentsCleanedFunc func(num int)

func (n notifySegmentsCleanedFunc) NotifyWrite() {
Expand All @@ -72,18 +130,19 @@ func (n notifySegmentsCleanedFunc) SeriesReset(segmentNum int) {
}

func TestWriter_OldSegmentsAreCleanedUp(t *testing.T) {
dir := t.TempDir()

maxSegmentAge := time.Second * 2

subscriber1 := []int{}
subscriber2 := []int{}
var (
dir = t.TempDir()
reg = prometheus.NewRegistry()
maxSegmentAge = time.Second * 2
subscriber1 = []int{}
subscriber2 = []int{}
)

writer, err := NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: maxSegmentAge,
}, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry())
}, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg))
require.NoError(t, err)
writer.Start(maxSegmentAge)
defer func() {
Expand Down Expand Up @@ -166,17 +225,18 @@ func TestWriter_OldSegmentsAreCleanedUp(t *testing.T) {
}

func TestWriter_NoSegmentIsCleanedUpIfTheresOnlyOne(t *testing.T) {
dir := t.TempDir()

maxSegmentAge := time.Second * 2

segmentsReclaimedNotificationsReceived := []int{}
var (
dir = t.TempDir()
reg = prometheus.NewRegistry()
maxSegmentAge = time.Second * 2
segmentsReclaimedNotificationsReceived = []int{}
)

writer, err := NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: maxSegmentAge,
}, autil.TestAlloyLogger(t).Slog(), prometheus.NewRegistry())
}, autil.TestAlloyLogger(t).Slog(), reg, NewWriterMetrics(reg))
require.NoError(t, err)
writer.Start(maxSegmentAge)
defer func() {
Expand Down Expand Up @@ -342,13 +402,16 @@ func BenchmarkWriter_WriteEntries(b *testing.B) {
}

func benchWriteEntries(b *testing.B, lines, labelSetCount int) {
dir := b.TempDir()
var (
dir = b.TempDir()
reg = prometheus.NewRegistry()
)

writer, err := NewWriter(Config{
Dir: dir,
Enabled: true,
MaxSegmentAge: time.Minute,
}, logging.NewSlogNop(), prometheus.NewRegistry())
}, logging.NewSlogNop(), reg, NewWriterMetrics(reg))
require.NoError(b, err)
writer.Start(time.Minute)
defer func() {
Expand Down
Loading