You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用Go与Operator SDK通过API调用管理Kubernetes Pod

基于Go Operator SDK实现API驱动的Pod创建/删除控制器

你的需求完全可行,但需要明确Operator SDK的核心定位:它是围绕自定义资源(CRD)+调和循环设计的框架,默认不直接处理HTTP请求。下面提供两种实现方案,分别符合Operator最佳实践和直接HTTP请求的需求。


方案一:用CRD作为中间层(推荐,符合Operator范式)

这种方式将用户的创建/删除请求转化为K8s自定义资源(CR)的增删操作,利用Operator的调和循环完成Pod的生命周期管理,同时复用K8s API Server的认证、授权、审计能力。

步骤1:定义自定义资源(CRD)

创建一个名为PodRequest的CRD,包含用户请求的imageTag、namespace字段,以及用于记录Pod状态的status字段:

// api/v1/podrequest_types.go
package v1

import (
	corev1 "k8s.io/api/core/v1"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)

// PodRequestSpec defines the desired state of PodRequest
type PodRequestSpec struct {
	ImageTag  string `json:"imageTag"`
	Namespace string `json:"namespace"`
}

// PodRequestStatus defines the observed state of PodRequest
type PodRequestStatus struct {
	PodID string           `json:"podID,omitempty"`
	Phase corev1.PodPhase `json:"phase,omitempty"`
}

//+kubebuilder:object:root=true
//+kubebuilder:subresource:status

// PodRequest is the Schema for the podrequests API
type PodRequest struct {
	metav1.TypeMeta   `json:",inline"`
	metav1.ObjectMeta `json:"metadata,omitempty"`

	Spec   PodRequestSpec   `json:"spec,omitempty"`
	Status PodRequestStatus `json:"status,omitempty"`
}

//+kubebuilder:object:root=true

// PodRequestList contains a list of PodRequest
type PodRequestList struct {
	metav1.TypeMeta `json:",inline"`
	metav1.ListMeta `json:"metadata,omitempty"`
	Items           []PodRequest `json:"items"`
}

func init() {
	SchemeBuilder.Register(&PodRequest{}, &PodRequestList{})
}

步骤2:实现控制器调和逻辑

在Reconcile函数中监听PodRequest的变化:

  • 当PodRequest被创建时,生成并创建对应的Pod,将Pod名称(即podId)写入PodRequest的status字段
  • 当PodRequest被删除时,通过OwnerReference自动关联删除Pod
  • 同步Pod状态到PodRequest的status中
// controllers/podrequest_controller.go
package controllers

import (
	"context"

	corev1 "k8s.io/api/core/v1"
	apierrors "k8s.io/apimachinery/pkg/api/errors"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/runtime"
	ctrl "sigs.k8s.io/controller-runtime"
	"sigs.k8s.io/controller-runtime/pkg/client"
	"sigs.k8s.io/controller-runtime/pkg/log"

	podapiv1 "your-operator/api/v1"
)

// PodRequestReconciler reconciles a PodRequest object
type PodRequestReconciler struct {
	client.Client
	Scheme *runtime.Scheme
}

//+kubebuilder:rbac:groups=podapi.example.com,resources=podrequests,verbs=get;list;watch;create;update;patch;delete
//+kubebuilder:rbac:groups=podapi.example.com,resources=podrequests/status,verbs=get;update;patch
//+kubebuilder:rbac:groups=podapi.example.com,resources=podrequests/finalizers,verbs=update
//+kubebuilder:rbac:groups="",resources=pods,verbs=get;list;watch;create;delete;update

