오래 걸리는 작업 취소하기

오래 걸리는 작업 취소하기 (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 작업에 대한 견고한 제어를 제공해서, 적절한 리소스 관리와 사용자 경험으로 프로덕션 배포가 가능하게 해줘요.

더 알아보기 (Learn more)