mirror of
https://github.com/milvus-io/milvus.git
synced 2026-07-21 10:15:43 +00:00
issue: #35917 - mlog package: move logger initialization, zap core, async buffered writes, field helpers, and scoped logger binding into pkg/mlog while removing pkg/log. - logging callsites: migrate Milvus logging usage to context-aware mlog APIs and simplify redundant With chains across components, utilities, tests, and tools. - trace propagation: replace logutil trace interceptors with mlog/tracer integration and add client_request_id fallback propagation for server stats handlers. --------- Signed-off-by: chyezh <chyezh@outlook.com>
327 lines
10 KiB
Go
327 lines
10 KiB
Go
package datacoord
|
|
|
|
import (
|
|
"context"
|
|
"math"
|
|
"testing"
|
|
|
|
"github.com/samber/lo"
|
|
"github.com/stretchr/testify/suite"
|
|
|
|
"github.com/milvus-io/milvus-proto/go-api/v3/commonpb"
|
|
"github.com/milvus-io/milvus-proto/go-api/v3/msgpb"
|
|
"github.com/milvus-io/milvus/pkg/v3/mlog"
|
|
"github.com/milvus-io/milvus/pkg/v3/proto/datapb"
|
|
"github.com/milvus-io/milvus/pkg/v3/util/paramtable"
|
|
)
|
|
|
|
func TestLevelZeroSegmentsViewSuite(t *testing.T) {
|
|
suite.Run(t, new(LevelZeroSegmentsViewSuite))
|
|
}
|
|
|
|
type LevelZeroSegmentsViewSuite struct {
|
|
suite.Suite
|
|
v *LevelZeroCompactionView
|
|
}
|
|
|
|
func genTestL0SegmentView(ID UniqueID, label *CompactionGroupLabel, posTime Timestamp) *SegmentView {
|
|
return &SegmentView{
|
|
ID: ID,
|
|
label: label,
|
|
dmlPos: &msgpb.MsgPosition{Timestamp: posTime},
|
|
Level: datapb.SegmentLevel_L0,
|
|
State: commonpb.SegmentState_Flushed,
|
|
}
|
|
}
|
|
|
|
func (s *LevelZeroSegmentsViewSuite) SetupTest() {
|
|
label := &CompactionGroupLabel{
|
|
CollectionID: 1,
|
|
PartitionID: 10,
|
|
Channel: "ch-1",
|
|
}
|
|
segments := []*SegmentView{
|
|
genTestL0SegmentView(100, label, 10000),
|
|
genTestL0SegmentView(101, label, 10000),
|
|
genTestL0SegmentView(102, label, 10000),
|
|
}
|
|
|
|
targetView := &LevelZeroCompactionView{
|
|
label: label,
|
|
l0Segments: segments,
|
|
latestDeletePos: &msgpb.MsgPosition{Timestamp: 10000},
|
|
triggerID: 10000,
|
|
}
|
|
|
|
s.True(label.Equal(targetView.GetGroupLabel()))
|
|
mlog.Info(context.TODO(), "LevelZeroSegmentsView", mlog.String("view", targetView.String()))
|
|
|
|
s.v = targetView
|
|
}
|
|
|
|
func (s *LevelZeroSegmentsViewSuite) TestTrigger() {
|
|
label := s.v.GetGroupLabel()
|
|
views := []*SegmentView{
|
|
genTestL0SegmentView(100, label, 20000),
|
|
genTestL0SegmentView(101, label, 10000),
|
|
genTestL0SegmentView(102, label, 30000),
|
|
genTestL0SegmentView(103, label, 40000),
|
|
}
|
|
|
|
s.v.l0Segments = views
|
|
tests := []struct {
|
|
description string
|
|
|
|
prepSizeEach float64
|
|
prepCountEach int
|
|
prepEarliestT Timestamp
|
|
|
|
expectedSegs []UniqueID
|
|
}{
|
|
{
|
|
"Not qualified",
|
|
1,
|
|
1,
|
|
30000,
|
|
nil,
|
|
},
|
|
{
|
|
"Trigger by > TriggerDeltaSize",
|
|
8 * 1024 * 1024,
|
|
1,
|
|
30000,
|
|
[]UniqueID{100, 101, 102, 103},
|
|
},
|
|
{
|
|
"Trigger by > TriggerDeltaCount",
|
|
1,
|
|
10,
|
|
30000,
|
|
[]UniqueID{100, 101, 102, 103},
|
|
},
|
|
{
|
|
"Trigger by > maxDeltaSize",
|
|
128 * 1024 * 1024,
|
|
1,
|
|
30000,
|
|
[]UniqueID{100},
|
|
},
|
|
{
|
|
"Trigger by > maxDeltaCount",
|
|
1,
|
|
800,
|
|
30000,
|
|
[]UniqueID{100},
|
|
},
|
|
}
|
|
|
|
for _, test := range tests {
|
|
s.Run(test.description, func() {
|
|
s.v.latestDeletePos.Timestamp = test.prepEarliestT
|
|
for _, view := range s.v.GetSegmentsView() {
|
|
if view.dmlPos.Timestamp < test.prepEarliestT {
|
|
view.DeltalogCount = test.prepCountEach
|
|
view.DeltaSize = test.prepSizeEach
|
|
view.DeltaRowCount = 1
|
|
}
|
|
}
|
|
mlog.Info(context.TODO(), "LevelZeroSegmentsView", mlog.String("view", s.v.String()))
|
|
|
|
gotView, reason := s.v.Trigger()
|
|
if len(test.expectedSegs) == 0 {
|
|
s.Nil(gotView)
|
|
} else {
|
|
levelZeroView, ok := gotView.(*LevelZeroCompactionView)
|
|
s.True(ok)
|
|
s.NotNil(levelZeroView)
|
|
|
|
gotSegIDs := lo.Map(levelZeroView.GetSegmentsView(), func(v *SegmentView, _ int) int64 {
|
|
return v.ID
|
|
})
|
|
s.ElementsMatch(gotSegIDs, test.expectedSegs)
|
|
mlog.Info(context.TODO(), "output view", mlog.String("view", levelZeroView.String()), mlog.String("trigger reason", reason))
|
|
}
|
|
})
|
|
}
|
|
}
|
|
|
|
func (s *LevelZeroSegmentsViewSuite) TestMinCountSizeTrigger() {
|
|
label := s.v.GetGroupLabel()
|
|
tests := []struct {
|
|
description string
|
|
segIDs []int64
|
|
segCounts []int
|
|
segSize []float64
|
|
|
|
expectedIDs []int64
|
|
}{
|
|
{"donot trigger", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{1, 1, 1}, nil},
|
|
{"trigger by count=15", []int64{100, 101, 102}, []int{5, 5, 5}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by count=10", []int64{100, 101, 102}, []int{5, 3, 2}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by count=50", []int64{100, 101, 102}, []int{32, 10, 8}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by size=24MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{8 * 1024 * 1024, 8 * 1024 * 1024, 8 * 1024 * 1024}, []int64{100, 101, 102}},
|
|
{"trigger by size=8MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{3 * 1024 * 1024, 3 * 1024 * 1024, 2 * 1024 * 1024}, []int64{100, 101, 102}},
|
|
{"trigger by size=128MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{100 * 1024 * 1024, 20 * 1024 * 1024, 8 * 1024 * 1024}, []int64{100}},
|
|
}
|
|
|
|
for _, test := range tests {
|
|
s.Run(test.description, func() {
|
|
views := []*SegmentView{}
|
|
for idx, ID := range test.segIDs {
|
|
seg := genTestL0SegmentView(ID, label, 10000)
|
|
seg.DeltaSize = test.segSize[idx]
|
|
seg.DeltalogCount = test.segCounts[idx]
|
|
|
|
views = append(views, seg)
|
|
}
|
|
|
|
picked, reason := s.v.minCountSizeTrigger(views)
|
|
s.ElementsMatch(lo.Map(picked, func(view *SegmentView, _ int) int64 {
|
|
return view.ID
|
|
}), test.expectedIDs)
|
|
|
|
if len(picked) > 0 {
|
|
s.NotEmpty(reason)
|
|
}
|
|
|
|
mlog.Info(context.TODO(), "test minCountSizeTrigger", mlog.Any("trigger reason", reason))
|
|
})
|
|
}
|
|
}
|
|
|
|
func (s *LevelZeroSegmentsViewSuite) TestForceTrigger() {
|
|
label := s.v.GetGroupLabel()
|
|
tests := []struct {
|
|
description string
|
|
segIDs []int64
|
|
segCounts []int
|
|
segSize []float64
|
|
|
|
expectedIDs []int64
|
|
}{
|
|
{"force trigger", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by count=15", []int64{100, 101, 102}, []int{5, 5, 5}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by count=10", []int64{100, 101, 102}, []int{5, 3, 2}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by count=50", []int64{100, 101, 102}, []int{32, 10, 8}, []float64{1, 1, 1}, []int64{100, 101, 102}},
|
|
{"trigger by size=24MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{8 * 1024 * 1024, 8 * 1024 * 1024, 8 * 1024 * 1024}, []int64{100, 101, 102}},
|
|
{"trigger by size=8MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{3 * 1024 * 1024, 3 * 1024 * 1024, 2 * 1024 * 1024}, []int64{100, 101, 102}},
|
|
{"trigger by size=128MB", []int64{100, 101, 102}, []int{1, 1, 1}, []float64{100 * 1024 * 1024, 20 * 1024 * 1024, 8 * 1024 * 1024}, []int64{100}},
|
|
}
|
|
|
|
for _, test := range tests {
|
|
s.Run(test.description, func() {
|
|
views := []*SegmentView{}
|
|
for idx, ID := range test.segIDs {
|
|
seg := genTestL0SegmentView(ID, label, 10000)
|
|
seg.DeltaSize = test.segSize[idx]
|
|
seg.DeltalogCount = test.segCounts[idx]
|
|
|
|
views = append(views, seg)
|
|
}
|
|
|
|
picked, reason := s.v.forceTrigger(views)
|
|
s.ElementsMatch(lo.Map(picked, func(view *SegmentView, _ int) int64 {
|
|
return view.ID
|
|
}), test.expectedIDs)
|
|
mlog.Info(context.TODO(), "test forceTrigger", mlog.Any("trigger reason", reason))
|
|
})
|
|
}
|
|
|
|
// Test that exceeding maxCount from paramtable picks only the first segment.
|
|
s.Run("trigger by exceeding maxCount from param", func() {
|
|
maxCount := paramtable.Get().DataCoordCfg.LevelZeroCompactionTriggerDeltalogMaxNum.GetAsInt()
|
|
views := []*SegmentView{
|
|
genTestL0SegmentView(100, label, 10000),
|
|
genTestL0SegmentView(101, label, 10000),
|
|
}
|
|
views[0].DeltaSize = 1
|
|
views[0].DeltalogCount = maxCount
|
|
views[1].DeltaSize = 1
|
|
views[1].DeltalogCount = maxCount
|
|
|
|
picked, reason := s.v.forceTrigger(views)
|
|
s.ElementsMatch(lo.Map(picked, func(view *SegmentView, _ int) int64 {
|
|
return view.ID
|
|
}), []int64{100})
|
|
mlog.Info(context.TODO(), "test forceTrigger", mlog.Any("trigger reason", reason))
|
|
})
|
|
}
|
|
|
|
func (s *LevelZeroSegmentsViewSuite) TestResolveLatestDeletePos() {
|
|
paramtable.Init()
|
|
dmlPos := &msgpb.MsgPosition{ChannelName: "ch-1", Timestamp: 12345}
|
|
|
|
s.Run("default_returns_dml_pos", func() {
|
|
paramtable.Get().Save("dataCoord.compaction.levelzero.forceSelectAllSegments", "false")
|
|
defer paramtable.Get().Reset("dataCoord.compaction.levelzero.forceSelectAllSegments")
|
|
|
|
got := resolveLatestDeletePos(dmlPos)
|
|
s.Equal(dmlPos, got)
|
|
})
|
|
|
|
s.Run("force_select_returns_max_timestamp", func() {
|
|
paramtable.Get().Save("dataCoord.compaction.levelzero.forceSelectAllSegments", "true")
|
|
defer paramtable.Get().Reset("dataCoord.compaction.levelzero.forceSelectAllSegments")
|
|
|
|
got := resolveLatestDeletePos(dmlPos)
|
|
s.Equal(uint64(math.MaxUint64), got.GetTimestamp())
|
|
s.Equal("ch-1", got.GetChannelName())
|
|
})
|
|
|
|
s.Run("force_select_with_nil_pos", func() {
|
|
paramtable.Get().Save("dataCoord.compaction.levelzero.forceSelectAllSegments", "true")
|
|
defer paramtable.Get().Reset("dataCoord.compaction.levelzero.forceSelectAllSegments")
|
|
|
|
got := resolveLatestDeletePos(nil)
|
|
s.Equal(uint64(math.MaxUint64), got.GetTimestamp())
|
|
s.Equal("", got.GetChannelName())
|
|
})
|
|
|
|
s.Run("trigger_uses_resolved_pos_when_force_select_enabled", func() {
|
|
paramtable.Get().Save("dataCoord.compaction.levelzero.forceSelectAllSegments", "true")
|
|
defer paramtable.Get().Reset("dataCoord.compaction.levelzero.forceSelectAllSegments")
|
|
|
|
label := s.v.GetGroupLabel()
|
|
views := []*SegmentView{
|
|
genTestL0SegmentView(100, label, 20000),
|
|
genTestL0SegmentView(101, label, 10000),
|
|
genTestL0SegmentView(102, label, 30000),
|
|
}
|
|
for _, v := range views {
|
|
v.DeltalogCount = 100
|
|
v.DeltaSize = 1
|
|
v.DeltaRowCount = 1
|
|
}
|
|
s.v.l0Segments = views
|
|
|
|
gotView, _ := s.v.Trigger()
|
|
s.Require().NotNil(gotView)
|
|
levelZeroView, ok := gotView.(*LevelZeroCompactionView)
|
|
s.Require().True(ok)
|
|
s.Equal(uint64(math.MaxUint64), levelZeroView.latestDeletePos.GetTimestamp())
|
|
})
|
|
|
|
s.Run("force_trigger_uses_resolved_pos_when_force_select_enabled", func() {
|
|
paramtable.Get().Save("dataCoord.compaction.levelzero.forceSelectAllSegments", "true")
|
|
defer paramtable.Get().Reset("dataCoord.compaction.levelzero.forceSelectAllSegments")
|
|
|
|
label := s.v.GetGroupLabel()
|
|
views := []*SegmentView{
|
|
genTestL0SegmentView(100, label, 20000),
|
|
genTestL0SegmentView(101, label, 10000),
|
|
}
|
|
for _, v := range views {
|
|
v.DeltalogCount = 1
|
|
v.DeltaSize = 1
|
|
v.DeltaRowCount = 1
|
|
}
|
|
s.v.l0Segments = views
|
|
|
|
gotView, _ := s.v.ForceTrigger()
|
|
s.Require().NotNil(gotView)
|
|
levelZeroView, ok := gotView.(*LevelZeroCompactionView)
|
|
s.Require().True(ok)
|
|
s.Equal(uint64(math.MaxUint64), levelZeroView.latestDeletePos.GetTimestamp())
|
|
})
|
|
}
|