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

如何将新Kubernetes Service关联至运行中的Job及Pod

问题场景与需求
  • 原有链路:Svc-1 创建并监控 Job-1,Job-1 生成 Pod-1,Svc-1 接收Job的ACTIVE/READY/FAIL/SUCCESS等状态并记录
  • 故障情况:Svc-1 下线后,新启动的 Svc-2 无法关联已在运行的 Job-1,无法接收其及对应Pod的事件
  • 目标:将 Svc-2 关联至运行中的 Job-1,达成链路 Svc-2 -> Job-1 -> Pod-1

现有配置与代码

Service配置文件

apiVersion: v1
kind: Service
metadata:
  name: svc-dev
  labels:
    app: svc-dev
spec:
  type: ClusterIP
  selector:
    app: svc-dev
  ports:
    - name: http
      protocol: TCP
      port: 8080

创建Kubernetes Job的Fabric8 Java代码

private JobDto createJob(KubernetesClient client, JobDto jobDto) throws JobExecutionException{
    try {
        MixedOperation<Job, JobList, ScalableResource<Job>> v1Job = client.batch().v1().jobs();
        Job kubernetesJob = v1Job.load(resource.getInputStream()).get();
        log.info("{} job has been loaded", jobType);

        //Set name for job
        kubernetesJob.getMetadata().setName(jobDto.getJobName());
        log.info("Set name for job {}",jobDto.getJobName());

        //Add label to job
        String executionType = jobDto.getExecutionType().name();
        Map<String, String> labels = MapUtils.of(LABEL1, jobType, LABEL2, executionType);
        kubernetesJob.getMetadata().setLabels(labels);
        log.info("Added labels for job {}", kubernetesJob.getMetadata().getLabels());

        //Add label to pod
        kubernetesJob.getSpec().getTemplate().getMetadata().setLabels(labels);
        log.info("Added labels for job pod template");

        //Add node selector role
        kubernetesJob.getSpec().getTemplate().getSpec().setNodeSelector(getNodeSelectorRole(jobType));
        log.info("Added NodeSelector for job pod template");

        //Add service account to job
        kubernetesJob.getSpec().getTemplate().getSpec().setServiceAccountName(config.getServiceAccount());
        log.info("Added ServiceAccountName for job pod template");

        //Add container configuration and volume is also mounted
        kubernetesJob.getSpec().getTemplate().getSpec().setContainers(getContainers(jobDto));
        log.info("Added Containers for job pod template");

        //Add Volume
        kubernetesJob.getSpec().getTemplate().getSpec().setVolumes(getVolumes());
        log.info("Added Volumes for job pod template");

        // Create Kubernetes Job
        v1Job.inNamespace(config.getNamespace()).resource(kubernetesJob).create();
        log.info("Created job with name {}",jobDto.getJobName());
        return jobDto;

    }catch(Exception exception ){
        log.error("Failed to create job", exception);
        throw new JobExecutionException("Job creation failed", exception);
    }
}

Deployment配置文件(Helm模板)

apiVersion: apps/v1
kind: Deployment
metadata:
  name: {{ .Values.application.name }}
  labels:
    app: {{ .Values.application.name }}
spec:
  replicas: {{ .Values.application.replicaCount }}
  selector:
    matchLabels:
      app: {{ .Values.application.name }}

  template:
    metadata:
      labels:
        app: {{ .Values.application.name }}
    spec:
      nodeSelector:
        role: {{ .Values.application.nodeSelectorRole }}
      serviceAccountName: {{ .Values.application.serviceAccountName }}
      serviceAccount: {{ .Values.application.serviceAccountName }}
      volumes:
        - name: {{ .Values.application.persistentVolume }}
          persistentVolumeClaim:
            claimName: {{ .Values.application.persistentVolumeClaim }}
      containers:
        - name: {{ .Values.application.name }}
          image: {{ .Values.application.image }}
          imagePullPolicy: {{ .Values.application.imagePullPolicy }}
          ports:
            - name: {{ .Values.application.portname }}
              containerPort: {{ .Values.application.port }}
          volumeMounts:
            - name: {{ .Values.application.persistentVolume }}
              mountPath: {{ .Values.application.containerMountPath }}

实现方案

基于Spring Boot + Fabric8 Kubernetes Client,按以下步骤完成:

1. 配置RBAC权限

确保Svc-2使用的ServiceAccount拥有监听Job和Pod的权限,添加以下RBAC资源到Helm模板:

apiVersion: rbac.authorization.k8s.io/v1
kind: ClusterRole
metadata:
  name: job-monitor-role
