오래 걸리는 작업 취소하기
오래 걸리는 작업 취소하기 (Cancelling Long-Running Tasks)
대규모 데이터셋이나 복잡한 평가에서 일부 Ragas 작업은 완료하는 데 상당한 시간이 걸릴 수 있어요. 취소 기능을 사용하면 필요할 때 이런 오래 걸리는 작업을 우아하게 종료할 수 있어요. 특히 프로덕션 환경에서 중요해요.
출처: 문서
본문
대규모 데이터셋이나 복잡한 평가에서 일부 Ragas 작업은 완료하는 데 상당한 시간이 걸릴 수 있어요. 취소 기능은 필요할 때 이렇게 오래 걸리는 작업을 우아하게 종료할 수 있게 해줘요. 특히 프로덕션 환경에서 중요해요.
개요
Ragas는 다음에 대한 취소 지원을 제공해요.
evaluate()- 메트릭으로 데이터셋 평가generate_with_langchain_docs()- 문서에서 테스트셋 생성
취소 메커니즘은 스레드 안전하며 가능하면 부분 결과와 함께 우아한 종료를 허용해요.
기본 사용법
취소 가능한 평가
평가를 직접 실행하는 대신 취소를 허용하는 executor를 얻을 수 있어요.
from ragas import evaluate
from ragas.dataset_schema import EvaluationDataset
# 데이터셋과 메트릭
dataset = EvaluationDataset(...)
metrics = [...]
# 즉시 평가를 실행하는 대신 executor 얻기
executor = evaluate(
dataset=dataset,
metrics=metrics,
return_executor=True # 핵심 파라미터
)
# 이제 할 수 있는 것:
# - 취소: executor.cancel()
# - 상태 확인: executor.is_cancelled()
# - 결과 얻기: executor.results() # 완료될 때까지 블로킹
취소 가능한 테스트셋 생성
테스트셋 생성에도 비슷한 접근 방식:
from ragas.testset.synthesizers.generate import TestsetGenerator
generator = TestsetGenerator(...)
# 취소 가능한 생성을 위한 executor 얻기
executor = generator.generate_with_langchain_docs(
documents=documents,
testset_size=100,
return_executor=True # 취소를 위해 Executor에 접근 허용
)
# 동일한 취소 인터페이스 사용
executor.cancel()
프로덕션 패턴
1. 타임아웃 패턴
시간 제한을 초과하는 작업을 자동으로 취소해요.
import threading
import time
def evaluate_with_timeout(dataset, metrics, timeout_seconds=300):
"""자동 타임아웃으로 평가 실행."""
# 취소 가능한 executor 얻기
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
results = None
exception = None
def run_evaluation():
nonlocal results, exception
try:
results = executor.results()
except Exception as e:
exception = e
# 백그라운드 스레드에서 평가 시작
thread = threading.Thread(target=run_evaluation)
thread.start()
# 완료 또는 타임아웃 대기
thread.join(timeout=timeout_seconds)
if thread.is_alive():
print(f"Evaluation exceeded {timeout_seconds}s timeout, cancelling...")
executor.cancel()
thread.join(timeout=10) # 필요에 따라 커스텀 타임아웃
return None, "timeout"
return results, exception
# 사용법
results, error = evaluate_with_timeout(dataset, metrics, timeout_seconds=600)
if error == "timeout":
print("Evaluation was cancelled due to timeout")
else:
print(f"Evaluation completed: {results}")
2. 시그널 핸들러 패턴 (Ctrl+C)
사용자가 키보드 인터럽트로 취소할 수 있게 해요.
import signal
import sys
def setup_cancellation_handler():
"""Ctrl+C에서 우아한 취소 설정."""
executor = None
def signal_handler(signum, frame):
if executor and not executor.is_cancelled():
print("\nReceived interrupt signal, cancelling evaluation...")
executor.cancel()
print("Cancellation requested. Waiting for graceful shutdown...")
sys.exit(0)
# 시그널 핸들러 등록
signal.signal(signal.SIGINT, signal_handler)
return lambda exec: setattr(signal_handler, 'executor', exec)
# 사용법
set_executor = setup_cancellation_handler()
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
set_executor(executor)
print("Running evaluation... Press Ctrl+C to cancel gracefully")
try:
results = executor.results()
print("Evaluation completed successfully")
except KeyboardInterrupt:
print("Evaluation was cancelled")
3. 웹 애플리케이션 패턴
웹 애플리케이션에서는 요청이 중단되면 작업을 취소해요.
from flask import Flask, request
import threading
import uuid
app = Flask(__name__)
active_evaluations = {}
@app.route('/evaluate', methods=['POST'])
def start_evaluation():
# 고유한 평가 ID 생성
eval_id = str(uuid.uuid4())
# 요청에서 데이터셋과 메트릭 얻기
dataset = get_dataset_from_request(request)
metrics = get_metrics_from_request(request)
# 취소 가능한 평가 시작
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
active_evaluations[eval_id] = executor
# 백그라운드에서 평가 시작
def run_eval():
try:
results = executor.results()
# 결과를 어딘가에 저장
store_results(eval_id, results)
except Exception as e:
store_error(eval_id, str(e))
finally:
active_evaluations.pop(eval_id, None)
threading.Thread(target=run_eval).start()
return {"evaluation_id": eval_id, "status": "started"}
@app.route('/evaluate/<eval_id>/cancel', methods=['POST'])
def cancel_evaluation(eval_id):
executor = active_evaluations.get(eval_id)
if executor:
executor.cancel()
return {"status": "cancelled"}
return {"error": "Evaluation not found"}, 404
고급 사용법
취소 상태 확인
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
# 백그라운드에서 시작
def monitor_evaluation():
while not executor.is_cancelled():
print("Evaluation still running...")
time.sleep(5)
print("Evaluation was cancelled")
threading.Thread(target=monitor_evaluation).start()
# 어떤 조건 후 취소
if some_condition():
executor.cancel()
부분 결과
실행 중 취소가 발생하면 부분 결과를 얻을 수 있어요.
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
try:
results = executor.results()
print(f"Completed {len(results)} evaluations")
except Exception as e:
if executor.is_cancelled():
print("Evaluation was cancelled - may have partial results")
else:
print(f"Evaluation failed: {e}")
커스텀 취소 로직
class EvaluationManager:
def __init__(self):
self.executors = []
def start_evaluation(self, dataset, metrics):
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
self.executors.append(executor)
return executor
def cancel_all(self):
"""실행 중인 모든 평가 취소."""
for executor in self.executors:
if not executor.is_cancelled():
executor.cancel()
print(f"Cancelled {len(self.executors)} evaluations")
def cleanup_completed(self):
"""완료된 executor 제거."""
self.executors = [ex for ex in self.executors if not ex.is_cancelled()]
# 사용법
manager = EvaluationManager()
# 여러 평가 시작
exec1 = manager.start_evaluation(dataset1, metrics)
exec2 = manager.start_evaluation(dataset2, metrics)
# 필요하면 모두 취소
manager.cancel_all()
모범 사례
1. 프로덕션에서 항상 타임아웃 사용
# 좋음: 항상 합리적인 타임아웃 설정
results, error = evaluate_with_timeout(dataset, metrics, timeout_seconds=1800) # 30 minutes
# 피하기: 무한 블로킹
results = executor.results() # 영원히 블로킹할 수 있음
2. 취소를 우아하게 처리
try:
results = executor.results()
process_results(results)
except Exception as e:
if executor.is_cancelled():
log_cancellation()
cleanup_partial_work()
else:
log_error(e)
handle_failure()
3. 사용자 피드백 제공
def run_with_progress_and_cancellation(executor):
print("Starting evaluation... Press Ctrl+C to cancel")
# 백그라운드에서 진행 상황 모니터링
def show_progress():
while not executor.is_cancelled():
# 진행 표시
print(".", end="", flush=True)
time.sleep(1)
progress_thread = threading.Thread(target=show_progress)
progress_thread.daemon = True
progress_thread.start()
try:
return executor.results()
except KeyboardInterrupt:
print("\nCancelling...")
executor.cancel()
return None
4. 리소스 정리
def managed_evaluation(dataset, metrics):
executor = None
try:
executor = evaluate(dataset=dataset, metrics=metrics, return_executor=True)
return executor.results()
except Exception as e:
if executor:
executor.cancel()
raise
finally:
# 임시 리소스 정리
cleanup_temp_files()
제한 사항
- 비동기 연산: 취소는 작업 단위로 작동하며 개별 LLM 호출 내부가 아님
- 부분 상태: 취소된 연산은 부분 결과나 임시 파일을 남길 수 있음
- 타이밍: 취소는 협조적임 - 작업이 주기적으로 취소를 확인해야 함
- 의존성: 일부 외부 서비스는 취소를 즉시 존중하지 않을 수 있음
문제 해결
취소가 작동하지 않음
# 취소가 설정되었는지 확인
if executor.is_cancelled():
print("Cancellation was requested")
else:
print("Cancellation not requested yet")
# cancel() 호출하는지 확인
executor.cancel()
assert executor.is_cancelled()
취소 후에도 작업이 계속 실행됨
# 우아한 종료 시간 제공
executor.cancel()
time.sleep(2) # 작업이 취소를 감지하도록 허용
# 필요하면 강제 정리
import asyncio
try:
loop = asyncio.get_running_loop()
for task in asyncio.all_tasks(loop):
task.cancel()
except RuntimeError:
pass # # 실행 중인 이벤트 루프 없음
취소 기능은 오래 걸리는 Ragas 작업에 대한 견고한 제어를 제공해서, 적절한 리소스 관리와 사용자 경험으로 프로덕션 배포가 가능하게 해줘요.