"""
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
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
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()
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
}
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']}"
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"]
}
}
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("=" * 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]] = []
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
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으로 분석/보강하여 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()