mirror of
https://github.com/zalando/postgres-operator.git
synced 2026-10-02 13:02:50 +02:00
Kube cluster upgrade
This commit is contained in:
committed by
Murat Kabilov
parent
1dbf259c76
commit
eba23279c8
@@ -41,6 +41,7 @@ type Controller struct {
|
||||
|
||||
postgresqlInformer cache.SharedIndexInformer
|
||||
podInformer cache.SharedIndexInformer
|
||||
nodesInformer cache.SharedIndexInformer
|
||||
podCh chan spec.PodEvent
|
||||
|
||||
clusterEventQueues []*cache.FIFO // [workerID]Queue
|
||||
@@ -111,6 +112,7 @@ func (c *Controller) initOperatorConfig() {
|
||||
func (c *Controller) initController() {
|
||||
c.initClients()
|
||||
c.initOperatorConfig()
|
||||
c.initSharedInformers()
|
||||
|
||||
c.logger.Infof("config: %s", c.opConfig.MustMarshal())
|
||||
|
||||
@@ -128,6 +130,23 @@ func (c *Controller) initController() {
|
||||
c.config.InfrastructureRoles = infraRoles
|
||||
}
|
||||
|
||||
c.clusterEventQueues = make([]*cache.FIFO, c.opConfig.Workers)
|
||||
c.workerLogs = make(map[uint32]ringlog.RingLogger, c.opConfig.Workers)
|
||||
for i := range c.clusterEventQueues {
|
||||
c.clusterEventQueues[i] = cache.NewFIFO(func(obj interface{}) (string, error) {
|
||||
e, ok := obj.(spec.ClusterEvent)
|
||||
if !ok {
|
||||
return "", fmt.Errorf("could not cast to ClusterEvent")
|
||||
}
|
||||
|
||||
return queueClusterKey(e.EventType, e.UID), nil
|
||||
})
|
||||
}
|
||||
|
||||
c.apiserver = apiserver.New(c, c.opConfig.APIPort, c.logger.Logger)
|
||||
}
|
||||
|
||||
func (c *Controller) initSharedInformers() {
|
||||
// Postgresqls
|
||||
c.postgresqlInformer = cache.NewSharedIndexInformer(
|
||||
&cache.ListWatch{
|
||||
@@ -162,31 +181,35 @@ func (c *Controller) initController() {
|
||||
DeleteFunc: c.podDelete,
|
||||
})
|
||||
|
||||
c.clusterEventQueues = make([]*cache.FIFO, c.opConfig.Workers)
|
||||
c.workerLogs = make(map[uint32]ringlog.RingLogger, c.opConfig.Workers)
|
||||
for i := range c.clusterEventQueues {
|
||||
c.clusterEventQueues[i] = cache.NewFIFO(func(obj interface{}) (string, error) {
|
||||
e, ok := obj.(spec.ClusterEvent)
|
||||
if !ok {
|
||||
return "", fmt.Errorf("could not cast to ClusterEvent")
|
||||
}
|
||||
|
||||
return queueClusterKey(e.EventType, e.UID), nil
|
||||
})
|
||||
// Kubernetes Nodes
|
||||
nodeLw := &cache.ListWatch{
|
||||
ListFunc: c.nodeListFunc,
|
||||
WatchFunc: c.nodeWatchFunc,
|
||||
}
|
||||
|
||||
c.apiserver = apiserver.New(c, c.opConfig.APIPort, c.logger.Logger)
|
||||
c.nodesInformer = cache.NewSharedIndexInformer(
|
||||
nodeLw,
|
||||
&v1.Node{},
|
||||
constants.QueueResyncPeriodNode,
|
||||
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc})
|
||||
|
||||
c.nodesInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
|
||||
AddFunc: c.nodeAdd,
|
||||
UpdateFunc: c.nodeUpdate,
|
||||
DeleteFunc: c.nodeDelete,
|
||||
})
|
||||
}
|
||||
|
||||
// Run starts background controller processes
|
||||
func (c *Controller) Run(stopCh <-chan struct{}, wg *sync.WaitGroup) {
|
||||
c.initController()
|
||||
|
||||
wg.Add(4)
|
||||
wg.Add(5)
|
||||
go c.runPodInformer(stopCh, wg)
|
||||
go c.runPostgresqlInformer(stopCh, wg)
|
||||
go c.clusterResync(stopCh, wg)
|
||||
go c.apiserver.Run(stopCh, wg)
|
||||
go c.kubeNodesInformer(stopCh, wg)
|
||||
|
||||
for i := range c.clusterEventQueues {
|
||||
wg.Add(1)
|
||||
@@ -212,3 +235,9 @@ func (c *Controller) runPostgresqlInformer(stopCh <-chan struct{}, wg *sync.Wait
|
||||
func queueClusterKey(eventType spec.EventType, uid types.UID) string {
|
||||
return fmt.Sprintf("%s-%s", eventType, uid)
|
||||
}
|
||||
|
||||
func (c *Controller) kubeNodesInformer(stopCh <-chan struct{}, wg *sync.WaitGroup) {
|
||||
defer wg.Done()
|
||||
|
||||
c.nodesInformer.Run(stopCh)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,162 @@
|
||||
package controller
|
||||
|
||||
import (
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/labels"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/watch"
|
||||
"k8s.io/client-go/pkg/api/v1"
|
||||
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/cluster"
|
||||
"github.com/zalando-incubator/postgres-operator/pkg/util"
|
||||
)
|
||||
|
||||
func (c *Controller) nodeListFunc(options metav1.ListOptions) (runtime.Object, error) {
|
||||
opts := metav1.ListOptions{
|
||||
Watch: options.Watch,
|
||||
ResourceVersion: options.ResourceVersion,
|
||||
TimeoutSeconds: options.TimeoutSeconds,
|
||||
}
|
||||
|
||||
return c.KubeClient.Nodes().List(opts)
|
||||
}
|
||||
|
||||
func (c *Controller) nodeWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
|
||||
opts := metav1.ListOptions{
|
||||
Watch: options.Watch,
|
||||
ResourceVersion: options.ResourceVersion,
|
||||
TimeoutSeconds: options.TimeoutSeconds,
|
||||
}
|
||||
|
||||
return c.KubeClient.Nodes().Watch(opts)
|
||||
}
|
||||
|
||||
func (c *Controller) nodeAdd(obj interface{}) {
|
||||
node, ok := obj.(*v1.Node)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
c.logger.Debugf("new node has been added: %q (%s)", util.NameFromMeta(node.ObjectMeta), node.Spec.ProviderID)
|
||||
}
|
||||
|
||||
func (c *Controller) nodeUpdate(prev, cur interface{}) {
|
||||
nodePrev, ok := prev.(*v1.Node)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
nodeCur, ok := cur.(*v1.Node)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
if util.MapContains(nodeCur.Labels, map[string]string{"master": "true"}) {
|
||||
return
|
||||
}
|
||||
|
||||
if nodePrev.Spec.Unschedulable && util.MapContains(nodePrev.Labels, c.opConfig.EOLNodeLabel) ||
|
||||
!nodeCur.Spec.Unschedulable || !util.MapContains(nodeCur.Labels, c.opConfig.EOLNodeLabel) {
|
||||
return
|
||||
}
|
||||
|
||||
c.logger.Infof("node %q became unschedulable and has EOL labels: %q", util.NameFromMeta(nodeCur.ObjectMeta),
|
||||
c.opConfig.EOLNodeLabel)
|
||||
|
||||
opts := metav1.ListOptions{
|
||||
LabelSelector: labels.Set(c.opConfig.ClusterLabels).String(),
|
||||
}
|
||||
podList, err := c.KubeClient.Pods(c.opConfig.Namespace).List(opts)
|
||||
if err != nil {
|
||||
c.logger.Errorf("could not fetch list of the pods: %v", err)
|
||||
return
|
||||
}
|
||||
|
||||
nodePods := make([]*v1.Pod, 0)
|
||||
for i, pod := range podList.Items {
|
||||
if pod.Spec.NodeName == nodeCur.Name {
|
||||
nodePods = append(nodePods, &podList.Items[i])
|
||||
}
|
||||
}
|
||||
|
||||
clusters := make(map[*cluster.Cluster]bool)
|
||||
masterPods := make(map[*v1.Pod]*cluster.Cluster)
|
||||
replicaPods := make(map[*v1.Pod]*cluster.Cluster)
|
||||
movedPods := 0
|
||||
for _, pod := range nodePods {
|
||||
podName := util.NameFromMeta(pod.ObjectMeta)
|
||||
|
||||
role, ok := pod.Labels[c.opConfig.PodRoleLabel]
|
||||
if !ok {
|
||||
c.logger.Warningf("could not move pod %q: pod has no role", podName)
|
||||
continue
|
||||
}
|
||||
|
||||
clusterName := c.podClusterName(pod)
|
||||
|
||||
c.clustersMu.RLock()
|
||||
cl, ok := c.clusters[clusterName]
|
||||
c.clustersMu.RUnlock()
|
||||
if !ok {
|
||||
c.logger.Warningf("could not move pod %q: pod does not belong to a known cluster", podName)
|
||||
continue
|
||||
}
|
||||
|
||||
movedPods++
|
||||
|
||||
if !clusters[cl] {
|
||||
clusters[cl] = true
|
||||
}
|
||||
|
||||
if cluster.PostgresRole(role) == cluster.Master {
|
||||
masterPods[pod] = cl
|
||||
} else {
|
||||
replicaPods[pod] = cl
|
||||
}
|
||||
}
|
||||
|
||||
for cl := range clusters {
|
||||
cl.Lock()
|
||||
}
|
||||
|
||||
for pod, cl := range masterPods {
|
||||
podName := util.NameFromMeta(pod.ObjectMeta)
|
||||
|
||||
if err := cl.MigrateMasterPod(podName); err != nil {
|
||||
c.logger.Errorf("could not move master pod %q: %v", podName, err)
|
||||
movedPods--
|
||||
}
|
||||
}
|
||||
|
||||
for pod, cl := range replicaPods {
|
||||
podName := util.NameFromMeta(pod.ObjectMeta)
|
||||
|
||||
if err := cl.MigrateReplicaPod(podName, nodeCur.Name); err != nil {
|
||||
c.logger.Errorf("could not move replica pod %q: %v", podName, err)
|
||||
movedPods--
|
||||
}
|
||||
}
|
||||
|
||||
for cl := range clusters {
|
||||
cl.Unlock()
|
||||
}
|
||||
|
||||
totalPods := len(nodePods)
|
||||
|
||||
c.logger.Infof("%d/%d pods have been moved out from the %q node",
|
||||
movedPods, totalPods, util.NameFromMeta(nodeCur.ObjectMeta))
|
||||
|
||||
if leftPods := totalPods - movedPods; leftPods > 0 {
|
||||
c.logger.Warnf("could not move %d/%d pods from the %q node",
|
||||
leftPods, totalPods, util.NameFromMeta(nodeCur.ObjectMeta))
|
||||
}
|
||||
}
|
||||
|
||||
func (c *Controller) nodeDelete(obj interface{}) {
|
||||
node, ok := obj.(*v1.Node)
|
||||
if !ok {
|
||||
return
|
||||
}
|
||||
|
||||
c.logger.Debugf("node has been deleted: %q (%s)", util.NameFromMeta(node.ObjectMeta), node.Spec.ProviderID)
|
||||
}
|
||||
@@ -11,12 +11,7 @@ import (
|
||||
)
|
||||
|
||||
func (c *Controller) podListFunc(options metav1.ListOptions) (runtime.Object, error) {
|
||||
var labelSelector string
|
||||
var fieldSelector string
|
||||
|
||||
opts := metav1.ListOptions{
|
||||
LabelSelector: labelSelector,
|
||||
FieldSelector: fieldSelector,
|
||||
Watch: options.Watch,
|
||||
ResourceVersion: options.ResourceVersion,
|
||||
TimeoutSeconds: options.TimeoutSeconds,
|
||||
@@ -26,12 +21,7 @@ func (c *Controller) podListFunc(options metav1.ListOptions) (runtime.Object, er
|
||||
}
|
||||
|
||||
func (c *Controller) podWatchFunc(options metav1.ListOptions) (watch.Interface, error) {
|
||||
var labelSelector string
|
||||
var fieldSelector string
|
||||
|
||||
opts := metav1.ListOptions{
|
||||
LabelSelector: labelSelector,
|
||||
FieldSelector: fieldSelector,
|
||||
Watch: options.Watch,
|
||||
ResourceVersion: options.ResourceVersion,
|
||||
TimeoutSeconds: options.TimeoutSeconds,
|
||||
|
||||
@@ -189,7 +189,7 @@ func (c *Controller) processEvent(event spec.ClusterEvent) {
|
||||
lg.Infoln("update of the cluster started")
|
||||
|
||||
if !clusterFound {
|
||||
lg.Warnln("cluster does not exist")
|
||||
lg.Warningln("cluster does not exist")
|
||||
return
|
||||
}
|
||||
c.curWorkerCluster.Store(event.WorkerID, cl)
|
||||
|
||||
Reference in New Issue
Block a user