26S26g

QK·2026년 9월 27일
#!/usr/bin/env python3
"""
Batch Markdown Preprocessor & Enricher for Vector DB (vLLM Backend)
- RCA(Root Cause Analysis) 및 SRE 최적화 버전

주요 기능:
1. 지정 디렉터리 내 모든 .md 파일 재귀 스캔 (YAML Frontmatter 자동 인식)
2. 헤더 계층 보존 (#, ##, ###) 및 코드/CLI 블록 절단 방지 청킹
3. vLLM 비동기 고속 배치를 통한 RCA 특화 메타데이터 추출:
   - doc_category (TROUBLESHOOTING, CONFIGURATION, ARCHITECTURE, etc.)
   - symptoms_and_errors (증상 및 에러 패턴 정규화)
   - related_metrics (Prometheus 메트릭 명칭 매핑)
   - target_resources (Node, Pod, DirectPVDrive 등 K8s/인프라 리소스)
   - is_destructive_action (데이터 삭제/포맷 등 위험 조치 플래그)
   - hypothetical questions (HyDE 가상 검색 질문)
4. 개별 백업 JSON 파일 및 Vector DB 일괄 적재(Bulk Upsert)용 단일 통합 JSON 생성
"""

import os
import re
import json
import glob
import asyncio
import argparse
from typing import List, Dict, Any, Tuple
from urllib.parse import urlparse

from openai import AsyncOpenAI
from langchain_text_splitters import MarkdownHeaderTextSplitter, RecursiveCharacterTextSplitter


# ==============================================================================
# 1. Frontmatter 및 Markdown 파싱 유틸리티
# ==============================================================================
def parse_frontmatter(content: str) -> Tuple[Dict[str, Any], str]:
    """YAML Frontmatter 메타데이터와 본문 텍스트를 분리합니다."""
    pattern = r"^---\s*\n(.*?)\n---\s*\n(.*)$"
    match = re.match(pattern, content, re.DOTALL)
    if not match:
        return {}, content

    yaml_block = match.group(1)
    body = match.group(2)
    metadata = {}

    for line in yaml_block.splitlines():
        line = line.strip()
        if not line or line.startswith("#"):
            continue
        if ":" in line:
            key, val = line.split(":", 1)
            val = val.strip().strip('"').strip("'")
            metadata[key.strip()] = val

    return metadata, body


def split_markdown_preserving_code(
    md_text: str, 
    file_path: str,
    base_metadata: Dict[str, Any],
    default_component: str
) -> List[Dict[str, Any]]:
    """
    1) Markdown 헤더(#, ##, ###) 계층 보존
    2) 코드 블록(```)이 중간에 잘리지 않도록 안전 분할
    3) Breadcrumb 경로 자동 생성
    """
    headers_to_split = [
        ("#", "h1"),
        ("##", "h2"),
        ("###", "h3"),
        ("####", "h4"),
    ]
    header_splitter = MarkdownHeaderTextSplitter(
        headers_to_split_on=headers_to_split, 
        strip_headers=False
    )
    sections = header_splitter.split_text(md_text)

    # 코드 블록/문단 우선 분할기
    text_splitter = RecursiveCharacterTextSplitter(
        chunk_size=1200,
        chunk_overlap=150,
        separators=[
            "\n```\n",
            "\n```bash\n",
            "\n```yaml\n",
            "\n```json\n",
            "\n\n",
            "\n",
            " "
        ],
        keep_separator=True
    )

    component = base_metadata.get("component", default_component)
    source_url = base_metadata.get("source_url", "")
    file_name = os.path.basename(file_path)

    chunks = []
    chunk_idx = 0

    for section in sections:
        sub_docs = text_splitter.split_documents([section])
        
        headers = [section.metadata.get(h) for h in ["h1", "h2", "h3", "h4"] if section.metadata.get(h)]
        breadcrumb = " > ".join(headers) if headers else base_metadata.get("title", file_name)

        for doc in sub_docs:
            chunk_idx += 1
            chunks.append({
                "chunk_id": f"{component}_{file_name}_{chunk_idx:04d}",
                "file_path": file_path,
                "component": component,
                "source_url": source_url,
                "breadcrumb": breadcrumb,
                "content": doc.page_content,
                "header_metadata": section.metadata
            })

    return chunks


