"""
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
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])
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
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": []
}
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']}"
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"]
}
}
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)
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]] = []
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
tasks = [enrich_chunk(client, model_name, c, semaphore) for c in initial_chunks]
file_enriched_chunks = await asyncio.gather(*tasks)
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)
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()