如何在Kubernetes Operator中监听特定DaemonSet所属的Pod实例?
Kubernetes Operator监听特定DaemonSet生成的Pod事件
场景描述
我有一个Kubernetes Operator,会创建多个DaemonSet,这些DaemonSet再生成Pod,层级关系如下:
- operator → daemonset1 → pod1_1
- operator → daemonset2 → {pod2_1, pod2_2}
现在需要捕获daemonset2的Pod发生变化的事件(比如旧Pod被替换为新Pod),要求:
- 仅监听与Operator同命名空间下的Pod
- 只关注由特定DaemonSet(daemonset2)创建的Pod
- 不监听所有DaemonSet或Pod
当前尝试的问题
Watches(&source.Kind{Type: &appsv1.DaemonSet{}}):会监听所有命名空间的DaemonSet,不符合需求(注:原表述中误写为监听Pod,实际该代码监听的是DaemonSet)Owns(&appsv1.DaemonSet{}):仅按资源类型过滤,无法针对特定DaemonSet实例EnqueueRequestsFromMapFunc方案:尝试后收到大量无关Reconcile请求,甚至来自其他命名空间,且每个请求被触发两次
尝试的代码片段
func (r *MyTest1Reconciler) SetupWithManager(mgr ctrl.Manager) error { log.Info("MapFunc ", "a.GetNamespace():", a.GetNamespace(), "a.GetName(): ", a.GetName()) mapFn := handler.MapFunc( func(a client.Object) []reconcile.Request { return []reconcile.Request{ {NamespacedName: types.NamespacedName{ Name: a.GetName(), Namespace: "test", }}, } }) p := predicate.Funcs{ // 尝试过滤无关事件 CreateFunc: func(e event.CreateEvent) bool { namespace := e.Object.GetNamespace() name := e.Object.GetName() ok := ((namespace == "test") && (strings.Contains(name, "pod2"))) return ok }, } mgr.GetFieldIndexer().IndexField() return ctrl.NewControllerManagedBy(mgr). For(&mytest1v1.MyTest1{}). Watches(&source.Kind{Type: &appsv1.DaemonSet{}}, handler.EnqueueRequestsFromMapFunc(mapFn), builder.WithPredicates(p)). Complete(r) }
错误日志输出
2022-09-12T02:43:59.735-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-cluster-csi-drivers", "Name: ": "aws-ebs-csi-driver-node"} 2022-09-12T02:43:59.779-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-cluster-node-tuning-operator", "a.GetName(): ": "tuned"} 2022-09-12T02:43:59.779-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-cluster-node-tuning-operator", "a.GetName(): ": "tuned"} 2022-09-12T02:43:59.779-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-cluster-node-tuning-operator", "Name: ": "tuned"} 2022-09-12T02:43:59.779-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-cluster-node-tuning-operator", "Name: ": "tuned"} 2022-09-12T02:43:59.884-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-monitoring", "Name: ": "node-exporter"} 2022-09-12T02:43:59.885-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-monitoring", "Name: ": "node-exporter"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-image-registry", "a.GetName(): ": "node-ca"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-image-registry", "a.GetName(): ": "node-ca"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-image-registry", "Name: ": "node-ca"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-multus", "a.GetName(): ": "network-metrics-daemon"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-multus", "a.GetName(): ": "network-metrics-daemon"} 2022-09-12T02:43:59.931-0700 INFO controller_xcrypt Reconcile: {"Namespace:": "openshift-multus", "Name: ": "network-metrics-daemon"} 2022-09-12T02:43:59.935-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-dns", "a.GetName(): ": "node-resolver"} 2022-09-12T02:43:59.935-0700 INFO controller_xcrypt MapFunc {"a.GetNamespace():": "openshift-dns", "a.GetName(): ": "node-resolver"}
解决方案
核心思路
要监听特定DaemonSet生成的Pod,需要通过Pod的OwnerReference关联到目标DaemonSet,同时结合谓词过滤和索引加速查询。
步骤1:为Pod的OwnerReference建立索引
在SetupWithManager中为Pod的OwnerReference建立索引,方便快速查询属于特定DaemonSet的Pod:
err := mgr.GetFieldIndexer().IndexField(context.TODO(), &corev1.Pod{}, "spec.ownerReferences.name", func(rawObj client.Object) []string { pod := rawObj.(*corev1.Pod) var names []string for _, owner := range pod.OwnerReferences { names = append(names, owner.Name) } return names }) if err != nil { return err }
步骤2:配置监听Pod并添加精准过滤谓词
直接监听Pod资源,通过谓词过滤出:
- 处于目标命名空间
- OwnerReference指向daemonset2
- 覆盖所有事件类型(Create/Update/Delete)
// 定义谓词过滤规则 podPredicate := predicate.Funcs{ CreateFunc: func(e event.CreateEvent) bool { return filterPod(e.Object.(*corev1.Pod)) }, UpdateFunc: func(e event.UpdateEvent) bool { return filterPod(e.ObjectNew.(*corev1.Pod)) }, DeleteFunc: func(e event.DeleteEvent) bool { return filterPod(e.Object.(*corev1.Pod)) }, } // 过滤函数 func filterPod(pod *corev1.Pod) bool { // 检查命名空间 if pod.Namespace != "test" { return false } // 检查OwnerReference是否指向daemonset2 for _, owner := range pod.OwnerReferences { if owner.Kind == "DaemonSet" && owner.Name == "daemonset2" { return true } } return false }
步骤3:配置Watches并关联到Operator的自定义资源
将Pod的事件映射到Operator的自定义资源(MyTest1)进行Reconcile:
mapFn := handler.MapFunc(func(a client.Object) []reconcile.Request { // 替换为你的自定义资源名称,若有多个CR实例,需根据Pod关联的DaemonSet匹配对应CR return []reconcile.Request{ { NamespacedName: types.NamespacedName{ Name: "your-cr-name", Namespace: "test", }, }, } }) // 组装控制器 return ctrl.NewControllerManagedBy(mgr). For(&mytest1v1.MyTest1{}). Watches(&source.Kind{Type: &corev1.Pod{}}, handler.EnqueueRequestsFromMapFunc(mapFn), builder.WithPredicates(podPredicate)). Complete(r)
额外说明:解决重复Reconcile问题
日志中每个请求触发两次的原因通常是Kubernetes Informers会发送Added和Updated事件(即使是新创建的资源),可以在Update谓词中添加判断,仅当Pod的状态或关键字段变化时才触发:
UpdateFunc: func(e event.UpdateEvent) bool { oldPod := e.ObjectOld.(*corev1.Pod) newPod := e.ObjectNew.(*corev1.Pod) // 仅当Pod的Phase变化时触发 return oldPod.Status.Phase != newPod.Status.Phase && filterPod(newPod) },
内容的提问来源于stack exchange,提问作者Victor Kim
相关产品推荐
相关产品推荐