// Reconcile is part of the main kubernetes reconciliation loop which aims to
// move the current state of the cluster closer to the desired state.
func (r *PodRequestReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
	log := log.FromContext(ctx)

	// 获取PodRequest实例
	podReq := &podapiv1.PodRequest{}
	if err := r.Get(ctx, req.NamespacedName, podReq); err != nil {
		return ctrl.Result{}, client.IgnoreNotFound(err)
	}

	// 若未关联Pod,创建Pod并更新状态
	if podReq.Status.PodID == "" {
		pod := &corev1.Pod{
			ObjectMeta: metav1.ObjectMeta{
				GenerateName: "podreq-",
				Namespace:    podReq.Spec.Namespace,
				// 设置OwnerReference,让Pod随PodRequest删除而自动删除
				OwnerReferences: []metav1.OwnerReference{
					*metav1.NewControllerRef(podReq, podapiv1.GroupVersion.WithKind("PodRequest")),
				},
			},
			Spec: corev1.PodSpec{
				Containers: []corev1.Container{
					{
						Name:  "app-container",
						Image: podReq.Spec.ImageTag,
					},
				},
			},
		}

		if err := r.Create(ctx, pod); err != nil {
			log.Error(err, "Failed to create Pod")
			return ctrl.Result{}, err
		}

		// 更新PodRequest状态记录PodID
		podReq.Status.PodID = pod.Name
		if err := r.Status().Update(ctx, podReq); err != nil {
			log.Error(err, "Failed to update PodRequest status")
			return ctrl.Result{}, err
		}
		return ctrl.Result{Requeue: true}, nil
	}

	// 同步Pod状态到PodRequest
	targetPod := &corev1.Pod{}
	err := r.Get(ctx, client.ObjectKey{
		Name:      podReq.Status.PodID,
		Namespace: podReq.Spec.Namespace,
	}, targetPod)
	if err != nil {
		if apierrors.IsNotFound(err) {
			podReq.Status.Phase = corev1.PodFailed
			if updateErr := r.Status().Update(ctx, podReq); updateErr != nil {
				return ctrl.Result{}, updateErr
			}
			return ctrl.Result{}, nil
		}
		return ctrl.Result{}, err
	}

	if podReq.Status.Phase != targetPod.Status.Phase {
		podReq.Status.Phase = targetPod.Status.Phase
		if err := r.Status().Update(ctx, podReq); err != nil {
			return ctrl.Result{}, err
		}
	}

	return ctrl.Result{}, nil
}

// SetupWithManager sets up the controller with the Manager.
func (r *PodRequestReconciler) SetupWithManager(mgr ctrl.Manager) error {
	return ctrl.NewControllerManagedBy(mgr).
		For(&podapiv1.PodRequest{}).
		Owns(&corev1.Pod{}).
		Complete(r)
}

步骤3:暴露API给用户

用户可以通过以下方式调用:

  • 直接调用K8s API:发送POST请求到/apis/podapi.example.com/v1/namespaces/{namespace}/podrequests,携带imageTag和namespace字段,返回的status.podID即为Pod的ID;发送DELETE请求到/apis/podapi.example.com/v1/namespaces/{namespace}/podrequests/{podrequest-name}即可删除关联Pod
  • 用Ingress/API网关暴露K8s API,或者写一个轻量代理服务,将用户的REST请求转化为K8s API调用

方案二:在Operator进程中嵌入HTTP服务器

如果需要直接接收用户的HTTP请求而不通过CRD,可以在Operator的主进程中启动一个HTTP服务,在Handler中直接调用K8s客户端创建/删除Pod。

核心代码示例(main.go)

package main

import (
	"encoding/json"
	"fmt"
	"net/http"
	"os"

	corev1 "k8s.io/api/core/v1"
	apierrors "k8s.io/apimachinery/pkg/api/errors"
	metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
	"k8s.io/apimachinery/pkg/runtime"
	clientgoscheme "k8s.io/client-go/kubernetes/scheme"
	ctrl "sigs.k8s.io/controller-runtime"
	"sigs.k8s.io/controller-runtime/pkg/client"
	"sigs.k8s.io/controller-runtime/pkg/log/zap"

	podapiv1 "your-operator/api/v1"
	"your-operator/controllers"
)

var (
	scheme   = runtime.NewScheme()
	setupLog = ctrl.Log.WithName("setup")
)

func init() {
	_ = clientgoscheme.AddToScheme(scheme)
	_ = podapiv1.AddToScheme(scheme)
}

