4b86b2a660
- server-dev start/stop/deploy 및 Gitea push 자동 배포 - local-dev 로컬 개발 환경 Co-authored-by: Cursor <cursoragent@cursor.com>
195 lines
8.1 KiB
Python
195 lines
8.1 KiB
Python
# 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)
|