# ==============================================================================
# 2. vLLM 확장 프롬프트 및 비동기 추론 엔진
# ==============================================================================
PROMPT_TEMPLATE = """당신은 대규모 분산 스토리지 및 쿠버네티스 인프라를 운영하는 최고 수준의 SRE/RCA(Root Cause Analysis) 전문가입니다.
주어진 기술 문서 청크를 정밀 분석하여, 장애 분석 시스템(RCA Agent)과 Vector DB 검색에 필요한 메타데이터를 정형 JSON 형태로 추출하십시오.

[문서 헤더 경로]
{breadcrumb}

[문서 내용]
{content}

[추출 규칙 및 필드 설명]
1. summary: 이 청크가 다루는 핵심 주제나 조치 내용 1줄 요약
2. doc_category: 다음 중 가장 적합한 카테고리 하나만 선택 ["TROUBLESHOOTING", "CONFIGURATION", "ARCHITECTURE", "COMMAND_REFERENCE", "INSTALLATION"]
3. questions: 운영자가 장애 대응이나 시스템 설정 시 검색할 만한 가상 질문 2~3개 (한국어/영어 혼용 권장)
4. symptoms_and_errors: 본문에 언급되거나 암시된 에러 메시지, HTTP 상태 코드, 장애 증상 문장 목록 (예: ["Drive formatted status stuck in pending", "Client HTTP 503 Slow Down", "Volume mount failed"])
5. related_metrics: 본문과 연관된 Prometheus/MinIO 메트릭 명칭 목록. 직접 언급이 없더라도 내용상 직결되는 표준 메트릭 유추 (예: ["minio_node_drive_total_duration_microseconds", "node_disk_io_time_seconds_total"])
6. target_resources: 연관된 쿠버네티스 또는 시스템 리소스 타입 목록 (예: ["Node", "Pod", "DirectPVDrive", "PersistentVolumeClaim"])
7. commands: 본문에 등장하는 핵심 CLI 명령어 목록 (예: ["kubectl directpv drives format", "mc admin drive offline"])
8. is_destructive_action: 본문의 조치 방법이 데이터 유실 위험이 있거나 드라이브 포맷/삭제/오프라인 격리 등의 파괴적 작업인지 여부 (true 또는 false)

반드시 유효한 JSON 형식만 출력하십시오. 마크다운 따옴표나 기타 설명 없이 오직 JSON 본문만 반환하십시오.
{{
  "summary": "1줄 요약",
  "doc_category": "TROUBLESHOOTING",
  "questions": ["질문 1", "질문 2"],
  "symptoms_and_errors": ["에러 1", "증상 2"],
  "related_metrics": ["metric_name_1"],
  "target_resources": ["DirectPVDrive", "Node"],
  "commands": ["kubectl directpv ..."],
  "is_destructive_action": false
}}"""

def clean_json_response(raw_text: str) -> Dict[str, Any]:
    """vLLM 응답에서 마크다운 코드블록을 제거하고 순수 JSON을 파싱합니다."""
    text = raw_text.strip()
    # 마크다운 ```json ... ``` 래핑 제거
    if text.startswith("```"):
        text = re.sub(r"^```[a-zA-Z]*\n", "", text)
        text = re.sub(r"\n```$", "", text).strip()
    return json.loads(text)


async def enrich_chunk(
    client: AsyncOpenAI, 
    model: str, 
    chunk: Dict[str, Any], 
    semaphore: asyncio.Semaphore
) -> Dict[str, Any]:
    """vLLM 비동기 호출을 통해 청크를 심층 보강합니다."""
    async with semaphore:
        prompt = PROMPT_TEMPLATE.format(
            breadcrumb=chunk["breadcrumb"],
            content=chunk["content"]
        )
        try:
            response = await client.chat.completions.create(
                model=model,
                messages=[{"role": "user", "content": prompt}],
                temperature=0.1,
                max_tokens=600,
                response_format={"type": "json_object"}
            )
            raw_text = response.choices[0].message.content
            enrichment_data = clean_json_response(raw_text)
        except Exception:
            enrichment_data = {
                "summary": chunk["breadcrumb"],
                "doc_category": "COMMAND_REFERENCE",
                "questions": [],
                "symptoms_and_errors": [],
                "related_metrics": [],
                "target_resources": [],
                "commands": [],
                "is_destructive_action": False
            }

        # 1) Vector 임베딩용 청크 헤더 구성 (코사인 유사도 검색 최적화)
        q_lines = "\n".join([f"- Q: {q}" for q in enrichment_data.get("questions", [])])
        symptoms_str = ", ".join(enrichment_data.get("symptoms_and_errors", [])) or "None"
        resources_str = ", ".join(enrichment_data.get("target_resources", [])) or "General"
        category = enrichment_data.get("doc_category", "GENERAL")

        header_block = (
            f"[Context: {chunk['breadcrumb']}]\n"
            f"[Category: {category}] [Target Resources: {resources_str}]\n"
            f"[Summary: {enrichment_data.get('summary', '')}]\n"
            f"[Symptoms/Errors: {symptoms_str}]"
        )
        if q_lines:
            header_block += f"\n{q_lines}"

        enriched_content = f"{header_block}\n\n{chunk['content']}"

        # 2) Vector DB 적재용 정규화 스키마 반환
        return {
            "id": chunk["chunk_id"],
            "page_content": enriched_content.strip(),
            "metadata": {
                "component": chunk["component"],
                "doc_category": category,
                "source_url": chunk["source_url"],
                "file_path": chunk["file_path"],
                "breadcrumb": chunk["breadcrumb"],
                "summary": enrichment_data.get("summary", ""),
                "questions": enrichment_data.get("questions", []),
                "symptoms_and_errors": enrichment_data.get("symptoms_and_errors", []),
                "related_metrics": enrichment_data.get("related_metrics", []),
                "target_resources": enrichment_data.get("target_resources", []),
                "commands": enrichment_data.get("commands", []),
                "is_destructive_action": bool(enrichment_data.get("is_destructive_action", False)),
                **chunk["header_metadata"]
            }
        }


