26S26e

QK·2026년 9월 27일
#!/usr/bin/env python3
"""
Batch Markdown Preprocessor & Enricher for Vector DB (vLLM Backend)

1. 특정 디렉터리 하위의 모든 .md 파일을 재귀 탐색 (os.walk)
2. 각 파일의 YAML Frontmatter(존재 시) 메타데이터 파싱
3. 마크다운 헤더 계층 보존 및 코드 블록 깨짐 방지 청킹
4. B300 vLLM 비동기(Async) 엔진을 통한 고속 배치 가상 질문/메타데이터 추출
5. 개별 파일 결과 보존 및 최종 단일 통합 JSON (Vector DB Bulk Upsert용) 생성
"""

import os
import re
import json
import glob
import asyncio
import argparse
from typing import List, Dict, Any, Optional
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) 헤더 계층 보존 (#, ##, ###)
    2) 코드 블록(\n```\n) 분할 보호
    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])
        
        # 헤더 Breadcrumbs 구성
        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/DevOps 전문가입니다.
주어진 기술 문서 청크를 분석하여, 운영자가 문제 해결 시 검색할 법한 가상 질문과 핵심 메타데이터를 JSON 형식으로 추출하십시오.

[문서 헤더 경로]
{breadcrumb}

[문서 내용]
{content}

[추출 규칙]
1. questions: 장애 조치, 에러 해결, 설정 질문 위주로 2~3개 한국어/영어 혼용 생성
2. error_keywords: 본문에 언급된 에러 코드, 상태(예: DriveNotReady, 503, Offline 등)
3. commands: 본문에 포함된 CLI 명령어 프리픽스 (예: `kubectl directpv`, `mc admin` 등)
4. summary: 이 청크의 핵심 내용 1줄 요약

반드시 유효한 JSON 문자열만 출력하십시오. 마크다운 코드블록이나 불필요한 설명은 금지합니다.
{{
  "summary": "1줄 요약",
  "questions": ["가상 질문 1", "가상 질문 2"],
  "error_keywords": ["키워드1", "키워드2"],
  "commands": ["명령어1"]
}}"""

async def enrich_chunk(
    client: AsyncOpenAI, 
    model: str, 
    chunk: Dict[str, Any], 
    semaphore: asyncio.Semaphore
) -> Dict[str, Any]:
    """청크에 가상 질문을 삽입하고 정형 메타데이터를 부여"""
    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=512,
                response_format={"type": "json_object"}
            )
            raw_json = response.choices[0].message.content
            enrichment_data = json.loads(raw_json)
        except Exception:
            enrichment_data = {
                "summary": chunk["breadcrumb"],
                "questions": [],
                "error_keywords": [],
                "commands": []
            }

        # 1) 청크 헤더에 가상 질문 및 Context 삽입 (HyDE 매칭 최적화)
        q_lines = "\n".join([f"- Q: {q}" for q in enrichment_data.get("questions", [])])
        header_context = f"[Context: {chunk['breadcrumb']}]\n[Summary: {enrichment_data.get('summary', '')}]"
        if q_lines:
            header_context += f"\n{q_lines}"

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

        # 2) Vector DB 저장용 표준 구조 생성
        return {
            "id": chunk["chunk_id"],
            "page_content": enriched_content.strip(),
            "metadata": {
                "component": chunk["component"],
                "source_url": chunk["source_url"],
                "file_path": chunk["file_path"],
                "breadcrumb": chunk["breadcrumb"],
                "summary": enrichment_data.get("summary", ""),
                "questions": enrichment_data.get("questions", []),
                "error_keywords": enrichment_data.get("error_keywords", []),
                "commands": enrichment_data.get("commands", []),
                **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(f"[*] 총 {len(md_files)}개의 Markdown 파일을 발견했습니다.")
    print(f"[*] vLLM 엔드포인트: {vllm_endpoint} (동시 처리: {concurrency})")

    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)}개")

        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 파일 저장
    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 보강 후 단일 JSON으로 병합")
    parser.add_argument("--input-dir", required=True, help="스캔할 Markdown 최상위 폴더 경로")
    parser.add_argument("--output-dir", default="./vectordb_chunks", help="개별 및 통합 JSON 저장 폴더")
    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-14B-Instruct", help="vLLM에 서빙 중인 모델 이름")
    parser.add_argument("--concurrency", type=int, default=32, help="vLLM 동시 요청 수 (B300 기준 32~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개의 댓글