Kubeflow 파이프라인 사용
파이프라인 실습 예제를 통해 실험 그리고 실행을 생성하고, 예측 모델을 학습해 봅니다.
지원 도구
| 도구 | 버전 | 설명 |
|---|---|---|
| KF Pipelines | 2.4.1 | - 간소화된 인터페이스 제공하며 신속한 실험 및 반복 가능한 기계 학습 워크플로우 도구 - 파라미터 튜닝, 실험 관리, 모델 버전 관리 등의 작업 지원 |
Step 1. 시작하기 전에
주요 개념
Kubeflow Pipeline(KFP)은 Docker 컨테이너 기반의 머신러닝 워크플로우를 정의하고 자동화할 수 있는 프레임워크입니다. Python SDK를 사용하여 파이프라인을 정의하고, KFP 백엔드를 통해 실행할 수 있습니다.
| 주요 개념 용어 | 설명 |
|---|---|
| 파이프라인(Pipeline) | 머신러닝 전체 워크플로우를 정의하는 단위로, 여러 개의 컴포넌트로 구성 |
| 컴포넌트(Component) | 워크플로우의 한 단계를 수행하는 실행 단위로 데이터 전처리, 모델 학습 등 각 작업이 해당 |
| 실험(Experiment) | 작성한 파이프라인을 다양한 설정으로 반복 실행하는 논리적 워크스페이스 |
| 런(Run) | Experiment 내에서 수행되는 단일 파이프라인 실행 인스턴스 |
| 아티팩트(Artifact) | 각 컴포넌트가 실행된 후 생성되는 출력 결과 (예: 모델 파일, 평가 점수, 로그 등) |
| TrainJob | Kubeflow Trainer를 통해 생성되는 분산 학습 작업 |
| InferenceService | KServe를 통해 생성되는 모델 서빙 서비스 |
파이프라인 구성
Kubeflow Pipeline(KFP)의 Python SDK를 활용해 JupyterLab 환경에서 파이프라인을 구성하고 실행합니다.
먼저 아래 예제 스크립트를 다운로드 후 노트북 환경에서 실행합니다.
(노트북 환경 구성은 Kubeflow를 이용한 Jupyter Notebook 환경 구성 문서를 확인해 주세요.)
- 실습 예제 스크립트 다운로드 : build_pipeline_for_deploy_kanana_natl_gpu
예제 스크립트 실행 후 아래 과정들을 차례대로 수행합니다.
Step 2. 환경 설정
1. 기본 설정
스크립트를 노트북에서 실행 후 아래 내용을 참고하여 필요한 모듈을 import하고 기본 설정을 수행합니다.
import os
import uuid
from kakaocloud_kbm import KbmPipelineClient
from kfp import kubernetes
import kfp.dsl as dsl
# KBM Kubeflow Pipeline 클라이언트 초기화
os.environ["KUBEFLOW_HOST"] = "https://nipagpu.kakaocloud.com"
os.environ["KUBEFLOW_USERNAME"] = "your-username@kakaoenterprise.com"
os.environ["KUBEFLOW_PASSWORD"] = "your-password"
client = KbmPipelineClient()
# 파이프라인 변수 설정
# 현재 노트북이 실행되는 Kubernetes 네임스페이스 추출
KBM_NAMESPACE = os.environ['NB_PREFIX'].split('/')[2]
# 컴포넌트 파일 저장 경로 설정
COMPONENT_PATH = 'components'
SERVE_ENPOINT_PATH = os.path.join(COMPONENT_PATH, 'kanana_finetune')
# 고유 작업 ID 생성 (각 파이프라인 실행을 구분하기 위함)
TASK_UUID = uuid.uuid1().hex[:8]
# 리소스 이름 생성
PVC_NAME = f"kanana-ft-pvc-{TASK_UUID}" # PersistentVolumeClaim 이름
MODEL_NAME = f"kanana-model-{TASK_UUID}" # 모델 이름
KSERVE_ISVC_NAME = f"kanana-isvc-{TASK_UUID}" # 모델 서빙 API 이름
EPOCH_NUM = 10 # 학습 에포크 수
print(f"Model Name: {MODEL_NAME}")
print(f"KServe InferenceService Name: {KSERVE_ISVC_NAME}")
print(f"Model PVC Name: {PVC_NAME}")
2. 컴포넌트 작성
KFP의 컴포넌트는 dsl 패키지의 데코레이터를 사용하여 작성할 수 있습니다. 각 컴포넌트는 하나의 함수로 구성되며, 해당 함수 내에서 필요한 라이브러리는 반드시 함수 내부에서 import해야 합니다.
2.1. 데이터 다운로드 컴포넌트
Object Storage에서 학습 데이터를 다운로드하는 컴포넌트입니다.
@dsl.component(
packages_to_install=['requests'],
base_image='python:3.11'
)
def download_dataset(kc_kbm_os_train_url: str):
"""
Object Storage에서 학습 데이터를 다운로드하는 컴포넌트
Args:
kc_kbm_os_train_url: 학습 데이터 CSV 파일의 Object Storage URL
"""
import os
from requests import get
def download(url, dist_dir, file_name=None):
"""URL에서 파일을 다운로드하여 지정된 디렉터리에 저장"""
if not file_name:
file_name = url.split('/')[-1]
file_path = os.path.join(dist_dir, file_name)
with open(file_path, "wb") as file:
response = get(url)
response.raise_for_status()
file.write(response.content)
print(f"Downloaded: {file_name} to {dist_dir}")
# PVC 마운트 경로
pvc_data_path = "/data"
# 기본 URL이 제공되지 않은 경우 샘플 데이터 URL 사용
if not kc_kbm_os_train_url:
kc_kbm_os_train_url = 'https://objectstorage.kr-central-2.kakaocloud.com/v1/c11fcba415bd4314b595db954e4d4422/public/tutorial/kubeflow/kubeflow-tensorboard/data/sample_train_data.csv'
# 학습 데이터 다운로드
download(kc_kbm_os_train_url, pvc_data_path, "sample_train_data.csv")
# 다운로드된 파일 목록 확인
print(f"Downloaded files in {pvc_data_path}:")
print(os.listdir(pvc_data_path))
2-2. 모델 파인튜닝 컴포넌트
Kubeflow Trainer를 사용하여 Kanana 모델을 LoRA 기반으로 파인튜닝하는 컴포넌트입니다.
- PEFT (Parameter-Efficient Fine-Tuning) 라이브러리의 LoRA 기법 사용
- Alpaca 형식의 프롬프트 템플릿 적용
- GPU 메모리 효율적인 학습 수행
@dsl.component(
packages_to_install=['kubeflow'],
install_kfp_package=True,
base_image='python:3.11',
output_component_file=f'{SERVE_ENPOINT_PATH}/train_component.yaml'
)
def finetune_kanana_model(
train_job_id: dsl.Output[dsl.Artifact],
epoch_num: str,
namespace: str,
pvc_name: str,
job_name: str,
):
"""
Kanana 모델을 LoRA 기반으로 파인튜닝하는 컴포넌트
Args:
train_job_id: TrainJob ID를 저장할 아티팩트 (출력)
epoch_num: 학습 에포크 수
namespace: Kubernetes 네임스페이스
pvc_name: 데이터 저장용 PVC 이름
job_name: TrainJob 이름
"""
from kubeflow.trainer import TrainerClient, CustomTrainer
import os
import time
def finetune_kanana(model_name: str, epoch_num: str, pvc_data_path: str):
"""
TrainJob Pod에서 실행될 학습 함수
LoRA 기반 파인튜닝을 수행합니다.
"""
from transformers import (
AutoModelForCausalLM,
AutoTokenizer,
TrainingArguments,
Trainer,
TrainerCallback,
)
from peft import LoraConfig, get_peft_model
from datasets import Dataset
import torch
import os
# 작업 디렉터리 설정
os.chdir("/")
# GPU 환경 설정
os.environ["NVIDIA_VISIBLE_DEVICES"] = "0"
os.environ["CUDA_VISIBLE_DEVICES"] = "0"
print(f"PyTorch version: {torch.__version__}")
print(f"CUDA available: {torch.cuda.is_available()}")
# 1. 모델 및 토크나이저 로드
print("=" * 80)
print("Step 1: Loading LLM Model and Tokenizer")
print("=" * 80)
tokenizer = AutoTokenizer.from_pretrained(model_name, padding_side="left")
# Llama 타입 모델은 pad_token을 eos_token으로 설정 필요
tokenizer.pad_token = tokenizer.eos_token
# 기본 모델 로드 (bfloat16으로 메모리 효율성 향상)
base_model = AutoModelForCausalLM.from_pretrained(
model_name,
torch_dtype=torch.bfloat16,
trust_remote_code=True,
device_map="auto"
)
# 2. LoRA 설정 및 적용
print("=" * 80)
print("Step 2: Setting up LoRA Configuration")
print("=" * 80)
lora_config = LoraConfig(
r=8, # LoRA rank (낮을수록 파라미터 수 감소)
lora_alpha=32, # LoRA alpha (스케일링 팩터)
lora_dropout=0.1, # Dropout 비율
target_modules=["q_proj", "k_proj", "v_proj"], # 적용할 모듈
task_type="CAUSAL_LM",
)
model = get_peft_model(base_model, lora_config)
model.print_trainable_parameters() # 학습 가능한 파라미터 수 출력
# 3. 데이터셋 로드 및 전처리
print("=" * 80)
print("Step 3: Loading and Processing Dataset")
print("=" * 80)
train_data_path = f"{pvc_data_path}/sample_train_data.csv"
dataset = Dataset.from_csv(train_data_path)
# Alpaca 형식 프롬프트 템플릿 적용
def formatting_prompts_func(examples):
"""Alpaca 형식으로 프롬프트 포맷팅"""
alpaca_prompt = """Below is an instruction that describes a task, paired with an input that provides further context. Write a response that appropriately completes the request.
### Instruction:
{}
### Input:
{}
### Response:
{}"""
instructions = examples["instruction"]
inputs = examples["input"]
outputs = examples["output"]
eos_token = tokenizer.eos_token
texts = []
for instruction, input_text, output in zip(instructions, inputs, outputs):
# EOS 토큰 추가 필수 (없으면 생성이 무한 반복될 수 있음)
text = alpaca_prompt.format(instruction, input_text, output) + eos_token
texts.append(text)
return {"text": texts}
# 프롬프트 포맷팅 적용
dataset = dataset.map(formatting_prompts_func, batched=True)
# 불필요한 칼럼 제거 (CSV 인덱스 칼럼 등)
if 'Unnamed: 0' in dataset.column_names:
dataset = dataset.remove_columns(['Unnamed: 0'])
# 토크나이징
def tokenize_function(examples):
"""텍스트를 토큰으로 변환"""
tokens = tokenizer(examples["text"], padding=True, return_tensors="pt")
tokens["labels"] = tokens["input_ids"] # Language modeling을 위한 labels
return tokens
dataset = dataset.map(tokenize_function, batched=True, remove_columns=["text"])
print("Dataset processing complete")
# 4. 학습 설정 및 Trainer 초기화
print("=" * 80)
print("Step 4: Setting up Trainer")
print("=" * 80)
class TrainingCallback(TrainerCallback):
"""학습 진행 상황 로깅용 콜백"""
def on_log(self, args, state, control, logs=None, **kwargs):
if logs:
print(f"Step {state.global_step}: {logs}")
trainer = Trainer(
model=model,
train_dataset=dataset,
args=TrainingArguments(
per_device_train_batch_size=2,
gradient_accumulation_steps=4, # 실제 배치 크기 = 2 * 4 = 8
warmup_steps=5,
max_steps=60, # 빠른 테스트를 위한 스텝 수
learning_rate=2e-4,
bf16=True, # bfloat16 사용으로 메모리 절약
logging_steps=1,
weight_decay=0.01,
lr_scheduler_type="linear",
seed=1234,
output_dir="outputs",
report_to="none" # 외부 로깅 서비스 사용 안 함
),
callbacks=[TrainingCallback()],
)
# GPU 메모리 상태 확인
gpu_stats = torch.cuda.get_device_properties(0)
start_gpu_memory = round(torch.cuda.max_memory_reserved() / 1024**3, 3)
max_memory = round(gpu_stats.total_memory / 1024**3, 3)
print(f"GPU: {gpu_stats.name}, Max memory: {max_memory} GB")
print(f"Initial reserved memory: {start_gpu_memory} GB")
# 5. 학습 실행
print("=" * 80)
print("Step 5: Starting Training")
print("=" * 80)
trainer_stats = trainer.train()
# 학습 완료 후 메모리 및 시간 통계 출력
used_memory = round(torch.cuda.max_memory_reserved() / 1024**3, 3)
used_memory_for_lora = round(used_memory - start_gpu_memory, 3)
used_percentage = round(used_memory / max_memory * 100, 3)
lora_percentage = round(used_memory_for_lora / max_memory * 100, 3)
print("=" * 80)
print("Training Statistics")
print("=" * 80)
print(f"Training time: {trainer_stats.metrics['train_runtime']:.2f} seconds ({trainer_stats.metrics['train_runtime']/60:.2f} minutes)")
print(f"Peak reserved memory: {used_memory} GB ({used_percentage}% of max)")
print(f"Memory for training: {used_memory_for_lora} GB ({lora_percentage}% of max)")
# 6. 모델 저장
print("=" * 80)
print("Step 6: Saving Model and Tokenizer")
print("=" * 80)
model_dir = f"{pvc_data_path}/kanana-2-1b-kcdocs"
# LoRA 가중치를 기본 모델과 병합 후 저장
model = model.merge_and_unload()
model.save_pretrained(model_dir)
tokenizer.save_pretrained(model_dir)
print(f"Model and tokenizer saved to: {model_dir}")
# 7. 분산 학습 프로세스 그룹 정리 (경고 방지)
print("=" * 80)
print("Step 7: Cleaning up distributed process group")
print("=" * 80)
if torch.distributed.is_initialized():
torch.distributed.destroy_process_group()
print("Distributed process group destroyed successfully")
else:
print("No distributed process group to clean up")
# 8. CUDA 컨텍스트 및 리소스 정리
print("=" * 80)
print("Step 8: Cleaning up CUDA resources")
print("=" * 80)
# CUDA 캐시 정리
if torch.cuda.is_available():
torch.cuda.empty_cache()
torch.cuda.synchronize()
print("CUDA cache cleared and synchronized")
# 모델과 토크나이저를 메모리에서 해제
del model
del tokenizer
del base_model
import gc
gc.collect()
if torch.cuda.is_available():
torch.cuda.empty_cache()
print("Model and tokenizer released from memory")
# 9. 명시적으로 프로세스 종료 (컨테이너가 종료되도록)
print("=" * 80)
print("Step 9: Training completed successfully, exiting...")
print("=" * 80)
# CustomTrainer 설정
pvc_data_path = "/data"
llm_model_name = "kakaocorp/kanana-nano-2.1b-base"
trainer = CustomTrainer(
func=finetune_kanana,
func_args={
"model_name": llm_model_name,
"epoch_num": epoch_num,
"pvc_data_path": pvc_data_path
},
num_nodes=1,
resources_per_node={
"nvidia.com/gpu": "1",
"cpu": "8",
"memory": "16Gi"
},
packages_to_install=[
"transformers",
"peft",
"datasets",
"torch",
"pandas",
"accelerate>=0.26.0",
],
)
trainer_client = TrainerClient()
train_kwargs = {"trainer": trainer}
# PVC 마운트를 위한 PodTemplateOverrides 설정
if pvc_name:
from kubeflow.trainer.options.kubernetes import (
PodTemplateOverrides,
PodTemplateOverride,
PodSpecOverride,
ContainerOverride,
)
options_list = [
PodTemplateOverrides(
PodTemplateOverride(
target_jobs=["node"],
metadata={
"annotations": {
"sidecar.istio.io/inject": "false"
}
},
spec=PodSpecOverride(
volumes=[
{
"name": "data",
"persistentVolumeClaim": {"claimName": pvc_name}
}
],
containers=[
ContainerOverride(
name="node",
volume_mounts=[
{
"name": "data",
"mountPath": pvc_data_path
}
]
)
]
)
)
)
]
train_kwargs["options"] = options_list
except ImportError:
print("Warning: Could not import PodTemplateOverrides. PVC mounting may not work.")
# TrainJob 생성
print("Creating TrainJob...")
import inspect
sig = inspect.signature(trainer_client.train)
if 'runtime' in sig.parameters and sig.parameters['runtime'].default != inspect.Parameter.empty:
job_id = trainer_client.train(**train_kwargs)
else:
try:
job_id = trainer_client.create_trainjob(trainer)
except AttributeError:
raise RuntimeError(
"TrainJob creation requires runtime configuration. "
"Please ensure ClusterTrainingRuntime is properly configured."
)
print(f"TrainJob created with job_id: {job_id}")
# TrainJob 시작 대기
print("Waiting for TrainJob to start...")
time.sleep(10)
# TrainJob 로그 확인
print("\n=== TrainJob Logs ===")
try:
logs = list(trainer_client.get_job_logs(job_id, follow=False))
if logs:
for logline in logs:
print(logline)
else:
print("No logs available yet")
except Exception as log_error:
print(f"Warning: Could not retrieve logs: {log_error}")
print("\nTraining job submitted successfully!")
# job_id를 컴포넌트 출력으로 반환
with open(train_job_id.path, 'w') as f:
f.write(job_id)
2-3. 모델 서빙 컴포넌트
KServe InferenceService를 생성하여 파인튜닝된 모델을 vLLM으로 서빙하는 컴포넌트입니다.
- vLLM 백엔드를 사용한 고성능 추론
- PVC에서 모델 로드
- GPU 리소스 자동 할당
@dsl.component(
packages_to_install=['kubeflow', 'kubernetes', 'pyyaml'],
install_kfp_package=True,
base_image='python:3.11',
output_component_file=f'{SERVE_ENPOINT_PATH}/deploy_component.yaml'
)
def deploy_kanana_op_func(
namespace: str,
pvc_name: str,
kserve_name: str,
train_job_id: dsl.Input[dsl.Artifact],
model_path_in_pvc: str = "kanana-2-1b-kcdocs",
served_model_name: str = "kanana-nano-2.1b-base",
max_model_len: str = "32768",
gpu_memory_utilization: str = "0.8",
):
"""
KServe InferenceService를 생성하여 파인튜닝된 모델을 vLLM으로 배포
Args:
namespace: Kubernetes 네임스페이스
pvc_name: 모델이 저장된 PVC 이름
kserve_name: InferenceService 이름
train_job_id: TrainJob ID 아티팩트 (이전 컴포넌트에서 전달받음)
model_path_in_pvc: PVC 내 모델 경로
served_model_name: 서빙될 모델 이름
max_model_len: 최대 시퀀스 길이
gpu_memory_utilization: GPU 메모리 사용률 (0.0-1.0)
"""
from kubernetes import client, config
from kubernetes.config import ConfigException
from kubeflow.trainer import TrainerClient
import time
import sys
# Kubernetes 클라이언트 초기화
try:
config.load_incluster_config()
except ConfigException:
config.load_kube_config() # 로컬에서 실행 시
custom_api = client.CustomObjectsApi()
group = "serving.kserve.io"
version = "v1beta1"
plural = "inferenceservices"
# KFP v2에서 아티팩트는 파일 경로를 통해 접근합니다
# 아티팩트 파일에서 실제 job_id 문자열을 읽어옵니다
with open(train_job_id.path, 'r') as f:
trainjob_id_str = f.read().strip()
# TrainerClient 초기화
trainer_client = TrainerClient()
# polling until trainjob is done
max_wait = 60 * 60 # 1 hour at most
interval = 15 # seconds
waited = 0
print(f"Waiting for TrainJob '{trainjob_id_str}' to complete...")
while waited < max_wait:
try:
trainjob = trainer_client.get_job(trainjob_id_str)
job_status = trainjob.status if hasattr(trainjob, 'status') else None
if job_status == 'Complete':
print(f"TrainJob '{trainjob_id_str}' completed successfully.")
break
if job_status == 'Failed':
raise RuntimeError(f"TrainJob '{trainjob_id_str}' failed!")
except RuntimeError:
raise
except Exception as e:
# TrainJob이 아직 생성되지 않았거나 다른 에러인 경우 계속 대기
pass
time.sleep(interval)
waited += interval
else:
raise TimeoutError(f"Timed out waiting for TrainJob '{trainjob_id_str}' to complete.")
# InferenceService 매니페스트 정의
inferenceservice_manifest = {
"apiVersion": f"{group}/{version}",
"kind": "InferenceService",
"metadata": {
"name": kserve_name,
"namespace": namespace,
},
"spec": {
"predictor": {
"annotations": {
"serving.knative.dev/progress-deadline": "1h"
},
"automountServiceAccountToken": False,
"maxReplicas": 1,
"minReplicas": 1,
"model": {
"args": [
f"--model_name={served_model_name}",
"--model_id=/mnt/models",
"--dtype=bfloat16",
"--backend=vllm"
],
"env": [
{
"name": "VLLM_LOGGING_LEVEL",
"value": "DEBUG"
},
{
"name": "MAX_MODEL_LEN",
"value": max_model_len
},
{
"name": "GPU_MEMORY_UTILIZATION",
"value": gpu_memory_utilization
},
{
"name": "PYTORCH_CUDA_ALLOC_CONF",
"value": "expandable_segments:True,max_split_size_mb:128"
}
],
"lifecycle": {
"preStop": {
"exec": {
"command": [""]
}
}
},
"modelFormat": {
"name": "huggingface"
},
"name": "",
"resources": {
"limits": {
"cpu": "23",
"memory": "180Gi",
"nvidia.com/gpu": "1"
},
"requests": {
"cpu": "1",
"memory": "2Gi",
"nvidia.com/gpu": "1"
}
},
"storageUri": f"pvc://{pvc_name}/{model_path_in_pvc}"
},
"timeout": 600
}
}
}
# InferenceService 생성 또는 업데이트
try:
try:
existing_isvc = custom_api.get_namespaced_custom_object(
group=group,
version=version,
namespace=namespace,
plural=plural,
name=kserve_name
)
print(f"InferenceService '{kserve_name}' already exists. Updating...")
inferenceservice_manifest["metadata"]["resourceVersion"] = existing_isvc["metadata"].get("resourceVersion")
updated_isvc = custom_api.patch_namespaced_custom_object(
group=group,
version=version,
namespace=namespace,
plural=plural,
name=kserve_name,
body=inferenceservice_manifest
)
print(f"InferenceService '{kserve_name}' updated successfully")
except client.rest.ApiException as e:
if e.status == 404:
print(f"Creating InferenceService '{kserve_name}'...")
created_isvc = custom_api.create_namespaced_custom_object(
group=group,
version=version,
namespace=namespace,
plural=plural,
body=inferenceservice_manifest
)
print(f"InferenceService '{kserve_name}' created successfully")
else:
raise
# InferenceService가 Ready 상태가 될 때까지 대기
print(f"Waiting for InferenceService '{kserve_name}' to be ready...")
max_wait_time = 1800
wait_interval = 10
elapsed_time = 0
while elapsed_time < max_wait_time:
try:
isvc_status = custom_api.get_namespaced_custom_object_status(
group=group,
version=version,
namespace=namespace,
plural=plural,
name=kserve_name
)
conditions = isvc_status.get("status", {}).get("conditions", [])
ready = False
for condition in conditions:
if condition.get("type") == "Ready":
if condition.get("status") == "True":
ready = True
break
if ready:
print(f"InferenceService '{kserve_name}' is ready!")
break
time.sleep(wait_interval)
elapsed_time += wait_interval
print(f"Still waiting... ({elapsed_time}/{max_wait_time} seconds)")
except Exception as status_error:
print(f"Error checking status: {status_error}")
time.sleep(wait_interval)
elapsed_time += wait_interval
if elapsed_time >= max_wait_time:
print(f"Warning: InferenceService '{kserve_name}' did not become ready within {max_wait_time} seconds")
print(f"KServe InferenceService '{kserve_name}' deployment completed successfully")
return
except Exception as e:
print(f"Error deploying InferenceService: {e}")
import traceback
print(f"Traceback: {traceback.format_exc()}")
raise
3. 파이프라인 정의
각 컴포넌트를 연결하여 전체 파이프라인을 정의합니다.
@dsl.pipeline(name="Kanana Model Finetuning Pipeline")
def kanana_model_finetuning_Pipeline(
kc_kbm_os_train_url: str = 'https://objectstorage.kr-central-2.kakaocloud.com/v1/c11fcba415bd4314b595db954e4d4422/public/tutorial/kubeflow/kubeflow-tensorboard/data/sample_train_data.csv',
epoch_num: str = "10",
job_name: str = None,
endpoint_name: str = None,
):
"""
Kanana 모델 파인튜닝 및 서빙 파이프라인
Args:
kc_kbm_os_train_url: 학습 데이터 CSV 파일의 Object Storage URL
epoch_num: 학습 에포크 수
job_name: TrainJob 이름 (선택적)
"""
# 1. PVC 생성 (데이터 및 모델 저장용)
pvc1 = kubernetes.CreatePVC(
pvc_name=PVC_NAME,
access_modes=['ReadWriteMany'],
size='10Gi',
storage_class_name='',
)
# 2. 데이터 다운로드 컴포넌트
download_data = download_dataset(kc_kbm_os_train_url=kc_kbm_os_train_url)
download_data.set_cpu_request(cpu="1").set_memory_request(memory="2G")
download_data.set_caching_options(enable_caching=False)
# PVC 마운트
kubernetes.mount_pvc(
download_data,
pvc_name=pvc1.outputs['name'],
mount_path='/data',
)
# 3. 모델 파인튜닝 컴포넌트
model_train = finetune_kanana_model(
epoch_num=epoch_num,
namespace=KBM_NAMESPACE,
pvc_name=pvc1.outputs['name'],
job_name=job_name,
)
model_train.set_cpu_request(cpu="1").set_memory_request(memory="2G")
# .set_cpu_limit(cpu="1").set_memory_limit(memory="2G")
model_train.set_caching_options(enable_caching=False)
model_train.after(download_data) # 데이터 다운로드 후 실행
# 4. 모델 서빙 컴포넌트
inference_model = deploy_kanana_op_func(
namespace=KBM_NAMESPACE,
kserve_name=endpoint_name,
pvc_name=pvc1.outputs['name'],
train_job_id=model_train.output,
served_model_name="kanana-nano-2.1b-base",
max_model_len="8192",
gpu_memory_utilization="0.8",
)
inference_model.set_cpu_request(cpu="1").set_memory_request(memory="2G")
# .set_cpu_limit(cpu="1").set_memory_limit(memory="2G")
inference_model.set_display_name("Serving Finetuned Kanana Model")
inference_model.after(model_train) # 파인튜닝 완료 후 실행
4. 파이프라인 컴파일 및 실행
experiment_name = kanana_model_finetuning_Pipeline.name + ' experiment'
run_name = kanana_model_finetuning_Pipeline.name + ' run'
arguments = {
"epoch_num": str(EPOCH_NUM),
"job_name": MODEL_NAME,
"endpoint_name": KSERVE_ISVC_NAME,
}
# 파이프라인 실행
run_result = client.create_run_from_pipeline_func(
kanana_model_finetuning_Pipeline,
experiment_name=experiment_name,
run_name=run_name,
arguments=arguments
)
print(f"Pipeline run created: {run_result.run_id}")
Step 3. 실행(Run) 확인하기
-
Kubeflow 대시보드에서 Pipelunes -> Runs 탭을 선택합니다.
-
생성한 실행을 선택하여 상세 화면으로 이동 후 실행에 대한 상세 정보를 확인합니다. 단, 실행이 완료되기까지는 수분이 소요될 수 있습니다.