# ==============================================================================
# 3. 디렉터리 일괄 처리 및 단일 통합 JSON 생성
# ==============================================================================
async def process_directory(
    input_dir: str,
    output_dir: str,
    final_output_file: str,
    default_component: str,
    vllm_endpoint: str,
    model_name: str,
    concurrency: int = 32
):
    os.makedirs(output_dir, exist_ok=True)
    client = AsyncOpenAI(base_url=vllm_endpoint, api_key="EMPTY")
    semaphore = asyncio.Semaphore(concurrency)

    # 1. 대상 디렉터리 내 모든 .md 파일 검색
    md_files = glob.glob(os.path.join(input_dir, "**", "*.md"), recursive=True)
    if not md_files:
        print(f"[!] {input_dir} 디렉터리 내에 .md 파일이 존재하지 않습니다.")
        return

    print("=" * 65)
    print(f"[*] 검색된 총 Markdown 파일 수 : {len(md_files)}개")
    print(f"[*] vLLM API 엔드포인트        : {vllm_endpoint}")
    print(f"[*] 서빙 중인 모델            : {model_name}")
    print(f"[*] 동시 비동기 요청 수 (B300) : {concurrency}")
    print("=" * 65)

    all_enriched_chunks: List[Dict[str, Any]] = []

    # 2. 파일별 순회 처리
    for idx, file_path in enumerate(md_files, 1):
        rel_path = os.path.relpath(file_path, input_dir)
        print(f"\n[{idx}/{len(md_files)}] 문서 분석 및 청킹: {rel_path}")

        with open(file_path, "r", encoding="utf-8") as f:
            raw_text = f.read()

        frontmatter, clean_body = parse_frontmatter(raw_text)
        initial_chunks = split_markdown_preserving_code(
            clean_body, 
            file_path=file_path, 
            base_metadata=frontmatter, 
            default_component=default_component
        )
        print(f"    - 생성된 청크 수: {len(initial_chunks)}개 (vLLM 메타데이터 추출 요청 중...)")

        if not initial_chunks:
            continue

        # vLLM 비동기 보강 실행
        tasks = [enrich_chunk(client, model_name, c, semaphore) for c in initial_chunks]
        file_enriched_chunks = await asyncio.gather(*tasks)

        # 개별 파일별 중간 JSON 저장 (단일 장애 복구용 백업)
        safe_rel_name = rel_path.replace(os.sep, "_").replace(".md", ".json")
        per_file_output = os.path.join(output_dir, safe_rel_name)
        with open(per_file_output, "w", encoding="utf-8") as f:
            json.dump(file_enriched_chunks, f, ensure_ascii=False, indent=2)

        all_enriched_chunks.extend(file_enriched_chunks)

    # 3. 전체 청크가 담긴 단일 통합 JSON 파일 저장 (Bulk Upsert용)
    final_path = os.path.join(output_dir, final_output_file)
    with open(final_path, "w", encoding="utf-8") as f:
        json.dump(all_enriched_chunks, f, ensure_ascii=False, indent=2)

    print("\n" + "=" * 65)
    print(f"[✓] 전체 일괄 변환 및 메타데이터 보강 완료!")
    print(f"[*] 총 생성된 Vector DB 청크 개수 : {len(all_enriched_chunks)}개")
    print(f"[*] 개별 백업 JSON 디렉터리       : {os.path.abspath(output_dir)}")
    print(f"[*] 최종 단일 통합 JSON 파일 위치 : {os.path.abspath(final_path)}")
    print("=" * 65)


def main():
    parser = argparse.ArgumentParser(description="Markdown 문서들을 vLLM으로 분석/보강하여 RCA Vector DB용 단일 통합 JSON 생성")
    parser.add_argument("--input-dir", required=True, help="스캔할 Markdown 최상위 폴더 경로")
    parser.add_argument("--output-dir", default="./vectordb_chunks", help="결과물 저장 폴더")
    parser.add_argument("--final-json", default="bulk_chunks_for_vectordb.json", help="최종 단일 통합 파일명")
    parser.add_argument("--component", default="minio-system", help="기본 컴포넌트 식별자")
    parser.add_argument("--vllm-url", default="http://localhost:8000/v1", help="vLLM API 엔드포인트")
    parser.add_argument("--model", default="Qwen/Qwen2.5-32B-Instruct", help="vLLM에 서빙 중인 모델 이름")
    parser.add_argument("--concurrency", type=int, default=48, help="vLLM 동시 요청 수 (B300 기준 48~64 권장)")

    args = parser.parse_args()

    asyncio.run(
        process_directory(
            input_dir=args.input_dir,
            output_dir=args.output_dir,
            final_output_file=args.final_json,
            default_component=args.component,
            vllm_endpoint=args.vllm_url,
            model_name=args.model,
            concurrency=args.concurrency
        )
    )


if __name__ == "__main__":
    main()
profile
engineer

0개의 댓글