使用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
相关产品推荐
相关产品推荐

