如何将新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
相关产品推荐
相关产品推荐

