# Kubeflow Trainer v2 - TensorFlow Custom Runtime Kubeflow Trainer v2에서 TensorFlow 분산학습을 위한 Custom ClusterTrainingRuntime 생성 및 Volcano 연동 가이드. ## 배경 Kubeflow Trainer v2 (v1alpha1)에는 다음 런타임만 기본 제공된다: | Runtime | Framework Label | mlPolicy 필드 | 엔트리포인트 주입 | |---------|----------------|---------------|-------------------| | torch-distributed | `torch` | `mlPolicy.torch` | torchrun + MASTER_ADDR/RANK 등 자동 설정 | | deepspeed-distributed | `deepspeed` | `mlPolicy.mpi` | MPI (OpenMPI) + SSH 자동 설정 | | mlx-distributed | `mlx` | `mlPolicy.mpi` | MPI (OpenMPI) + SSH 자동 설정 | | torchtune-* | `torchtune` | `mlPolicy.torch` | tune run + rdzv 자동 설정 | **TensorFlow runtime은 없다.** 따라서 Custom ClusterTrainingRuntime을 직접 생성해야 한다. ## 실행 결과 요약 2노드 × 8 GPU (Tesla V100-SXM2-32GB) 분산학습 성공: ``` 환경: - 노드: 2개 - GPU: 노드당 8장 (총 16장) - GPU 모델: Tesla V100-SXM2-32GB - 통신: NCCL + InfiniBand + GPUDirect RDMA - 이미지: nvcr.io/nvidia/tensorflow:24.03-tf2-py3 학습 결과: - 데이터셋: CIFAR-10 (로컬 hostPath) - 모델: CNN (Conv2D + BatchNorm + Dropout) - Strategy: MultiWorkerMirroredStrategy (NCCL) - num_replicas_in_sync: 16 - Global Batch Size: 64 × 16 = 1024 - Epochs: 10 - Test Loss: 1.6169 - Test Accuracy: 44.59% NCCL 초기화 로그: NCCL INFO Connected all trees NCCL INFO 16 coll channels, 0 nvls channels, ... NCCL INFO comm 0x... rank 0 nranks 16 ... busId ... COMPLETE → 16개 rank 전부 NCCL Init COMPLETE, GPUDirect RDMA 활성화 확인 ``` ## PyTorch vs TensorFlow Runtime 핵심 차이 ``` ┌─────────────────────────────────────────────────────────────────────┐ │ PyTorch Runtime (built-in) │ │ │ │ mlPolicy.torch → 컨트롤러가 자동으로: │ │ ├── torchrun 엔트리포인트 주입 │ │ ├── MASTER_ADDR, MASTER_PORT 설정 │ │ ├── RANK, LOCAL_RANK, WORLD_SIZE 설정 │ │ └── rdzv (rendezvous) 자동 구성 │ │ │ │ numNodes → replica별 별도 Job 생성 │ │ Pod hostname: {trainjob}-node-{replica_index}-0 │ │ │ │ → 학습 코드에서 분산환경 설정 불필요 │ └─────────────────────────────────────────────────────────────────────┘ ┌─────────────────────────────────────────────────────────────────────┐ │ TensorFlow Runtime (custom) │ │ │ │ mlPolicy에 framework 없음 → 컨트롤러가 주입하는 것 없음 │ │ │ │ numNodes → 단일 Job에 completions=numNodes 로 매핑 │ │ Pod hostname: {trainjob}-node-0-{completion_index} │ │ completion_index = worker_index │ │ │ │ 학습 스크립트에서 자체 구성: │ │ ├── Pod hostname 파싱 → TrainJob 이름, worker index 추출 │ │ ├── DNS 프로빙으로 워커 자동 감지 │ │ ├── TF_CONFIG 환경변수 자동 구성 (import tensorflow 전에!) │ │ └── tf.distribute.MultiWorkerMirroredStrategy 사용 │ │ │ │ 필수 설정: │ │ ├── network.publishNotReadyAddresses: true (DNS 조기 해석) │ │ └── 워커 수는 DNS 프로빙으로 자동 감지 (수동 설정 불필요) │ └─────────────────────────────────────────────────────────────────────┘ ``` ## Pod 네이밍 패턴 (중요) mlPolicy에 framework(torch/mpi) 설정 여부에 따라 Pod 네이밍이 달라진다: ``` ┌─ mlPolicy.torch 사용 (PyTorch) ──────────────────────────────────┐ │ numNodes: 2 → replicatedJob replicas=2 → Job 2개 생성 │ │ │ │ Job: {trainjob}-node-0 → Pod: {trainjob}-node-0-0 │ │ Job: {trainjob}-node-1 → Pod: {trainjob}-node-1-0 │ │ ↑ replica_index ↑ completion=0 │ └────────────────────────────────────────────────────────────────────┘ ┌─ mlPolicy에 framework 없음 (TensorFlow custom) ─────────────────┐ │ numNodes: 2 → 단일 Job에 completions=2 로 매핑 │ │ │ │ Job: {trainjob}-node-0 → Pod: {trainjob}-node-0-0 │ │ Pod: {trainjob}-node-0-1 │ │ ↑ replica=0 ↑ completion_index │ │ │ │ completion_index = worker_index │ └────────────────────────────────────────────────────────────────────┘ ``` 이 차이를 모르면 DNS 프로빙 패턴이 잘못되어 워커를 1개만 발견하는 문제가 발생한다. ## TF_CONFIG 자동 구성 원리 ### 1. Pod hostname에서 정보 추출 ``` hostname: tf-distributed-training-node-0-1 │ │ │ └─ completion_index (= worker_index = 1) └─── replica_index (항상 0) ``` hostname 파싱 정규식: ```python match = re.match(r'^(.+)-node-(\d+)-(\d+)$', hostname) job_name = match.group(1) # tf-distributed-training replica_index = match.group(2) # 0 (항상) worker_index = int(match.group(3)) # completion_index = worker_index ``` ### 2. Headless Service DNS로 워커 주소 구성 ``` Worker 0: tf-distributed-training-node-0-0.tf-distributed-training.{namespace}.svc.cluster.local:12345 Worker 1: tf-distributed-training-node-0-1.tf-distributed-training.{namespace}.svc.cluster.local:12345 ``` ### 3. DNS 프로빙으로 워커 수 자동 감지 ``` discover_workers() 동작: node-0-0 DNS 조회 → 성공 → workers에 추가 node-0-1 DNS 조회 → 성공 → workers에 추가 node-0-2 DNS 조회 → 실패 → 탐색 종료 → num_workers = 2 안정성: 2개 이상 발견 + 연속 2회 동일 결과 → 확정 ``` ### 4. TF_CONFIG 생성 ```json { "cluster": { "worker": [ "tf-distributed-training-node-0-0.tf-distributed-training.{namespace}.svc.cluster.local:12345", "tf-distributed-training-node-0-1.tf-distributed-training.{namespace}.svc.cluster.local:12345" ] }, "task": {"type": "worker", "index": 0} } ``` ### 5. TF_CONFIG 설정 순서 (중요) ```python # 반드시 이 순서를 지켜야 한다: setup_tf_config() # 1. TF_CONFIG 환경변수 설정 import tensorflow as tf # 2. 그 다음 tensorflow import # 이유: # TF는 import 시점에 GPU를 감지하고, TF_CONFIG가 있으면 gRPC 서버를 시작한다. # import 후에 TF_CONFIG를 설정하면 gRPC 서버가 잘못된 상태로 시작되어 # "different incarnation" 에러가 발생한다. ``` ## TFJob (Trainer v1) vs TrainJob (Trainer v2) 비교 ### chief/worker 패턴 비교 Trainer v1의 TFJob은 `chief`와 `worker` role을 명시적으로 분리했다: ```yaml # TFJob (Trainer v1) - chief/worker 명시적 분리 apiVersion: kubeflow.org/v1 kind: TFJob spec: tfReplicaSpecs: Chief: replicas: 1 template: ... # chief 전용 설정 Worker: replicas: 3 template: ... # worker 설정 PS: replicas: 2 template: ... # Parameter Server (선택) ``` TrainJob (Trainer v2)은 단일 role 타입(`node`) 기반이다: ```yaml # TrainJob (Trainer v2) - 단일 role apiVersion: trainer.kubeflow.org/v1alpha1 kind: TrainJob spec: trainer: numNodes: 4 # 모든 노드가 동일한 role ``` ### 구현 가능성 분석 | 기능 | TFJob (v1) | TrainJob (v2) | 비고 | |------|-----------|---------------|------| | chief/worker 분리 | `tfReplicaSpecs.Chief/Worker` | 미지원 (단일 role) | TrainJob은 replicatedJobs로 부분 구현 가능 | | Parameter Server | `tfReplicaSpecs.PS` | 미지원 | PS 전략 자체가 레거시 | | Evaluator | `tfReplicaSpecs.Evaluator` | 미지원 | - | | worker index 기반 분기 | role별로 자동 | `worker_index == 0`으로 직접 구현 | 실질적으로 동일한 효과 | | 분산 전략 | PS Strategy / MirroredStrategy | MultiWorkerMirroredStrategy | v2는 AllReduce 기반 | ### worker_index == 0 패턴 (현재 구현) TFJob의 chief가 하던 역할(체크포인트 저장, 로그 출력, 평가 결과 출력 등)을 `worker_index == 0`으로 분기하여 동일하게 구현: ```python # worker_index == 0 이 사실상 chief 역할 if worker_index == 0: print(f"Test Loss: {loss:.4f}") print(f"Test Accuracy: {accuracy:.4f}") # 체크포인트 저장, TensorBoard 로그 등도 여기서 ``` **주의**: `model.evaluate()`는 모든 워커가 함께 호출해야 한다 (collective operation). Worker 0만 호출하면 다른 워커가 heartbeat timeout으로 실패한다. ```python # 올바른 패턴: 모든 워커가 evaluate 참여, 출력만 worker 0 loss, accuracy = model.evaluate(test_dataset, verbose=eval_verbose) # 모든 워커 if worker_index == 0: # 출력만 worker 0 print(f"Test Accuracy: {accuracy:.4f}") ``` ### replicatedJobs를 이용한 multi-role 구현 (참고) TrainJob의 Runtime에서 replicatedJobs를 여러 개 정의하면 role 분리가 이론적으로 가능하다: ```yaml # 이론적 구현 (권장하지 않음) replicatedJobs: - name: chief # chief role template: spec: template: spec: containers: - name: node env: - name: TF_ROLE value: "chief" - name: worker # worker role template: spec: template: spec: containers: - name: node env: - name: TF_ROLE value: "worker" ``` 하지만 TrainJob API의 `trainer` 필드는 단일 role을 가정하므로, 여러 replicatedJobs에 대한 `numNodes`, `resourcesPerNode` 등의 매핑이 명확하지 않다. **결론**: 현재 TrainJob에서는 `worker_index == 0` 패턴으로 chief 역할을 대체하는 것이 가장 실용적이다. TFJob의 PS(Parameter Server) 전략 자체가 레거시이므로, MultiWorkerMirroredStrategy 기반 AllReduce가 현대적 표준이다. ## CIFAR-10 데이터 로딩 `tf.keras.datasets.cifar10.load_data()`는 외부 URL에서 다운로드를 시도한다. 에어갭 환경이나 빠른 로딩을 위해 로컬 hostPath를 마운트하여 pickle 파일에서 직접 로드한다. ```yaml # Runtime에서 hostPath 볼륨 마운트 volumeMounts: - name: cifar-data mountPath: /workspace/data/cifar-10-batches-py volumes: - name: cifar-data hostPath: path: /home/ubuntu/cifar-10-batches-py # 노드에 미리 준비 type: Directory ``` ## 통신 구조 (InfiniBand + NCCL) TensorFlow도 PyTorch와 동일하게 NCCL + InfiniBand를 사용할 수 있다. NCCL은 NVIDIA에서 제공하는 프레임워크 독립적인 GPU 간 통신 라이브러리이므로, TF/PyTorch 구분 없이 동일하게 동작한다. ``` PyTorch: torchrun → torch.distributed → NCCL → InfiniBand/GPUDirect RDMA TensorFlow: gRPC (조정) → MultiWorkerMirroredStrategy → NCCL → InfiniBand/GPUDirect RDMA ``` 동일한 NCCL 환경변수가 양쪽 모두에 적용된다: ```yaml env: - name: NCCL_IB_DISABLE value: "0" # InfiniBand 활성화 - name: NCCL_SOCKET_IFNAME value: "net" # OOB 통신용 인터페이스 - name: NCCL_DEBUG value: "INFO" # NCCL 디버그 로그 - name: NCCL_DEBUG_SUBSYS value: "INIT,NET,IB" # NCCL 서브시스템별 디버그 ``` 실제 학습 시 확인된 NCCL 초기화 로그: ``` NCCL INFO NET/IB : Using [0]mlx5_0:1/RoCE ... ; OOB net0: NCCL INFO Connected all trees NCCL INFO 16 coll channels, 0 nvls channels, ... NCCL INFO comm 0x... rank 0 nranks 16 ... COMPLETE ``` ## 파일 구조 ``` . ├── README.md # 이 문서 ├── tensorflow-custom-runtime.yaml # TF Custom Runtime (기본, Volcano 미포함) └── tensorflow-volcano-trainjob-integration.yaml # TF + Volcano 전체 통합 (ConfigMap + Queue + Runtime + TrainJob) ``` ## 적용 방법 ### 1단계: 기본 TensorFlow Runtime만 적용 Volcano 없이 기본 TF 런타임만 필요한 경우: ```bash kubectl apply -f tensorflow-custom-runtime.yaml ``` 확인: ```bash kubectl get clustertrainingruntime # tensorflow-distributed 가 목록에 나타나야 함 ``` ### 2단계: Volcano 연동 전체 적용 Volcano gang scheduling + 토폴로지 인식 스케줄링이 포함된 전체 버전: ```bash # ConfigMap (학습 스크립트) + Queue + Runtime + TrainJob 한번에 적용 kubectl apply -f tensorflow-volcano-trainjob-integration.yaml ``` 확인: ```bash # Runtime 확인 kubectl get clustertrainingruntime tensorflow-distributed-volcano # TrainJob 상태 확인 kubectl get trainjob tf-distributed-training # Pod 상태 확인 (namespace 지정 필요) kubectl get pods -l jobset.sigs.k8s.io/jobset-name=tf-distributed-training -n {namespace} # Volcano PodGroup 확인 kubectl get podgroup # 로그 확인 (worker 0) kubectl -n {namespace} logs -l job-name=tf-distributed-training-node-0 -f ``` ### 3단계: TrainJob만 변경하여 재실행 Runtime은 유지하고 TrainJob만 변경할 경우: ```bash # 기존 TrainJob 삭제 kubectl delete trainjob tf-distributed-training -n {namespace} # 수정 후 재적용 kubectl apply -f tensorflow-volcano-trainjob-integration.yaml ``` ## 주의사항 ### 워커 수 자동 감지 (DNS 프로빙) 학습 스크립트의 `discover_workers()` 함수가 DNS 프로빙으로 워커 수를 자동 감지한다. `NUM_WORKERS` 같은 환경변수를 수동으로 관리할 필요가 없으므로, TrainJob에서 `numNodes`만 변경하면 된다: ```yaml apiVersion: trainer.kubeflow.org/v1alpha1 kind: TrainJob spec: runtimeRef: name: tensorflow-distributed-volcano trainer: numNodes: 4 # 이것만 변경하면 됨. 스크립트가 자동으로 4개 워커 감지. ``` 비교: | | PyTorch Runtime | TensorFlow Custom Runtime | |---|---|---| | 워커 수 감지 | 컨트롤러가 `WORLD_SIZE` 자동 주입 | 스크립트가 DNS 프로빙으로 자동 감지 | | numNodes 변경 시 | 추가 작업 없음 | 추가 작업 없음 | ### InfiniBand 설정 클러스터에 InfiniBand가 없는 경우 다음 항목을 제거해야 한다: ```yaml # 제거 대상: metadata: annotations: k8s.v1.cni.cncf.io/networks: hostdevice-net # ← 제거 resources: requests: nvidia.com/hostdev: "2" # ← 제거 limits: nvidia.com/hostdev: "2" # ← 제거 ``` ### network.publishNotReadyAddresses TensorFlow의 `MultiWorkerMirroredStrategy`는 시작 시 모든 worker에 gRPC 연결을 시도한다. 모든 Pod이 Ready 상태가 되기 전에도 DNS가 해석되어야 하므로 이 설정이 **필수**다: ```yaml spec: template: spec: network: publishNotReadyAddresses: true # 반드시 true ``` PyTorch runtime에서는 torchrun의 rdzv(rendezvous) 메커니즘이 이를 처리하므로 불필요하다. ### TrainJob env 오버라이드 Runtime에 정의된 환경변수는 TrainJob의 `trainer.env`에서 오버라이드할 수 있다: ```yaml apiVersion: trainer.kubeflow.org/v1alpha1 kind: TrainJob spec: runtimeRef: name: tensorflow-distributed-volcano trainer: env: - name: TF_WORKER_PORT # Runtime 기본값 "12345" → 오버라이드 value: "23456" - name: EPOCHS # Runtime 기본값 "10" → 오버라이드 value: "50" - name: NCCL_DEBUG # NCCL 디버그 레벨 변경 value: "WARN" ``` ### model.evaluate()는 모든 워커가 참여해야 함 MultiWorkerMirroredStrategy는 collective operation 기반이다. `model.evaluate()`를 worker 0만 호출하면 다른 워커가 heartbeat timeout으로 실패한다. ```python # 잘못된 패턴 (worker 0만 evaluate) if worker_index == 0: loss, accuracy = model.evaluate(test_dataset) # ← 다른 워커 timeout 발생 # 올바른 패턴 (모든 워커가 evaluate, 출력만 worker 0) eval_verbose = 2 if worker_index == 0 else 0 loss, accuracy = model.evaluate(test_dataset, verbose=eval_verbose) # 모든 워커 참여 if worker_index == 0: print(f"Test Accuracy: {accuracy:.4f}") # 출력만 worker 0 ``` ## 트러블슈팅 ### "different incarnation" 에러 ``` E tensorflow/...coordination_service.cc: /job:worker/replica:0/task:0 unexpectedly tried to connect with a different incarnation. It has likely restarted. ``` **원인**: `TF_CONFIG`가 `import tensorflow` 이후에 설정됨 **해결**: 학습 스크립트에서 `setup_tf_config()`를 `import tensorflow` 전에 호출 ### Pod 네이밍 패턴 불일치로 워커 감지 실패 ``` [Discovery] attempt 1/60: found 1 workers, retrying in 5s... (영원히 1개만 발견) ``` **원인**: DNS 프로빙 패턴이 `node-{i}-0` (PyTorch 패턴)으로 되어있으나, 실제는 `node-0-{i}` (completion 패턴) **해결**: mlPolicy에 framework가 없으면 `numNodes`가 단일 Job의 `completions`로 매핑됨을 이해하고, DNS 패턴을 `{trainjob}-node-0-{i}`로 수정 ### DNS 프로빙에서 워커를 1개만 발견 (정상 대기) ``` [Discovery] attempt 1/60: found 1 workers, retrying in 5s... [Discovery] attempt 2/60: found 2 workers, retrying in 5s... [Discovery] Discovered 2 workers via DNS ``` **원인**: 다른 워커 Pod이 아직 시작되지 않았거나 DNS 전파 지연 **해결**: Volcano gang scheduling 사용 시 모든 Pod이 동시에 스케줄링되지만, DNS 전파에 시간이 걸릴 수 있다. `discover_workers()`가 연속 2회 동일 결과를 확인한 후 진행하므로 정상적으로 기다린다. ### model.evaluate() heartbeat timeout ``` W tensorflow/core/distributed_runtime/...coordination_service_agent.cc: Heartbeat timeout from ... ``` **원인**: Worker 0만 `model.evaluate()` 호출, 다른 워커는 이미 종료 **해결**: MultiWorkerMirroredStrategy는 collective operation이므로, 모든 워커가 `model.evaluate()`에 참여해야 함. 출력만 `worker_index == 0`에서 처리. ### CIFAR-10 데이터 다운로드 시도 ``` Downloading data from https://www.cs.toronto.edu/~kriz/cifar-10-python.tar.gz ``` **원인**: `tf.keras.datasets.cifar10.load_data()` 사용 또는 hostPath 볼륨 미마운트 **해결**: `load_cifar10_local()` 함수 사용 + cifar-data hostPath 볼륨 마운트 확인