kubeflow-trainer-v2-tf-cust.../README.md

533 lines
20 KiB
Markdown
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

# 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:<ip>
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 볼륨 마운트 확인