rules:
- apiGroups: ["batch"]
  resources: ["jobs", "jobs/status"]
  verbs: ["get", "list", "watch"]
- apiGroups: [""]
  resources: ["pods", "pods/status"]
  verbs: ["get", "list", "watch"]
---
apiVersion: rbac.authorization.k8s.io/v1
kind: RoleBinding
metadata:
  name: job-monitor-binding
  namespace: {{ .Values.namespace }}
subjects:
- kind: ServiceAccount
  name: {{ .Values.application.serviceAccountName }}
  namespace: {{ .Values.namespace }}
roleRef:
  kind: ClusterRole
  name: job-monitor-role
  apiGroup: rbac.authorization.k8s.io

2. 实现Job与Pod的事件监听逻辑

在Svc-2中添加监听组件,完成初始化关联已运行Job和实时监听事件的功能:

@Component
public class JobMonitorListener implements ApplicationListener<ApplicationReadyEvent> {

    private final KubernetesClient kubernetesClient;
    private final JobStatusRepository jobStatusRepository;
    private final String namespace;
    private final String targetJobName;

    public JobMonitorListener(KubernetesClient kubernetesClient, 
                              JobStatusRepository jobStatusRepository,
                              @Value("${kubernetes.namespace}") String namespace,
                              @Value("${target.job.name}") String targetJobName) {
        this.kubernetesClient = kubernetesClient;
        this.jobStatusRepository = jobStatusRepository;
        this.namespace = namespace;
        this.targetJobName = targetJobName;
    }

    @Override
    public void onApplicationEvent(ApplicationReadyEvent event) {
        // 查询已运行的Job-1
        Job existingJob = kubernetesClient.batch().v1().jobs()
                .inNamespace(namespace)
                .withName(targetJobName)
                .get();

        if (existingJob != null) {
            // 同步当前状态到数据库
            syncJobStatusToDb(existingJob);
            // 注册Job状态监听器
            registerJobListener(existingJob);
            // 注册对应Pod的监听器
            registerPodListener(existingJob.getMetadata().getLabels());
        }
    }

    private void syncJobStatusToDb(Job job) {
        JobStatus status = new JobStatus();
        status.setJobName(job.getMetadata().getName());
        status.setNamespace(job.getMetadata().getNamespace());
        status.setStatus(job.getStatus().getConditions().stream()
                .map(JobCondition::getType)
                .findFirst()
                .orElse("UNKNOWN"));
        status.setActive(job.getStatus().getActive() != null ? job.getStatus().getActive() : 0);
        status.setSucceeded(job.getStatus().getSucceeded() != null ? job.getStatus().getSucceeded() : 0);
        status.setFailed(job.getStatus().getFailed() != null ? job.getStatus().getFailed() : 0);
        jobStatusRepository.save(status);
    }

    private void registerJobListener(Job job) {
        kubernetesClient.batch().v1().jobs()
                .inNamespace(namespace)
                .withName(job.getMetadata().getName())
                .watch(new Watcher<>() {
                    @Override
                    public void eventReceived(Action action, Job job) {
                        syncJobStatusToDb(job);
                        log.info("Job {} status updated: {}", job.getMetadata().getName(), action);
                    }

                    @Override
                    public void onClose(WatcherException cause) {
                        log.error("Job watcher closed", cause);
                        registerJobListener(job);
                    }
                });
    }

    private void registerPodListener(Map<String, String> jobLabels) {
        kubernetesClient.pods()
                .inNamespace(namespace)
                .withLabels(jobLabels)
                .watch(new Watcher<>() {
                    @Override
                    public void eventReceived(Action action, Pod pod) {
                        String jobName = pod.getMetadata().getLabels().get(LABEL1);
                        JobStatus jobStatus = jobStatusRepository.findByJobName(jobName);
                        if (jobStatus != null) {
                            jobStatus.setPodStatus(pod.getStatus().getPhase().name());
                            jobStatusRepository.save(jobStatus);
                        }
                        log.info("Pod {} status updated: {}", pod.getMetadata().getName(), action);
                    }

                    @Override
                    public void onClose(WatcherException cause) {
                        log.error("Pod watcher closed", cause);
                        registerPodListener(jobLabels);
                    }
                });
    }
}

3. 配置参数

在Svc-2的application.yml中添加必要配置:

kubernetes:
  namespace: your-namespace
target:
  job:
    name: Job-1

4. 验证与测试

  • 部署Svc-2后,检查日志确认监听器已成功注册
  • 查看数据库,确认Job-1的当前状态已被同步
  • 触发Job-1的状态变化(如Pod完成),验证状态是否自动更新到数据库

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 11:40:51