func main() {
	var metricsAddr string
	var enableLeaderElection bool
	var probeAddr string
	ctrl.SetLogger(zap.New(zap.UseDevMode(true)))

	flag.StringVar(&metricsAddr, "metrics-bind-address", ":8080", "The address the metric endpoint binds to.")
	flag.StringVar(&probeAddr, "health-probe-address", ":8081", "The address the health probe endpoint binds to.")
	flag.BoolVar(&enableLeaderElection, "leader-elect", false,
		"Enable leader election for controller manager. "+
			"Enabling this will ensure there is only one active controller manager.")
	flag.Parse()

	mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), ctrl.Options{
		Scheme:                 scheme,
		MetricsBindAddress:     metricsAddr,
		Port:                   9443,
		HealthProbeBindAddress: probeAddr,
		LeaderElection:         enableLeaderElection,
		LeaderElectionID:       "pod-operator-lock.example.com",
	})
	if err != nil {
		setupLog.Error(err, "unable to start manager")
		os.Exit(1)
	}

	// 注册控制器(如果需要保留CRD相关逻辑)
	if err = (&controllers.PodRequestReconciler{
		Client: mgr.GetClient(),
		Scheme: mgr.GetScheme(),
	}).SetupWithManager(mgr); err != nil {
		setupLog.Error(err, "unable to create controller", "controller", "PodRequest")
		os.Exit(1)
	}

	// 启动HTTP服务处理请求
	go func() {
		// 创建Pod接口
		http.HandleFunc("/create-pod", func(w http.ResponseWriter, r *http.Request) {
			if r.Method != http.MethodPost {
				w.WriteHeader(http.StatusMethodNotAllowed)
				fmt.Fprintf(w, "only POST method is supported")
				return
			}

			var reqBody struct {
				ImageTag  string `json:"imageTag"`
				Namespace string `json:"namespace"`
			}
			if err := json.NewDecoder(r.Body).Decode(&reqBody); err != nil {
				w.WriteHeader(http.StatusBadRequest)
				fmt.Fprintf(w, "invalid request body: %v", err)
				return
			}

			pod := &corev1.Pod{
				ObjectMeta: metav1.ObjectMeta{
					GenerateName: "direct-pod-",
					Namespace:    reqBody.Namespace,
				},
				Spec: corev1.PodSpec{
					Containers: []corev1.Container{
						{
							Name:  "app",
							Image: reqBody.ImageTag,
						},
					},
				},
			}

			if err := mgr.GetClient().Create(r.Context(), pod); err != nil {
				w.WriteHeader(http.StatusInternalServerError)
				fmt.Fprintf(w, "failed to create pod: %v", err)
				return
			}

			w.WriteHeader(http.StatusOK)
			json.NewEncoder(w).Encode(map[string]string{"podId": pod.Name})
		})

		// 删除Pod接口
		http.HandleFunc("/delete-pod", func(w http.ResponseWriter, r *http.Request) {
			if r.Method != http.MethodDelete {
				w.WriteHeader(http.StatusMethodNotAllowed)
				fmt.Fprintf(w, "only DELETE method is supported")
				return
			}

			podId := r.URL.Query().Get("podId")
			namespace := r.URL.Query().Get("namespace")
			if podId == "" || namespace == "" {
				w.WriteHeader(http.StatusBadRequest)
				fmt.Fprintf(w, "podId and namespace are required")
				return
			}

			pod := &corev1.Pod{
				ObjectMeta: metav1.ObjectMeta{
					Name:      podId,
					Namespace: namespace,
				},
			}

			if err := mgr.GetClient().Delete(r.Context(), pod); err != nil {
				if apierrors.IsNotFound(err) {
					w.WriteHeader(http.StatusNotFound)
					fmt.Fprintf(w, "pod not found")
					return
				}
				w.WriteHeader(http.StatusInternalServerError)
				fmt.Fprintf(w, "failed to delete pod: %v", err)
				return
			}

			w.WriteHeader(http.StatusOK)
			fmt.Fprintf(w, "pod deleted successfully")
		})

		setupLog.Info("starting HTTP server on :8082")
		if err := http.ListenAndServe(":8082", nil); err != nil {
			setupLog.Error(err, "failed to start HTTP server")
			os.Exit(1)
		}
	}()

	setupLog.Info("starting manager")
	if err := mgr.Start(ctrl.SetupSignalHandler()); err != nil {
		setupLog.Error(err, "problem running manager")
		os.Exit(1)
	}
}

注意事项

  • 该方案需要自行处理认证授权(比如添加API Key验证、集成K8s RBAC),否则存在未授权访问风险
  • 需为Operator的ServiceAccount配置足够的Pod操作权限(参考方案一中的RBAC规则)

内容的提问来源于stack exchange,提问作者Mahesh

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 11:44:59