feat: get os traffic from minio (#4968)

* get os traffic from minio

* optimize get obj traffic within the last hour

---------

Co-authored-by: jiahui <bxy4543@gmail.com>
This commit is contained in:
xuziyi
2024-08-20 18:04:22 +08:00
committed by GitHub
co-authored by jiahui
parent 3f6343238e
commit cca4da935f
9 changed files with 514 additions and 51 deletions
+6
View File
@@ -36,6 +36,7 @@ type Interface interface {
Account
Traffic
CVM
Creator
}
type CVM interface {
@@ -49,6 +50,10 @@ type Account interface {
GetBillingHistoryNamespaceList(ns *accountv1.NamespaceBillingHistorySpec, owner string) ([]string, error)
GetBillingHistoryNamespaces(startTime, endTime *time.Time, billType int, owner string) ([]string, error)
SaveBillings(billing ...*resources.Billing) error
SaveObjTraffic(obs ...*types.ObjectStorageTraffic) error
GetAllLatestObjTraffic(startTime, endTime time.Time) ([]types.ObjectStorageTraffic, error)
HandlerTimeObjBucketSentTraffic(startTime, endTime time.Time, bucket string) (int64, error)
GetTimeObjBucketBucket(startTime, endTime time.Time) ([]string, error)
QueryBillingRecords(billingRecordQuery *accountv1.BillingRecordQuery, owner string) error
GetUnsettingBillingHandler(owner string) ([]resources.BillingHandler, error)
UpdateBillingStatus(orderID string, status resources.BillingStatus) error
@@ -117,6 +122,7 @@ type Creator interface {
CreateBillingIfNotExist() error
//suffix by day, eg: monitor_20200101
CreateMonitorTimeSeriesIfNotExist(collTime time.Time) error
CreateTTLTrafficTimeSeries() error
}
type MeteringOwnerTimeResult struct {
+147
View File
@@ -22,6 +22,8 @@ import (
"strings"
"time"
"github.com/labring/sealos/controllers/pkg/types"
"github.com/labring/sealos/controllers/pkg/utils/env"
"github.com/labring/sealos/controllers/pkg/common"
@@ -57,6 +59,7 @@ const (
DefaultMeteringConn = "metering"
DefaultMonitorConn = "monitor"
DefaultBillingConn = "billing"
DefaultObjTrafficConn = "objectstorage-traffic"
DefaultUserConn = "user"
DefaultPricesConn = "prices"
DefaultPropertiesConn = "properties"
@@ -75,6 +78,7 @@ type mongoDB struct {
MonitorConnPrefix string
MeteringConn string
BillingConn string
ObjTrafficConn string
PropertiesConn string
TrafficConn string
}
@@ -246,6 +250,130 @@ func (m *mongoDB) SaveBillings(billing ...*resources.Billing) error {
return err
}
func (m *mongoDB) SaveObjTraffic(obs ...*types.ObjectStorageTraffic) error {
traffic := make([]interface{}, len(obs))
for i, ob := range obs {
traffic[i] = ob
}
_, err := m.getObjTrafficCollection().InsertMany(context.Background(), traffic)
return err
}
func (m *mongoDB) GetAllLatestObjTraffic(startTime, endTime time.Time) ([]types.ObjectStorageTraffic, error) {
pipeline := []bson.M{
{
"$match": bson.M{
"time": bson.M{
"$gt": startTime,
"$lte": endTime,
},
},
},
{
"$sort": bson.M{"time": -1},
},
{
"$group": bson.M{
"_id": bson.M{"user": "$user", "bucket": "$bucket"},
"latestDoc": bson.M{"$first": "$$ROOT"},
},
},
{
"$replaceRoot": bson.M{"newRoot": "$latestDoc"},
},
}
cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline)
if err != nil {
return nil, err
}
defer cursor.Close(context.Background())
var results []types.ObjectStorageTraffic
if err = cursor.All(context.Background(), &results); err != nil {
return nil, err
}
return results, nil
}
func (m *mongoDB) HandlerTimeObjBucketSentTraffic(startTime, endTime time.Time, bucket string) (int64, error) {
pipeline := []bson.M{
{
"$match": bson.M{
"time": bson.M{
"$gt": startTime,
"$lte": endTime,
},
"bucket": bucket,
},
},
{
"$group": bson.M{
"_id": nil,
"totalSent": bson.M{"$sum": "$sent"},
},
},
}
cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline)
if err != nil {
return 0, err
}
defer cursor.Close(context.Background())
var result struct {
TotalSent int64 `bson:"totalSent"`
}
if cursor.Next(context.Background()) {
if err := cursor.Decode(&result); err != nil {
return 0, err
}
return result.TotalSent, nil
}
if err := cursor.Err(); err != nil {
return 0, err
}
return 0, nil
}
func (m *mongoDB) GetTimeObjBucketBucket(startTime, endTime time.Time) ([]string, error) {
pipeline := []bson.M{
{
"$match": bson.M{
"time": bson.M{
"$gt": startTime,
"$lte": endTime,
},
},
},
{
"$group": bson.M{
"_id": nil,
"buckets": bson.M{"$addToSet": "$bucket"},
},
},
}
cursor, err := m.getObjTrafficCollection().Aggregate(context.Background(), pipeline)
if err != nil {
return nil, err
}
defer cursor.Close(context.Background())
var result struct {
Buckets []string `bson:"buckets"`
}
if cursor.Next(context.Background()) {
if err := cursor.Decode(&result); err != nil {
return nil, err
}
return result.Buckets, nil
}
if err := cursor.Err(); err != nil {
return nil, err
}
return nil, nil
}
// InsertMonitor insert monitor data to mongodb collection monitor + time (eg: monitor_20200101)
// The monitor data is saved daily 2020-12-01 00:00:00 - 2020-12-01 23:59:59 => monitor_20201201
func (m *mongoDB) InsertMonitor(ctx context.Context, monitors ...*resources.Monitor) error {
@@ -885,6 +1013,9 @@ func (m *mongoDB) getMonitorCollectionName(collTime time.Time) string {
func (m *mongoDB) getBillingCollection() *mongo.Collection {
return m.Client.Database(m.AccountDB).Collection(m.BillingConn)
}
func (m *mongoDB) getObjTrafficCollection() *mongo.Collection {
return m.Client.Database(m.AccountDB).Collection(m.ObjTrafficConn)
}
func (m *mongoDB) getPropertiesCollection() *mongo.Collection {
return m.Client.Database(m.AccountDB).Collection(m.PropertiesConn)
@@ -941,6 +1072,21 @@ func (m *mongoDB) CreateTimeSeriesIfNotExist(dbName, collectionName string) erro
return m.Client.Database(dbName).RunCommand(context.TODO(), cmd).Err()
}
func (m *mongoDB) CreateTTLTrafficTimeSeries() error {
// Check if the collection already exists
if exist, err := m.collectionExist(m.AccountDB, m.ObjTrafficConn); exist || err != nil {
return err
}
// If the collection does not exist, create it
cmd := bson.D{
primitive.E{Key: "create", Value: m.ObjTrafficConn},
primitive.E{Key: "timeseries", Value: bson.D{{Key: "timeField", Value: "time"}}},
//default ttl set 30 days
primitive.E{Key: "expireAfterSeconds", Value: 30 * 24 * 60 * 60},
}
return m.Client.Database(m.AccountDB).RunCommand(context.TODO(), cmd).Err()
}
func (m *mongoDB) DropMonitorCollectionsOlderThan(days int) error {
db := m.Client.Database(m.AccountDB)
// Get the current time minus the number of days
@@ -985,6 +1131,7 @@ func NewMongoInterface(ctx context.Context, URL string) (database.Interface, err
MeteringConn: DefaultMeteringConn,
MonitorConnPrefix: DefaultMonitorConn,
BillingConn: DefaultBillingConn,
ObjTrafficConn: DefaultObjTrafficConn,
PropertiesConn: DefaultPropertiesConn,
TrafficConn: env.GetEnvWithDefault(EnvTrafficConn, DefaultTrafficConn),
CvmConn: env.GetEnvWithDefault(EnvCVMConn, DefaultCVMConn),
@@ -22,6 +22,8 @@ import (
"testing"
"time"
"github.com/labring/sealos/controllers/pkg/types"
"github.com/labring/sealos/controllers/pkg/resources"
"github.com/dustin/go-humanize"
@@ -595,3 +597,120 @@ func Test_mongoDB_GetDistinctMonitorCombinations(t *testing.T) {
}
t.Logf("monitorCombinations: %v", monitorCombinations)
}
func Test_mongoDB_CreateTTLTrafficTimeSeries(t *testing.T) {
dbCTX := context.Background()
m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI"))
if err != nil {
t.Errorf("failed to connect mongo: error = %v", err)
}
defer func() {
if err = m.Disconnect(dbCTX); err != nil {
t.Errorf("failed to disconnect mongo: error = %v", err)
}
}()
if err = m.CreateTTLTrafficTimeSeries(); err != nil {
t.Fatalf("failed to create TTL traffic time series: %v", err)
}
t.Logf("create TTL traffic time series success")
}
func Test_mongoDB_SaveObjTraffic(t *testing.T) {
dbCTX := context.Background()
m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI"))
if err != nil {
t.Errorf("failed to connect mongo: error = %v", err)
}
defer func() {
if err = m.Disconnect(dbCTX); err != nil {
t.Errorf("failed to disconnect mongo: error = %v", err)
}
}()
var traffic []*types.ObjectStorageTraffic
for i := 0; i < 10; i++ {
traffic = append(traffic, &types.ObjectStorageTraffic{
Time: time.Now().UTC(),
User: "user-" + fmt.Sprint(i),
Bucket: "bucket-" + fmt.Sprint(i),
TotalSent: int64(1000 + i),
Sent: int64(100 + i),
})
}
if err = m.SaveObjTraffic(traffic...); err != nil {
t.Fatalf("failed to save object storage traffic: %v", err)
}
t.Logf("save object storage traffic success")
}
func Test_mongoDB_GetAllLatestObjTraffic(t *testing.T) {
dbCTX := context.Background()
m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI"))
if err != nil {
t.Errorf("failed to connect mongo: error = %v", err)
}
defer func() {
if err = m.Disconnect(dbCTX); err != nil {
t.Errorf("failed to disconnect mongo: error = %v", err)
}
}()
traffic, err := m.GetAllLatestObjTraffic(time.Now().Add(-time.Hour), time.Now())
if err != nil {
t.Fatalf("failed to save object storage traffic: %v", err)
}
t.Logf("save object storage traffic success")
for _, tf := range traffic {
t.Logf("traffic: %#+v", tf)
}
}
func Test_mongoDB_HandlerTimeObjBucketSentTraffic(t *testing.T) {
dbCTX := context.Background()
m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI"))
if err != nil {
t.Errorf("failed to connect mongo: error = %v", err)
}
defer func() {
if err = m.Disconnect(dbCTX); err != nil {
t.Errorf("failed to disconnect mongo: error = %v", err)
}
}()
bytes, err := m.HandlerTimeObjBucketSentTraffic(time.Now().UTC().Add(-time.Hour), time.Now().UTC(), "bucket-6")
if err != nil {
t.Fatalf("failed to handle time object bucket usage: %v", err)
}
t.Logf("handle time object bucket usage success: %v", bytes)
}
func init() {
os.Setenv("MONGODB_URI", "")
}
func Test_mongoDB_GetTimeObjBucketBucket(t *testing.T) {
dbCTX := context.Background()
m, err := NewMongoInterface(dbCTX, os.Getenv("MONGODB_URI"))
if err != nil {
t.Errorf("failed to connect mongo: error = %v", err)
}
defer func() {
if err = m.Disconnect(dbCTX); err != nil {
t.Errorf("failed to disconnect mongo: error = %v", err)
}
}()
buckets, err := m.GetTimeObjBucketBucket(time.Now().UTC().Add(-10*time.Hour), time.Now().UTC())
if err != nil {
t.Fatalf("failed to get time object bucket bucket: %v", err)
}
t.Logf("get time object bucket bucket success: len: %v", len(buckets))
for _, bucket := range buckets {
t.Logf("bucket: %#+v", bucket)
}
}
+23 -11
View File
@@ -75,13 +75,18 @@ func NewMetricsClient(endpoint string, accessKeyID, secretAccessKey string, secu
return privateNewMetricsClient(endpointURL, jwtToken, secure)
}
// BucketUsageTotalBytesMetrics - returns Bucket Metrics in Prometheus format
// BucketUsageTotalBytesMetrics - returns Bucket Usage Total Metrics in Prometheus format
func (client *MetricsClient) BucketUsageTotalBytesMetrics(ctx context.Context) ([]*prom2json.Family, error) {
return client.fetchMetrics(ctx, "bucket", "minio_bucket_usage_total_bytes")
return client.fetchMetrics(ctx, "bucket", []string{"minio_bucket_usage_total_bytes"})
}
// BucketUsageAndTrafficBytesMetrics - returns Bucket Usage And Traffic Metrics in Prometheus format
func (client *MetricsClient) BucketUsageAndTrafficBytesMetrics(ctx context.Context) ([]*prom2json.Family, error) {
return client.fetchMetrics(ctx, "bucket", []string{"minio_bucket_usage_total_bytes", "minio_bucket_traffic_sent_bytes", "minio_bucket_traffic_received_bytes"})
}
// fetchMetrics - returns Metrics of given subsystem in Prometheus format
func (client *MetricsClient) fetchMetrics(ctx context.Context, subSystem string, metricsName string) ([]*prom2json.Family, error) {
func (client *MetricsClient) fetchMetrics(ctx context.Context, subSystem string, metricsName []string) ([]*prom2json.Family, error) {
reqData := metricsRequestData{
relativePath: "/v2/metrics/" + subSystem,
}
@@ -119,7 +124,7 @@ func closeResponse(resp *http.Response) {
}
}
func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2json.Family, err error) {
func parsePrometheusResults(reader io.Reader, prefix []string) (results []*prom2json.Family, err error) {
filteredReader, err := filterMetricsByPrefix(reader, prefix)
if err != nil {
return nil, err
@@ -136,10 +141,11 @@ func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2js
}()
for mf := range mfChan {
if !strings.Contains(mf.GetName(), prefix) {
continue
for i := range prefix {
if strings.Contains(mf.GetName(), prefix[i]) {
results = append(results, prom2json.NewFamily(mf))
}
}
results = append(results, prom2json.NewFamily(mf))
}
if err := <-errChan; err != nil {
return nil, err
@@ -147,7 +153,7 @@ func parsePrometheusResults(reader io.Reader, prefix string) (results []*prom2js
return results, nil
}
func filterMetricsByPrefix(reader io.Reader, prefix string) (io.Reader, error) {
func filterMetricsByPrefix(reader io.Reader, prefix []string) (io.Reader, error) {
var buf bytes.Buffer
for {
line, err := readLine(reader)
@@ -156,11 +162,17 @@ func filterMetricsByPrefix(reader io.Reader, prefix string) (io.Reader, error) {
} else if err != nil {
return nil, err
}
if bytes.HasPrefix(line, []byte("#")) || !bytes.HasPrefix(line, []byte(prefix)) {
if bytes.HasPrefix(line, []byte("#")) {
continue
}
if _, err := buf.Write(line); err != nil {
return nil, err
for i := range prefix {
if bytes.HasPrefix(line, []byte(prefix[i])) {
if _, err := buf.Write(line); err != nil {
return nil, err
}
}
}
}
return &buf, nil
+96 -22
View File
@@ -199,6 +199,10 @@ func extractValues(result1, result2 model.Value) (int64, int64) {
type MetricData struct {
// key: bucket name, value: usage
Usage map[string]int64
// key: bucket name, value: traffic sent
Sent map[string]int64
// key: bucket name, value: traffic received
Received map[string]int64
}
type Metrics map[string]MetricData
@@ -237,6 +241,57 @@ func QueryUserUsage(client *MetricsClient) (Metrics, error) {
return obMetrics, err
}
func QueryUserUsageAndTraffic(client *MetricsClient) (Metrics, error) {
obMetrics := make(Metrics)
bucketMetrics, err := client.BucketUsageAndTrafficBytesMetrics(context.TODO())
if err != nil {
return nil, fmt.Errorf("failed to get bucket traffic metrics: %w", err)
}
for _, bucketMetric := range bucketMetrics {
if !isUsageAndTrafficBytesTargetMetric(bucketMetric.Name) {
continue
}
for _, metrics := range bucketMetric.Metrics {
promMetrics := metrics.(prom2json.Metric)
floatValue, err := strconv.ParseFloat(promMetrics.Value, 64)
if err != nil {
return nil, fmt.Errorf("failed to parse %s to float value", promMetrics.Value)
}
intValue := int64(floatValue)
if bucket := promMetrics.Labels["bucket"]; bucket != "" {
//fmt.Println("debug info", "type:", bucketMetric.Name, "promMetrics:", promMetrics)
user := getUserWithBucket(bucket)
if user == "" {
//fmt.Println("debug info", "false bucket:", bucket)
continue
}
//fmt.Println("debug info", "true bucket:", bucket, "user:", user)
metricData, exists := obMetrics[user]
if !exists {
metricData = MetricData{
Usage: make(map[string]int64),
Sent: make(map[string]int64),
Received: make(map[string]int64),
}
}
if bucketMetric.Name == "minio_bucket_usage_total_bytes" {
metricData.Usage[bucket] += intValue
}
if bucketMetric.Name == "minio_bucket_traffic_sent_bytes" {
metricData.Sent[bucket] += intValue
}
if bucketMetric.Name == "minio_bucket_traffic_received_bytes" {
metricData.Received[bucket] += intValue
}
obMetrics[user] = metricData
}
}
}
return obMetrics, err
}
func isUsageBytesTargetMetric(name string) bool {
targetMetrics := []string{
"minio_bucket_usage_total_bytes",
@@ -249,8 +304,27 @@ func isUsageBytesTargetMetric(name string) bool {
return false
}
func isUsageAndTrafficBytesTargetMetric(name string) bool {
targetMetrics := []string{
"minio_bucket_usage_total_bytes",
"minio_bucket_traffic_sent_bytes",
"minio_bucket_traffic_received_bytes",
}
for _, target := range targetMetrics {
if name == target {
return true
}
}
return false
}
func getUserWithBucket(bucket string) string {
return strings.Split(bucket, "-")[0]
re := regexp.MustCompile(`^([a-zA-Z0-9]{8})-(.*)$`)
matches := re.FindStringSubmatch(bucket)
if len(matches) == 3 {
return matches[1]
}
return ""
}
/*
@@ -268,25 +342,25 @@ func getUserWithBucket(bucket string) string {
6/halo-faxdridb-pg-yhxnjm.tar.gz
*/
func GetUserBakFileSize(client *minio.Client) map[string]int64 {
bucket := "file-backup"
userUsageMap := make(map[string]int64)
objectsCh := client.ListObjects(context.Background(), bucket, minio.ListObjectsOptions{Recursive: true})
for object := range objectsCh {
user := extractNamespace(object.Key)
if user != "" {
userUsageMap[user] += object.Size
}
}
//func GetUserBakFileSize(client *minio.Client) map[string]int64 {
// bucket := "file-backup"
// userUsageMap := make(map[string]int64)
// objectsCh := client.ListObjects(context.Background(), bucket, minio.ListObjectsOptions{Recursive: true})
// for object := range objectsCh {
// user := extractNamespace(object.Key)
// if user != "" {
// userUsageMap[user] += object.Size
// }
// }
//
// return userUsageMap
//}
return userUsageMap
}
func extractNamespace(input string) string {
re := regexp.MustCompile(`ns-(\w+)`)
matches := re.FindStringSubmatch(input)
if len(matches) < 2 {
return ""
}
return matches[1]
}
//func extractNamespace(input string) string {
// re := regexp.MustCompile(`ns-(\w+)`)
// matches := re.FindStringSubmatch(input)
// if len(matches) < 2 {
// return ""
// }
// return matches[1]
//}
@@ -62,14 +62,20 @@ func TestQueryUserUsage(t *testing.T) {
}
}
func TestGetUserBakFileSize(t *testing.T) {
objStorageClient, err := objectstoragev1.NewOSClient("objectstorageapi.192.168.0.55.nip.io", "username", "passw0rd")
func TestQueryUserTraffic(t *testing.T) {
obClient, err := NewMetricsClient("objectstorageapi.192.168.0.55.nip.io", "username", "passw0rd", false)
if err != nil {
t.Fatalf("NewOSClient error: %v", err)
t.Error(err)
}
bakFileSize := GetUserBakFileSize(objStorageClient)
metrics, err := QueryUserUsageAndTraffic(obClient)
if err != nil {
t.Fatalf("GetUserBakFileSize error: %v", err)
t.Error(err)
}
for user, metric := range metrics {
fmt.Println("user:", user)
fmt.Println("usage:", metric.Usage)
fmt.Println("sent:", metric.Sent)
fmt.Println("received:", metric.Received)
}
fmt.Println(bakFileSize)
}
+28
View File
@@ -0,0 +1,28 @@
// Copyright © 2024 sealos.
//
// 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.
package types
import "time"
type ObjectStorageTraffic struct {
Time time.Time `json:"time" bson:"time"`
User string `json:"user" bson:"user"`
Bucket string `json:"bucket" bson:"bucket"`
//bytes
TotalSent int64 `json:"totalSent" bson:"totalSent"`
//The sent traffic since the last time
Sent int64 `json:"sent" bson:"sent"`
}
@@ -25,6 +25,8 @@ import (
"sync"
"time"
"github.com/labring/sealos/controllers/pkg/types"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
appv1 "github.com/labring/sealos/controllers/app/api/v1"
@@ -80,7 +82,8 @@ type MonitorReconciler struct {
TrafficClient database.Interface
Properties *resources.PropertyTypeLS
PromURL string
currentObjectMetrics map[string]objstorage.MetricData
lastObjectMetrics objstorage.Metrics
currentObjectMetrics objstorage.Metrics
ObjStorageClient *minio.Client
ObjStorageMetricsClient *objstorage.MetricsClient
ObjStorageUserBackupSize map[string]int64
@@ -288,22 +291,52 @@ func (r *MonitorReconciler) processNamespaceList(namespaceList *corev1.Namespace
}(&namespaceList.Items[i])
}
wg.Wait()
if err := r.monitorObjectStorageTraffic(); err != nil {
r.Logger.Error(err, "failed to monitor object storage traffic")
}
logger.Info("end processNamespaceList", "time", time.Now().Format("2006-01-02 15:04:05"))
return nil
}
func (r *MonitorReconciler) preMonitorResourceUsage() error {
if r.ObjStorageMetricsClient != nil {
metrics, err := objstorage.QueryUserUsage(r.ObjStorageMetricsClient)
metrics, err := objstorage.QueryUserUsageAndTraffic(r.ObjStorageMetricsClient)
if err != nil {
r.lastObjectMetrics = r.currentObjectMetrics
return fmt.Errorf("failed to query object storage metrics: %w", err)
}
if r.currentObjectMetrics != nil {
r.lastObjectMetrics = r.currentObjectMetrics
} else {
latestObjTrafficSentMetrics := make(objstorage.Metrics)
startTime, endTime := time.Now().UTC().Add(-time.Hour), time.Now().UTC()
traffic, err := r.DBClient.GetAllLatestObjTraffic(startTime, endTime)
if err != nil {
return fmt.Errorf("failed to get all latest object storage traffic: %w", err)
}
for i := range traffic {
user := traffic[i].User
bucket := traffic[i].Bucket
if _, ok := metrics[user]; !ok {
continue
}
if traffic[i].Time.Before(time.Now().Add(-time.Hour)) {
continue
}
if _, ok := latestObjTrafficSentMetrics[user]; !ok {
latestObjTrafficSentMetrics[user] = objstorage.MetricData{
Sent: make(map[string]int64),
}
}
latestObjTrafficSentMetrics[user].Sent[bucket] = traffic[i].TotalSent
}
r.lastObjectMetrics = latestObjTrafficSentMetrics
}
r.currentObjectMetrics = metrics
logger.Info("success query object storage resource usage", "time", time.Now().Format("2006-01-02 15:04:05"))
}
if r.ObjStorageClient != nil {
r.ObjStorageUserBackupSize = objstorage.GetUserBakFileSize(r.ObjStorageClient)
logger.Info("success query object storage backup size", "time", time.Now().Format("2006-01-02 15:04:05"))
logger.Info("success query object storage usage and traffic metrics", "time", time.Now().Format("2006-01-02 15:04:05"))
}
return nil
}
@@ -457,7 +490,7 @@ func (r *MonitorReconciler) monitorDatabaseBackupUsage(namespace string, resUsed
for i := range backupList.Items {
backup := &backupList.Items[i]
backupRes := resources.NewResourceNamed(backup)
fmt.Printf("backup name: %v, backup size: %v, backupRes: %s \n", backupList.Items[i].Name, backupList.Items[i].Status.TotalSize, backupRes.String())
//fmt.Printf("backup name: %v, backup size: %v, backupRes: %s \n", backupList.Items[i].Name, backupList.Items[i].Status.TotalSize, backupRes.String())
if resUsed[backupRes.String()] == nil {
resNamed[backupRes.String()] = backupRes
resUsed[backupRes.String()] = initResources()
@@ -533,6 +566,41 @@ func (r *MonitorReconciler) monitorObjectStorageUsage(namespace string, resMap m
return nil
}
func (r *MonitorReconciler) monitorObjectStorageTraffic() error {
if r.currentObjectMetrics == nil {
return nil
}
var objTraffic []*types.ObjectStorageTraffic
now := time.Now().UTC()
for user, metric := range r.currentObjectMetrics {
if len(metric.Sent) == 0 {
continue
}
for bucket, m := range metric.Sent {
sent := int64(0)
if r.lastObjectMetrics != nil && r.lastObjectMetrics[user].Sent != nil {
if _, ok := r.lastObjectMetrics[user].Sent[bucket]; ok {
ss := m - r.lastObjectMetrics[user].Sent[bucket]
if ss > 0 {
sent = ss
}
}
}
objTraffic = append(objTraffic, &types.ObjectStorageTraffic{
Time: now,
User: user,
Bucket: bucket,
TotalSent: m,
Sent: sent,
})
}
}
if err := r.DBClient.SaveObjTraffic(objTraffic...); err != nil {
return fmt.Errorf("failed to save object storage traffic: %w", err)
}
return nil
}
func (r *MonitorReconciler) MonitorTrafficUsed(startTime, endTime time.Time) error {
logger.Info("start getTrafficUsed", "startTime", startTime.Format(time.RFC3339), "endTime", endTime.Format(time.RFC3339))
execTime := time.Now().UTC()
@@ -551,9 +619,9 @@ func (r *MonitorReconciler) MonitorTrafficUsed(startTime, endTime time.Time) err
}
func (r *MonitorReconciler) monitorObjectStorageTrafficUsed(startTime, endTime time.Time) error {
buckets, err := objstorage.ListAllObjectStorageBucket(r.ObjStorageClient)
buckets, err := r.DBClient.GetTimeObjBucketBucket(startTime, endTime)
if err != nil {
return fmt.Errorf("failed to list object storage buckets: %w", err)
return fmt.Errorf("failed to get object storage buckets: %w", err)
}
r.Logger.Info("object storage buckets", "buckets len", len(buckets))
wg, _ := errgroup.WithContext(context.Background())
@@ -571,7 +639,7 @@ func (r *MonitorReconciler) monitorObjectStorageTrafficUsed(startTime, endTime t
}
func (r *MonitorReconciler) handlerObjectStorageTrafficUsed(startTime, endTime time.Time, bucket string) error {
bytes, err := objstorage.GetObjectStorageFlow(r.PromURL, bucket, r.ObjectStorageInstance, startTime, endTime)
bytes, err := r.DBClient.HandlerTimeObjBucketSentTraffic(startTime, endTime, bucket)
if err != nil {
return fmt.Errorf("failed to get object storage flow: %w", err)
}
@@ -633,7 +701,6 @@ func (r *MonitorReconciler) handlerTrafficUsed(startTime, endTime time.Time, mon
Time: endTime.Add(-1 * time.Minute),
Type: monitor.Type,
}
r.Logger.Info("monitor traffic used", "monitor", ro)
err = r.DBClient.InsertMonitor(context.Background(), &ro)
if err != nil {
return fmt.Errorf("failed to insert monitor: %w", err)
+4
View File
@@ -185,6 +185,10 @@ func main() {
} else {
reconciler.Logger.Info("minio info not found, please check env: MINIO_ENDPOINT, MINIO_AK, MINIO_SK, MINIO_METRICS_ADDR")
}
err = reconciler.DBClient.CreateTTLTrafficTimeSeries()
if err != nil {
reconciler.Logger.Error(err, "failed to create ttl traffic time series")
}
// timer creates tomorrow's timing table in advance to ensure that tomorrow's table exists
// Execute immediately and then every 24 hours.
time.AfterFunc(time.Until(getNextMidnight()), func() {