mirror of
https://github.com/kubevela/kubevela.git
synced 2026-08-18 03:56:36 +00:00
Feat: init multi-cluster metrics management (#3365)
* Feat: init multi-cluster metrics management Signed-off-by: zhukunshuai <jookunshuai@gmail.com> * Refactor and add metrics collection for metrics-server Signed-off-by: zhukunshuai <jookunshuai@gmail.com> * Modify comment Signed-off-by: zhukunshuai <jookunshuai@gmail.com> * Fix staticcheck Signed-off-by: zhukunshuai <jookunshuai@gmail.com> * Fix some error Signed-off-by: zhukunshuai <jookunshuai@gmail.com> * Fix Signed-off-by: zhukunshuai <jookunshuai@gmail.com>
This commit is contained in:
@@ -82,6 +82,7 @@ require (
|
||||
k8s.io/kube-aggregator v0.22.1
|
||||
k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e
|
||||
k8s.io/kubectl v0.21.0
|
||||
k8s.io/metrics v0.21.0
|
||||
k8s.io/utils v0.0.0-20210802155522-efc7438f0176
|
||||
open-cluster-management.io/api v0.0.0-20210804091127-340467ff6239
|
||||
rsc.io/letsencrypt v0.0.3 // indirect
|
||||
|
||||
@@ -2589,6 +2589,7 @@ k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e h1:KLHHjkdQFomZy8+06csTWZ
|
||||
k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e/go.mod h1:vHXdDvt9+2spS2Rx9ql3I8tycm3H9FDfdUoIuKCefvw=
|
||||
k8s.io/kubectl v0.21.0 h1:WZXlnG/yjcE4LWO2g6ULjFxtzK6H1TKzsfaBFuVIhNg=
|
||||
k8s.io/kubectl v0.21.0/go.mod h1:EU37NukZRXn1TpAkMUoy8Z/B2u6wjHDS4aInsDzVvks=
|
||||
k8s.io/metrics v0.21.0 h1:uwS3CgheLKaw3PTpwhjMswnm/PMqeLbdLH88VI7FMQQ=
|
||||
k8s.io/metrics v0.21.0/go.mod h1:L3Ji9EGPP1YBbfm9sPfEXSpnj8i24bfQbAFAsW0NueQ=
|
||||
k8s.io/utils v0.0.0-20190221042446-c2654d5206da/go.mod h1:8k8uAuAQ0rXslZKaEWd0c3oVhZz7sSzSiPnVZayjIX0=
|
||||
k8s.io/utils v0.0.0-20190801114015-581e00157fb1/go.mod h1:sZAwmy6armz5eXlNoLmJcl4F1QuKu7sr+mFQ0byX7Ew=
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
/*
|
||||
Copyright 2022 The KubeVela 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.
|
||||
*/
|
||||
|
||||
package multicluster
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"k8s.io/klog/v2"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
)
|
||||
|
||||
// metricsMap records the metrics of clusters
|
||||
var metricsMap map[string]*ClusterMetrics
|
||||
|
||||
// ClusterMetricsMgr manage metrics of clusters
|
||||
type ClusterMetricsMgr struct {
|
||||
kubeClient client.Client
|
||||
}
|
||||
|
||||
// ClusterMetricsHelper is the interface that provides operations for cluster metrics
|
||||
type ClusterMetricsHelper interface {
|
||||
Refresh() error
|
||||
}
|
||||
|
||||
// NewClusterMetricsMgr will create a cluster metrics manager
|
||||
func NewClusterMetricsMgr(kubeClient client.Client) (*ClusterMetricsMgr, error) {
|
||||
mgr := &ClusterMetricsMgr{
|
||||
kubeClient: kubeClient,
|
||||
}
|
||||
err := mgr.Refresh()
|
||||
return mgr, err
|
||||
}
|
||||
|
||||
// Refresh will re-collect cluster metrics and refresh cache
|
||||
func (cmm *ClusterMetricsMgr) Refresh() error {
|
||||
clusters, _ := ListVirtualClusters(context.Background(), cmm.kubeClient)
|
||||
m := make(map[string]*ClusterMetrics)
|
||||
|
||||
// retrieves metrics by cluster-gateway
|
||||
for _, cluster := range clusters {
|
||||
isConnected := true
|
||||
clusterInfo, err := GetClusterInfo(context.Background(), cmm.kubeClient, cluster.Name)
|
||||
if err != nil {
|
||||
klog.Warningf("failed to get cluster info of cluster-(%s)", cluster.Name)
|
||||
isConnected = false
|
||||
}
|
||||
clusterUsageMetrics, err := GetClusterMetricsFromMetricsAPI(context.Background(), cmm.kubeClient, cluster.Name)
|
||||
if err != nil {
|
||||
klog.Warningf("failed to request metrics api of cluster-(%s)", cluster.Name)
|
||||
}
|
||||
metrics := &ClusterMetrics{
|
||||
IsConnected: isConnected,
|
||||
ClusterInfo: clusterInfo,
|
||||
ClusterUsageMetrics: clusterUsageMetrics,
|
||||
}
|
||||
m[cluster.Name] = metrics
|
||||
}
|
||||
metricsMap = m
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,155 @@
|
||||
/*
|
||||
Copyright 2022 The KubeVela 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.
|
||||
*/
|
||||
|
||||
package multicluster
|
||||
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"strconv"
|
||||
"testing"
|
||||
|
||||
"gotest.tools/assert"
|
||||
|
||||
"github.com/oam-dev/cluster-gateway/pkg/apis/cluster/v1alpha1"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
"k8s.io/apimachinery/pkg/api/resource"
|
||||
metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1"
|
||||
clusterv1 "open-cluster-management.io/api/cluster/v1"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client/fake"
|
||||
|
||||
"github.com/oam-dev/kubevela/pkg/utils/common"
|
||||
)
|
||||
|
||||
const (
|
||||
NormalClusterName = "normal-cluster"
|
||||
DisconnectedClusterName = "disconnected-cluster"
|
||||
|
||||
NodeName1 = "node-1"
|
||||
NodeName2 = "node-2"
|
||||
)
|
||||
|
||||
func TestRefresh(t *testing.T) {
|
||||
fakeClient := NewFakeClient(fake.NewClientBuilder().
|
||||
WithScheme(common.Scheme).
|
||||
WithRuntimeObjects(FakeManagedCluster("managed-cluster")).
|
||||
WithObjects(FakeSecret(NormalClusterName), FakeSecret(DisconnectedClusterName)).
|
||||
Build())
|
||||
|
||||
normalCluster := fake.NewClientBuilder().
|
||||
WithScheme(common.Scheme).
|
||||
WithObjects(FakeNode(NodeName1, "8", strconv.FormatInt(16*1024*1024*1024, 10)),
|
||||
FakeNode(NodeName2, "7", strconv.FormatInt(32*1024*1024*1024, 10)),
|
||||
FakeNodeMetrics(NodeName1, "4", strconv.FormatInt(8*1024*1024*1024, 10)),
|
||||
FakeNodeMetrics(NodeName2, "1", strconv.FormatInt(3*1024*1024*1024, 10))).
|
||||
Build()
|
||||
|
||||
disconnectedCluster := &disconnectedClient{}
|
||||
|
||||
fakeClient.AddCluster(NormalClusterName, normalCluster)
|
||||
fakeClient.AddCluster(DisconnectedClusterName, disconnectedCluster)
|
||||
|
||||
mgr, err := NewClusterMetricsMgr(fakeClient)
|
||||
assert.NilError(t, err)
|
||||
|
||||
err = mgr.Refresh()
|
||||
assert.NilError(t, err)
|
||||
|
||||
clusters, err := ListVirtualClusters(context.Background(), fakeClient)
|
||||
assert.NilError(t, err)
|
||||
|
||||
for _, cluster := range clusters {
|
||||
assertClusterMetrics(t, &cluster)
|
||||
}
|
||||
|
||||
disCluster, err := GetVirtualCluster(context.Background(), fakeClient, DisconnectedClusterName)
|
||||
assert.NilError(t, err)
|
||||
assertClusterMetrics(t, disCluster)
|
||||
|
||||
norCluster, err := GetVirtualCluster(context.Background(), fakeClient, NormalClusterName)
|
||||
assert.NilError(t, err)
|
||||
assertClusterMetrics(t, norCluster)
|
||||
}
|
||||
|
||||
func assertClusterMetrics(t *testing.T, cluster *VirtualCluster) {
|
||||
metrics := cluster.Metrics
|
||||
switch cluster.Name {
|
||||
case DisconnectedClusterName:
|
||||
assert.Equal(t, metrics.IsConnected, false)
|
||||
assert.Assert(t, metrics.ClusterInfo == nil)
|
||||
assert.Assert(t, metrics.ClusterUsageMetrics == nil)
|
||||
case NormalClusterName:
|
||||
assert.Equal(t, metrics.IsConnected, true)
|
||||
|
||||
assert.Assert(t, resource.MustParse("15").Equal(metrics.ClusterInfo.CPUCapacity))
|
||||
assert.Assert(t, resource.MustParse(strconv.FormatInt(48*1024*1024*1024, 10)).Equal(metrics.ClusterInfo.MemoryCapacity))
|
||||
assert.Assert(t, resource.MustParse("15").Equal(metrics.ClusterInfo.CPUAllocatable))
|
||||
assert.Assert(t, resource.MustParse(strconv.FormatInt(48*1024*1024*1024, 10)).Equal(metrics.ClusterInfo.MemoryAllocatable))
|
||||
|
||||
assert.Assert(t, resource.MustParse("5").Equal(metrics.ClusterUsageMetrics.CPUUsage))
|
||||
assert.Assert(t, resource.MustParse(strconv.FormatInt(11*1024*1024*1024, 10)).Equal(metrics.ClusterUsageMetrics.MemoryUsage))
|
||||
}
|
||||
}
|
||||
|
||||
func FakeNodeMetrics(name string, cpu string, memory string) *metricsV1beta1api.NodeMetrics {
|
||||
nodeMetrics := &metricsV1beta1api.NodeMetrics{}
|
||||
nodeMetrics.Name = name
|
||||
nodeMetrics.Usage = corev1.ResourceList{
|
||||
corev1.ResourceCPU: resource.MustParse(cpu),
|
||||
corev1.ResourceMemory: resource.MustParse(memory),
|
||||
}
|
||||
return nodeMetrics
|
||||
}
|
||||
|
||||
func FakeNode(name string, cpu string, memory string) *corev1.Node {
|
||||
node := &corev1.Node{}
|
||||
node.Name = name
|
||||
node.Status = corev1.NodeStatus{
|
||||
Allocatable: map[corev1.ResourceName]resource.Quantity{
|
||||
corev1.ResourceCPU: resource.MustParse(cpu),
|
||||
corev1.ResourceMemory: resource.MustParse(memory),
|
||||
},
|
||||
Capacity: map[corev1.ResourceName]resource.Quantity{
|
||||
corev1.ResourceCPU: resource.MustParse(cpu),
|
||||
corev1.ResourceMemory: resource.MustParse(memory),
|
||||
},
|
||||
}
|
||||
return node
|
||||
}
|
||||
|
||||
func FakeSecret(name string) *corev1.Secret {
|
||||
secret := &corev1.Secret{}
|
||||
secret.Name = name
|
||||
secret.Labels = map[string]string{
|
||||
v1alpha1.LabelKeyClusterCredentialType: "ServiceAccountToken",
|
||||
}
|
||||
return secret
|
||||
}
|
||||
|
||||
func FakeManagedCluster(name string) *clusterv1.ManagedCluster {
|
||||
managedCluster := &clusterv1.ManagedCluster{}
|
||||
managedCluster.Name = name
|
||||
return managedCluster
|
||||
}
|
||||
|
||||
type disconnectedClient struct {
|
||||
client.Client
|
||||
}
|
||||
|
||||
func (cli *disconnectedClient) List(ctx context.Context, list client.ObjectList, opts ...client.ListOption) error {
|
||||
return errors.New("no such host")
|
||||
}
|
||||
@@ -19,6 +19,8 @@ package multicluster
|
||||
import (
|
||||
"context"
|
||||
|
||||
metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
corev1 "k8s.io/api/core/v1"
|
||||
storagev1 "k8s.io/api/storage/v1"
|
||||
@@ -26,6 +28,13 @@ import (
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
)
|
||||
|
||||
// ClusterMetrics describes the metrics of a cluster
|
||||
type ClusterMetrics struct {
|
||||
IsConnected bool
|
||||
ClusterInfo *ClusterInfo
|
||||
ClusterUsageMetrics *ClusterUsageMetrics
|
||||
}
|
||||
|
||||
// ClusterInfo describes the basic information of a cluster
|
||||
type ClusterInfo struct {
|
||||
Nodes *corev1.NodeList
|
||||
@@ -40,6 +49,12 @@ type ClusterInfo struct {
|
||||
StorageClasses *storagev1.StorageClassList
|
||||
}
|
||||
|
||||
// ClusterUsageMetrics describes the usage metrics of a cluster
|
||||
type ClusterUsageMetrics struct {
|
||||
CPUUsage resource.Quantity
|
||||
MemoryUsage resource.Quantity
|
||||
}
|
||||
|
||||
// GetClusterInfo retrieves current cluster info from cluster
|
||||
func GetClusterInfo(_ctx context.Context, k8sClient client.Client, clusterName string) (*ClusterInfo, error) {
|
||||
ctx := ContextWithClusterName(_ctx, clusterName)
|
||||
@@ -81,3 +96,20 @@ func GetClusterInfo(_ctx context.Context, k8sClient client.Client, clusterName s
|
||||
StorageClasses: storageClasses,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// GetClusterMetricsFromMetricsAPI retrieves current cluster metrics based on GetNodeMetricsFromMetricsAPI
|
||||
func GetClusterMetricsFromMetricsAPI(ctx context.Context, k8sClient client.Client, clusterName string) (*ClusterUsageMetrics, error) {
|
||||
nodeMetricsList := metricsV1beta1api.NodeMetricsList{}
|
||||
if err := k8sClient.List(ContextWithClusterName(ctx, clusterName), &nodeMetricsList); err != nil {
|
||||
return nil, errors.Wrapf(err, "failed to list node metrics")
|
||||
}
|
||||
var memoryUsage, cpuUsage resource.Quantity
|
||||
for _, nm := range nodeMetricsList.Items {
|
||||
cpuUsage.Add(*nm.Usage.Cpu())
|
||||
memoryUsage.Add(*nm.Usage.Memory())
|
||||
}
|
||||
return &ClusterUsageMetrics{
|
||||
CPUUsage: cpuUsage,
|
||||
MemoryUsage: memoryUsage,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -46,6 +46,7 @@ type VirtualCluster struct {
|
||||
EndPoint string
|
||||
Accepted bool
|
||||
Labels map[string]string
|
||||
Metrics *ClusterMetrics
|
||||
}
|
||||
|
||||
// NewVirtualClusterFromSecret extract virtual cluster from cluster secret
|
||||
@@ -68,6 +69,7 @@ func NewVirtualClusterFromSecret(secret *corev1.Secret) (*VirtualCluster, error)
|
||||
EndPoint: endpoint,
|
||||
Accepted: true,
|
||||
Labels: labels,
|
||||
Metrics: metricsMap[secret.Name],
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -82,6 +84,7 @@ func NewVirtualClusterFromManagedCluster(managedCluster *clusterv1.ManagedCluste
|
||||
EndPoint: "-",
|
||||
Accepted: managedCluster.Spec.HubAcceptsClient,
|
||||
Labels: managedCluster.GetLabels(),
|
||||
Metrics: metricsMap[managedCluster.Name],
|
||||
}, nil
|
||||
}
|
||||
|
||||
|
||||
@@ -56,6 +56,8 @@ import (
|
||||
"sigs.k8s.io/controller-runtime/pkg/client/config"
|
||||
"sigs.k8s.io/yaml"
|
||||
|
||||
metricsV1beta1api "k8s.io/metrics/pkg/apis/metrics/v1beta1"
|
||||
|
||||
oamcore "github.com/oam-dev/kubevela/apis/core.oam.dev"
|
||||
"github.com/oam-dev/kubevela/apis/core.oam.dev/common"
|
||||
"github.com/oam-dev/kubevela/apis/core.oam.dev/v1beta1"
|
||||
@@ -96,6 +98,7 @@ func init() {
|
||||
_ = ocmclusterv1.Install(Scheme)
|
||||
_ = ocmworkv1.Install(Scheme)
|
||||
_ = clustergatewayapi.AddToScheme(Scheme)
|
||||
_ = metricsV1beta1api.AddToScheme(Scheme)
|
||||
// +kubebuilder:scaffold:scheme
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user