-
Run 상세 페이지에서 각 컴포넌트의 실행 상태, 로그, 입력 및 출력 결과를 확인할 수 있습니다.

-
Kubeflow 대시보드에서 Trainjob 탭을 선택합니다.
-
파이프라인 실행 중 모델 파인튜닝 컴포넌트에서 생성된 TrainJob의 상태를 확인합니다.


-
Kubeflow 대시보드에서 KServe Endpoints 탭을 선택합니다.
-
파이프라인 실행 완료 후 생성된 KServe InferenceService의 상태와 엔드포인트를 확인합니다.
- 상태: Ready, Pending, Failed 등
- 엔드포인트 URL: 모델 서빙에 사용할 수 있는 URL
- 리소스 사용량: CPU, Memory, GPU 사용량
- 로그: InferenceService 실행 로그


Step 4. 모델 서빙 API 테스트
파이프라인 실행 후 Kserve InferenceService가 Ready 상태가 되면, 배포된 모델을 테스트할 수 있습니다.
방법 1. requests 라이브러리를 사용한 API 테스트
import os
import requests
# Kubeflow 인증을 위한 세션 쿠키 획득
host = os.environ.get("KUBEFLOW_HOST", "https://nipagpu.kakaocloud.com")
username = os.environ.get("KUBEFLOW_USERNAME", "")
password = os.environ.get("KUBEFLOW_PASSWORD", "")
session = requests.Session()
_kargs = {"verify": False} if host.startswith("https") else {}
response = session.get(host, **_kargs)
response = session.post(
response.url,
headers={"Content-Type": "application/x-www-form-urlencoded"},
data={"login": username, "password": password},
**_kargs
)
# 이메일로 받은 MFA 인증 코드 입력
mfa_code = input("이메일 인증 코드: ")
response = session.post(
response.url,
headers={"Content-Type": "application/x-www-form-urlencoded"},
data={"code": mfa_code},
**_kargs
)
response.raise_for_status()
session_cookie = session.cookies.get_dict().get("authservice_session", "")
if not session_cookie:
raise RuntimeError("MFA 인증에 실패했습니다.")
# OpenAI API 형식으로 InferenceService 테스트
NAMESPACE = KBM_NAMESPACE
KUBEFLOW_PUBLIC_DOMAIN = host.split("//")[1]
SERVED_MODEL_NAME = "kanana-nano-2.1b-base"
# 테스트 프롬프트
prompt_text = "카카오엔터프라이즈에 대해서 설명해줘"
# OpenAI completions API 형식 요청
data = {
"model": SERVED_MODEL_NAME,
"prompt": prompt_text,
"stream": False,
"max_tokens": 1000
}
# API 요청
response = requests.post(
url=f"{host}/openai/v1/completions",
cookies={'authservice_session': session_cookie},
headers={
"Host": f"{KSERVE_ISVC_NAME}-{NAMESPACE}.{KUBEFLOW_PUBLIC_DOMAIN}",
"Content-Type": "application/json",
},
json=data,
**_kargs
)
response_json = response.json()
print(f"입력 프롬프트: {prompt_text}")
print(f"상태 코드: {response.status_code}")
print(f"응답: {response_json['choices'][0]['text']}")
방법 2. LangChain을 사용한 API 테스트
LangChain의 ChatOpenAI를 사용하여 InferenceService와 통신할 수 있습니다.
from langchain_openai import ChatOpenAI
# InferenceService 내부 서비스 URL (클러스터 내부에서 접근)
llm_svc_url = f"http://{KSERVE_ISVC_NAME}.{NAMESPACE}.svc.cluster.local/"
llm = ChatOpenAI(
model_name=SERVED_MODEL_NAME,
base_url=f"{llm_svc_url}openai/v1",
openai_api_key="empty" # KServe는 API 키를 요구하지 않음
)
input_text = "카카오엔터프라이즈에 대해서 설명해줘"
result = llm.invoke(input_text)
print(result.content)
Step 5. 실행(Run) 보관하기
-
대시보드에 접속하여 Runs 탭을 클릭한 후, 목록에서 보관할 실행을 선택하고 [Archive] 버튼을 클릭합니다.

-
보관한 실행은 Runs 탭 목록 화면의 Archived 항목에서 확인할 수 있으며, 실행을 선택하고 [Restore] 버튼을 클릭하면 원래대로 복구할 수 있습니다.
Step 6. 실행(Run) 삭제하기
완료 또는 미사용 실행은 리소스 관리를 위해 삭제를 권고드립니다
-
대시보드에 접속하여 Runs 탭을 클릭한 후, 목록에서 삭제할 실행을 선택하고 [Archive] 버튼을 클릭합니다.
-
보관한 실행은 Runs 탭 목록 화면의 Archived 항목에서 확인할 수 있으며, 실행을 선택하고 [Delete] 버튼을 클릭하면 실행을 삭제할 수 있습니다.

-
실행 삭제 시, 파드까지 삭제된 것을 확인할 수 있습니다.
파이프라인 생성에 대한 자세한 설명은 Kubeflow > Kubeflow Pipeline > Quick Start 문서를 확인해 주세요.