milvus/internal/querynodev2/segments/count_reducer.go
Jiquan Long c002745902
enhance: retrieve output fields after local reduce (#32346)
issue: #31822

---------

Signed-off-by: longjiquan <jiquan.long@zilliz.com>
2024-04-25 09:49:26 +08:00

51 lines
1.4 KiB
Go

package segments
import (
"context"
"github.com/milvus-io/milvus/internal/proto/internalpb"
"github.com/milvus-io/milvus/internal/proto/segcorepb"
"github.com/milvus-io/milvus/internal/util/funcutil"
)
type cntReducer struct{}
func (r *cntReducer) Reduce(ctx context.Context, results []*internalpb.RetrieveResults) (*internalpb.RetrieveResults, error) {
cnt := int64(0)
allRetrieveCount := int64(0)
relatedDataSize := int64(0)
for _, res := range results {
allRetrieveCount += res.GetAllRetrieveCount()
relatedDataSize += res.GetCostAggregation().GetTotalRelatedDataSize()
c, err := funcutil.CntOfInternalResult(res)
if err != nil {
return nil, err
}
cnt += c
}
res := funcutil.WrapCntToInternalResult(cnt)
res.AllRetrieveCount = allRetrieveCount
res.CostAggregation = &internalpb.CostAggregation{
TotalRelatedDataSize: relatedDataSize,
}
return res, nil
}
type cntReducerSegCore struct{}
func (r *cntReducerSegCore) Reduce(ctx context.Context, results []*segcorepb.RetrieveResults, _ []Segment, _ *RetrievePlan) (*segcorepb.RetrieveResults, error) {
cnt := int64(0)
allRetrieveCount := int64(0)
for _, res := range results {
allRetrieveCount += res.GetAllRetrieveCount()
c, err := funcutil.CntOfSegCoreResult(res)
if err != nil {
return nil, err
}
cnt += c
}
res := funcutil.WrapCntToSegCoreResult(cnt)
res.AllRetrieveCount = allRetrieveCount
return res, nil
}