# rag-demo/scripts/ingest_qa.py import json, pathlib, os, sys import numpy as np from api_clients import TEIEmbeddingClient from vector_store import get_vector_store, source_created_at_cmp def print_progress(current, total, prefix='', suffix='', step=10): """진행률 출력 (Docker 환경 대응 - 일정 간격으로 새 줄 출력)""" percent = int(100 * current / total) # step% 간격으로만 출력 (10%, 20%, ... 또는 완료 시) if percent % step == 0 or current == total: # 이전에 이 퍼센트를 출력했는지 체크 (중복 방지) if not hasattr(print_progress, '_last_percent'): print_progress._last_percent = {} key = f"{prefix}_{total}" if key not in print_progress._last_percent or print_progress._last_percent[key] != percent: print_progress._last_percent[key] = percent # 프로그레스 바 생성 length = 40 filled = int(length * current / total) bar = '█' * filled + '░' * (length - filled) print(f'{prefix} [{bar}] {percent}% ({current}/{total}) {suffix}', flush=True) # 완료 시 초기화 if current == total: key = f"{prefix}_{total}" if hasattr(print_progress, '_last_percent') and key in print_progress._last_percent: del print_progress._last_percent[key] # ⚠️ 변경: qa.jsonl 대신 qa_raw.jsonl 직접 사용 (preprocess 단계 생략) QA_FILE = pathlib.Path("/app/data/qa_raw.jsonl") # 원본 파일 직접 사용 DATA_DIR = pathlib.Path("/app/data") VECS_FILE = DATA_DIR / "qa_vecs.jsonl" # 배치 크기 설정 (임베딩 API 호출 단위) # - TEI API 제한: 최대 32개까지 한 번에 처리 가능 # - 권장값: 16~32 (안정성을 위해 32 이하 권장) EMBED_BATCH_SIZE = int(os.getenv("EMBED_BATCH_SIZE", "32")) # API 호출 시 한 번에 임베딩할 개수 (최대 32) VECTOR_BATCH_SIZE = int(os.getenv("VECTOR_BATCH_SIZE", "500")) # 벡터 스토어 저장 단위 # ────────────────────────────────────────────── # ✅ 스킵 로직: 이미 임베딩이 완료되었는지 확인 # ────────────────────────────────────────────── def should_skip_embedding(): """임베딩 작업을 스킵해야 하는지 판단""" if not VECS_FILE.exists(): return False # qa_vecs.jsonl의 라인 수와 qa_raw.jsonl의 라인 수 비교 try: with open(QA_FILE, 'r', encoding='utf-8') as f: raw_lines = sum(1 for line in f if line.strip()) with open(VECS_FILE, 'r', encoding='utf-8') as f: vec_lines = sum(1 for line in f if line.strip()) if raw_lines == vec_lines: print(f"[Ingest] ✅ 임베딩 이미 완료됨 (qa_raw: {raw_lines}개, qa_vecs: {vec_lines}개)") print(f"[Ingest] ⏩ 스킵합니다. 재임베딩이 필요하면 'rm {VECS_FILE}'을 실행하세요.") return True else: print(f"[Ingest] ⚠️ 라인 수 불일치 (qa_raw: {raw_lines}개, qa_vecs: {vec_lines}개) → 재임베딩") return False except Exception as e: print(f"[Ingest] ⚠️ 스킵 체크 실패: {e} → 임베딩 진행") return False if should_skip_embedding(): print("[Ingest] 🎉 임베딩 작업 완료 (스킵)") exit(0) # API 클라이언트 및 벡터 스토어 초기화 embed_client = TEIEmbeddingClient() vector_store = get_vector_store() print(f"[Ingest] 임베딩 배치 크기: {EMBED_BATCH_SIZE}개 (API 호출 단위)") print(f"[Ingest] 벡터 저장 배치 크기: {VECTOR_BATCH_SIZE}개 (벡터 스토어 저장 단위)") print("") # 빈 줄 # ── Step 1: QA 데이터 로드 ───────────────────── print("[Ingest] QA 데이터 로딩 중...") qa_data = [] skipped_no_date = 0 skipped_no_qa = 0 for ln, line in enumerate(QA_FILE.open(encoding="utf-8"), 1): line = line.strip() if not line: continue try: obj = json.loads(line) except json.JSONDecodeError as e: raise RuntimeError(f"❌ JSON 오류 (line {ln}): {e.msg}\n> {line}") from None question = obj.get("question") or obj.get("q") answer = obj.get("answer") or obj.get("a") if not question or not answer: skipped_no_qa += 1 print(f"[Ingest] ⚠️ 필수 필드 누락으로 스킵 (line {ln}): q/a 필요") continue source_created_at = obj.get("source_created_at") if source_created_at is None or str(source_created_at).strip() == "": skipped_no_date += 1 print(f"[Ingest] ⚠️ source_created_at 없음으로 스킵 (line {ln})") continue meta = dict(obj) meta["q"] = question meta["a"] = answer meta["source_created_at"] = str(source_created_at).strip() meta.pop("question", None) meta.pop("answer", None) qa_data.append({"question": str(question).strip(), "answer": answer, "meta": meta}) # 동일 질문 중 source_created_at 최신만 유지 deduped = {} dedup_skipped = 0 for item in qa_data: q = item["question"] prev = deduped.get(q) if prev is None: deduped[q] = item continue if source_created_at_cmp(item["meta"]["source_created_at"], prev["meta"]["source_created_at"]) > 0: deduped[q] = item dedup_skipped += 1 else: dedup_skipped += 1 qa_data = list(deduped.values()) print( f"[Ingest] 총 {len(qa_data)}개 QA 쌍 로드 완료 " f"(q/a 스킵 {skipped_no_qa}, 날짜 스킵 {skipped_no_date}, 중복 제거 {dedup_skipped})" ) # ── Step 2: 배치 임베딩 ───────────────────────── print(f"[Ingest] 임베딩 시작 (배치 크기: {EMBED_BATCH_SIZE})") all_vectors = [] all_metadatas = [] for i in range(0, len(qa_data), EMBED_BATCH_SIZE): batch_data = qa_data[i:i + EMBED_BATCH_SIZE] batch_texts = [item["question"] for item in batch_data] # 배치 임베딩 (한 번에 여러 개) embeddings = embed_client.embed(batch_texts, normalize=True, is_query=False) # 벡터 및 메타데이터 수집 for j, emb in enumerate(embeddings): all_vectors.append(emb) all_metadatas.append(batch_data[j]["meta"]) # 진행률 바 표시 (10% 간격) processed = min(i + EMBED_BATCH_SIZE, len(qa_data)) print_progress(processed, len(qa_data), prefix='[Ingest] 임베딩 진행', suffix='✨', step=10) print(f"[Ingest] ✅ 임베딩 완료: {len(all_vectors)}개 벡터") print("") # 빈 줄 # ── Step 3: 벡터 스토어에 배치 저장 ────────────── print(f"[Ingest] 벡터 저장 시작 (배치 크기: {VECTOR_BATCH_SIZE})") total_inserted = total_updated = total_skipped = 0 for i in range(0, len(all_vectors), VECTOR_BATCH_SIZE): batch_vectors = all_vectors[i:i + VECTOR_BATCH_SIZE] batch_metas = all_metadatas[i:i + VECTOR_BATCH_SIZE] vectors_array = np.array(batch_vectors, dtype="float32") if hasattr(vector_store, "upsert_vectors"): stats = vector_store.upsert_vectors(vectors_array, batch_metas, skip_if_older=True) total_inserted += stats.get("inserted", 0) total_updated += stats.get("updated", 0) total_skipped += stats.get("skipped", 0) else: vector_store.add_vectors(vectors_array, batch_metas, skip_if_older=True) # 진행률 바 표시 (10% 간격) processed = min(i + VECTOR_BATCH_SIZE, len(all_vectors)) print_progress(processed, len(all_vectors), prefix='[Ingest] 저장 진행', suffix='💾', step=10) print("") # 빈 줄 if total_inserted or total_updated or total_skipped: print( f"[Ingest] 저장 결과: 신규 {total_inserted}, 갱신 {total_updated}, " f"날짜 구버전 스킵 {total_skipped}" ) # ── Step 4: 최종 저장 ───────────────────────────── vector_store.save(str(DATA_DIR)) print(f"✅ 임베딩 및 인덱싱 완료: {len(all_vectors)}개 벡터", flush=True)