886 lines
35 KiB
Python
886 lines
35 KiB
Python
"""
|
|
Weaviate 통합 데이터 업로드 스크립트
|
|
=====================================
|
|
모든 법률 데이터를 Weaviate에 한 번에 업로드합니다.
|
|
|
|
데이터 유형:
|
|
1. Check_List - 체크리스트 문서
|
|
2. Analyzed_Cases - 분석된 판례 (multi-tenant)
|
|
3. Law_and_rules - 법률 및 규칙 (multi-tenant)
|
|
4. Legal_Books - 법률 서적 (multi-tenant)
|
|
5. Past_Cases - 과거 판례 PDF (multi-tenant)
|
|
"""
|
|
|
|
import weaviate
|
|
import weaviate.classes as wvc
|
|
from glob import glob
|
|
import os
|
|
import re
|
|
import json
|
|
import copy
|
|
from typing import Optional, Dict, List
|
|
|
|
# PDF 처리를 위한 선택적 import
|
|
try:
|
|
import fitz # PyMuPDF
|
|
HAS_FITZ = True
|
|
except ImportError:
|
|
HAS_FITZ = False
|
|
print("Warning: PyMuPDF(fitz) not installed. PDF processing will be skipped.")
|
|
|
|
# 토큰 계산을 위한 선택적 import
|
|
try:
|
|
import tiktoken
|
|
tokenizer = tiktoken.get_encoding("cl100k_base")
|
|
HAS_TIKTOKEN = True
|
|
except ImportError:
|
|
HAS_TIKTOKEN = False
|
|
print("Warning: tiktoken not installed. Token counting will use character estimation.")
|
|
|
|
|
|
# ============================================================
|
|
# 설정
|
|
# ============================================================
|
|
WEAVIATE_HOST = "localhost"
|
|
WEAVIATE_PORT = 7080
|
|
MAX_TOKENS = 8000
|
|
BATCH_SIZE = 3 # 배치 크기 (타임아웃 방지)
|
|
|
|
# 데이터 경로 설정 (필요에 따라 수정)
|
|
DATA_PATHS = {
|
|
"check_list": [
|
|
# Check_List 컬렉션용 파일 (현재 없음)
|
|
],
|
|
"analyzed_cases": {
|
|
"Cases_Actio_Pauliana": "/mnt/nas/Weaviate/Analyzed_Cases/Cases_Actio_Pauliana",
|
|
},
|
|
"law_and_rules": {
|
|
"How_to_Compute_value_of_claim_and_damages": "/mnt/nas/Weaviate/How_to_Compute_value_of_claim_and_damages",
|
|
},
|
|
"legal_books": {
|
|
"Information_for_Actio_Pauliana_Suit": "/mnt/nas/Weaviate/Information_for_Claim/Information_for_Actio_Pauliana_Suit/",
|
|
"Information_for_Loan_Claim_Suit": "/mnt/nas/Weaviate/Information_for_Claim/Information_for_Loan_Claim_Suit/",
|
|
"Information_for_Indemnity_Claim_Suit": "/mnt/nas/Weaviate/Information_for_Claim/Information_for_Indemnity_Claim_Suit/",
|
|
},
|
|
# Legal_Books에 개별 파일로 업로드할 문서들 (파일명이 테넌트가 됨)
|
|
"legal_books_files": [
|
|
"/mnt/nas/Weaviate/Criteria_for_Individual_Consolidated_Claim/Criteria_individual_consolidated_claim.txt",
|
|
"/mnt/nas/Weaviate/Information_for_Indemnity_Claim_Suit/Indemnity_Claim_legally_required_facts.md",
|
|
],
|
|
"past_cases": {
|
|
"Cases_Loan_Claim": "/mnt/nas/Weaviate/Past_Cases/Cases_Loan_Claim",
|
|
"Cases_Indemnity_Claim": "/mnt/nas/Weaviate/Past_Cases/Cases_Indemnity_Claim",
|
|
"Cases_Guarantee_Claim": "/mnt/nas/Weaviate/Past_Cases/Cases_Guarantee_Claim",
|
|
},
|
|
}
|
|
|
|
|
|
# ============================================================
|
|
# 유틸리티 함수
|
|
# ============================================================
|
|
def count_tokens(text: str) -> int:
|
|
"""텍스트의 토큰 수를 계산합니다."""
|
|
if not text:
|
|
return 0
|
|
if HAS_TIKTOKEN:
|
|
return len(tokenizer.encode(text))
|
|
else:
|
|
# 대략적인 추정 (한글 기준 2자당 1토큰)
|
|
return len(text) // 2
|
|
|
|
|
|
def split_text_by_tokens(text: str, max_tokens: int) -> list:
|
|
"""긴 텍스트를 max_tokens 단위로 자릅니다."""
|
|
if HAS_TIKTOKEN:
|
|
tokens = tokenizer.encode(text)
|
|
chunks = []
|
|
for i in range(0, len(tokens), max_tokens):
|
|
chunk_tokens = tokens[i:i + max_tokens]
|
|
chunks.append(tokenizer.decode(chunk_tokens))
|
|
return chunks
|
|
else:
|
|
# 대략적인 분할 (문자 기반)
|
|
char_limit = max_tokens * 2
|
|
chunks = []
|
|
for i in range(0, len(text), char_limit):
|
|
chunks.append(text[i:i + char_limit])
|
|
return chunks
|
|
|
|
|
|
def split_text_by_token_limit(header: str, body_text: str, limit: int = 8000) -> list:
|
|
"""
|
|
본문이 토큰 제한을 넘을 경우, 의미 단위(줄바꿈)로 분할하여 리스트로 반환합니다.
|
|
"""
|
|
full_text = f"{header}\n{body_text}"
|
|
|
|
if count_tokens(full_text) <= limit:
|
|
return [full_text]
|
|
|
|
splitted_chunks = []
|
|
current_chunk = [header]
|
|
current_length = count_tokens(header)
|
|
|
|
lines = body_text.split('\n')
|
|
|
|
for line in lines:
|
|
line_token_count = count_tokens(line + "\n")
|
|
|
|
if current_length + line_token_count > limit:
|
|
if len(current_chunk) > 1:
|
|
splitted_chunks.append("\n".join(current_chunk))
|
|
current_chunk = [header]
|
|
current_length = count_tokens(header)
|
|
|
|
current_chunk.append(line)
|
|
current_length += line_token_count
|
|
|
|
if len(current_chunk) > 1:
|
|
splitted_chunks.append("\n".join(current_chunk))
|
|
|
|
return splitted_chunks
|
|
|
|
|
|
# ============================================================
|
|
# 스키마 정의
|
|
# ============================================================
|
|
SCHEMAS = {
|
|
"Check_List": [
|
|
wvc.config.Property(name="check_list_name", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="check_list_content", data_type=wvc.config.DataType.TEXT),
|
|
],
|
|
"Analyzed_Cases": [
|
|
wvc.config.Property(name="raw_case_id", data_type=wvc.config.DataType.TEXT, description="판결정보"),
|
|
wvc.config.Property(name="raw_case_abstract", data_type=wvc.config.DataType.TEXT, description="판결요약"),
|
|
wvc.config.Property(name="raw_case_holding_and_reasoning", data_type=wvc.config.DataType.TEXT, description="판시사항"),
|
|
wvc.config.Property(name="raw_case_summary_of_decision", data_type=wvc.config.DataType.TEXT, description="판결요지"),
|
|
],
|
|
"Law_and_rules": [
|
|
wvc.config.Property(name="chapter", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="section", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="article", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="path", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="content", data_type=wvc.config.DataType.TEXT),
|
|
],
|
|
"Legal_Books": [
|
|
wvc.config.Property(name="section", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="level_1", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="level_2", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="level_3", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="level_4", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="level_5", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="path", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="footnotes", data_type=wvc.config.DataType.TEXT),
|
|
wvc.config.Property(name="content", data_type=wvc.config.DataType.TEXT),
|
|
],
|
|
"Past_Cases": [
|
|
wvc.config.Property(name="raw_case_id", data_type=wvc.config.DataType.TEXT, description="판결정보"),
|
|
wvc.config.Property(name="raw_case_holding_and_reasoning", data_type=wvc.config.DataType.TEXT, description="판시사항"),
|
|
wvc.config.Property(name="raw_case_summary_of_decision", data_type=wvc.config.DataType.TEXT, description="판결요지"),
|
|
wvc.config.Property(name="raw_case_reference_statutes", data_type=wvc.config.DataType.TEXT, description="참조조문"),
|
|
wvc.config.Property(name="raw_case_reference_cases", data_type=wvc.config.DataType.TEXT, description="참조판례"),
|
|
wvc.config.Property(name="raw_case_full_header", data_type=wvc.config.DataType.TEXT, description="전문"),
|
|
wvc.config.Property(name="raw_case_ruling_order", data_type=wvc.config.DataType.TEXT, description="주문"),
|
|
wvc.config.Property(name="raw_case_claim_and_appeal", data_type=wvc.config.DataType.TEXT, description="청구취지 및 항소취지"),
|
|
wvc.config.Property(name="raw_case_reasoning", data_type=wvc.config.DataType.TEXT, description="이유"),
|
|
],
|
|
}
|
|
|
|
|
|
# ============================================================
|
|
# 파서 함수들
|
|
# ============================================================
|
|
def parse_analyzed_cases(content: str) -> list:
|
|
"""분석된 판례 텍스트를 파싱합니다."""
|
|
regex = (r"# 판례 \d{1,4}: (?P<raw_case_id>[가-힣 \d\.]+)\s+"
|
|
r"(?P<raw_case_abstract>[가-힣 \d\.\(\)]+)\s+"
|
|
r"## 판시사항\s+"
|
|
r"(?P<raw_case_holding_and_reasoning>[가-힣 \d\.\(\)\[\]\s,]+)"
|
|
r"## 판결요지\s+"
|
|
r"(?P<raw_case_summary_of_decision>[가-힣 \d\.\(\)\[\]\s,]+)")
|
|
|
|
parsed_chunks = []
|
|
for section in content.split("---"):
|
|
matches = re.search(regex, section.strip(), re.MULTILINE)
|
|
if matches:
|
|
raw_case_json = {
|
|
"raw_case_id": matches.group("raw_case_id").strip(),
|
|
"raw_case_abstract": matches.group("raw_case_abstract").strip(),
|
|
"raw_case_holding_and_reasoning": matches.group("raw_case_holding_and_reasoning").strip(),
|
|
"raw_case_summary_of_decision": matches.group("raw_case_summary_of_decision").strip()
|
|
}
|
|
parsed_chunks.append(raw_case_json)
|
|
return parsed_chunks
|
|
|
|
|
|
def parse_law_and_rule(text: str) -> list:
|
|
"""법률 및 규칙 텍스트를 '조(Article)' 단위로 청킹합니다."""
|
|
REGEX_ARTICLE = re.compile(r'^##\s+(.+)')
|
|
|
|
chunks = []
|
|
state = {"chapter": None, "section": None, "article": None}
|
|
current_content_buffer = []
|
|
|
|
def save_chunk():
|
|
if state["article"] and current_content_buffer:
|
|
body_text = "\n".join(current_content_buffer).strip()
|
|
if body_text:
|
|
split_contents = split_text_by_token_limit(state["article"], body_text, limit=8000)
|
|
path_parts = [p for p in [state["chapter"], state["section"], state["article"]] if p]
|
|
path_str = " > ".join(path_parts)
|
|
|
|
raw_metadata = {
|
|
"chapter": state["chapter"],
|
|
"section": state["section"],
|
|
"article": state["article"],
|
|
"path": path_str
|
|
}
|
|
clean_metadata = {k: v for k, v in raw_metadata.items() if v is not None}
|
|
|
|
for content in split_contents:
|
|
chunk_data = {"metadata": clean_metadata, "content": content}
|
|
chunks.append(chunk_data)
|
|
current_content_buffer.clear()
|
|
|
|
lines = text.split('\n')
|
|
for line in lines:
|
|
line = line.rstrip()
|
|
stripped_line = line.strip()
|
|
|
|
if not stripped_line:
|
|
continue
|
|
|
|
if REGEX_ARTICLE.match(stripped_line):
|
|
save_chunk()
|
|
state["article"] = REGEX_ARTICLE.match(stripped_line).group(1)
|
|
else:
|
|
if state["article"]:
|
|
current_content_buffer.append(line)
|
|
|
|
save_chunk()
|
|
return chunks
|
|
|
|
|
|
def parse_legal_document(text: str, footnotes: Optional[dict] = None) -> list:
|
|
"""법률 문서를 깊은 계층 구조로 파싱합니다."""
|
|
REGEX_SECTION = re.compile(r'^\s*제\s*\d+\s*절')
|
|
REGEX_LVL1 = re.compile(r'^\s*([1-9][0-9]?|100)\.\s+')
|
|
REGEX_LVL2 = re.compile(r'^\s*[가-하]\.\s+')
|
|
REGEX_LVL3 = re.compile(r'^\s*\(\d+\)')
|
|
REGEX_LVL4 = re.compile(r'^\s*\([가-하]\)')
|
|
REGEX_LVL5 = re.compile(r'^\s*\d+\)')
|
|
|
|
chunks = []
|
|
state = {"section": None, "lvl1": None, "lvl2": None, "lvl3": None, "lvl4": None, "lvl5": None}
|
|
current_content_buffer = []
|
|
|
|
def save_chunk():
|
|
if current_content_buffer:
|
|
content_text = "\n".join(current_content_buffer).strip()
|
|
if content_text:
|
|
path_elements = [state["section"], state["lvl1"], state["lvl2"],
|
|
state["lvl3"], state["lvl4"], state["lvl5"]]
|
|
valid_path = [p for p in path_elements if p is not None]
|
|
path_str = " > ".join(valid_path)
|
|
|
|
meta_footnotes = {}
|
|
if footnotes:
|
|
for k, v in footnotes.items():
|
|
if k in content_text:
|
|
meta_footnotes[k] = v
|
|
|
|
chunk_data = {
|
|
"metadata": {
|
|
"section": state["section"],
|
|
"level_1": state["lvl1"],
|
|
"level_2": state["lvl2"],
|
|
"level_3": state["lvl3"],
|
|
"level_4": state["lvl4"],
|
|
"level_5": state["lvl5"],
|
|
"path": path_str,
|
|
"footnotes": "\n".join(meta_footnotes.values()) if meta_footnotes else None
|
|
},
|
|
"content": content_text,
|
|
}
|
|
chunk_data["metadata"] = {k: v for k, v in chunk_data["metadata"].items() if v is not None}
|
|
chunks.append(chunk_data)
|
|
current_content_buffer.clear()
|
|
|
|
lines = text.split('\n')
|
|
for line in lines:
|
|
stripped_line = line.strip()
|
|
if not stripped_line:
|
|
continue
|
|
|
|
if REGEX_SECTION.match(stripped_line):
|
|
save_chunk()
|
|
state["section"] = stripped_line
|
|
state["lvl1"] = state["lvl2"] = state["lvl3"] = state["lvl4"] = state["lvl5"] = None
|
|
current_content_buffer.append(stripped_line)
|
|
elif REGEX_LVL1.match(stripped_line):
|
|
save_chunk()
|
|
state["lvl1"] = stripped_line
|
|
state["lvl2"] = state["lvl3"] = state["lvl4"] = state["lvl5"] = None
|
|
current_content_buffer.append(stripped_line)
|
|
elif REGEX_LVL2.match(stripped_line):
|
|
save_chunk()
|
|
state["lvl2"] = stripped_line
|
|
state["lvl3"] = state["lvl4"] = state["lvl5"] = None
|
|
current_content_buffer.append(stripped_line)
|
|
elif REGEX_LVL3.match(stripped_line):
|
|
save_chunk()
|
|
state["lvl3"] = stripped_line
|
|
state["lvl4"] = state["lvl5"] = None
|
|
current_content_buffer.append(stripped_line)
|
|
elif REGEX_LVL4.match(stripped_line):
|
|
save_chunk()
|
|
state["lvl4"] = stripped_line
|
|
state["lvl5"] = None
|
|
current_content_buffer.append(stripped_line)
|
|
elif REGEX_LVL5.match(stripped_line):
|
|
save_chunk()
|
|
state["lvl5"] = stripped_line
|
|
current_content_buffer.append(stripped_line)
|
|
else:
|
|
current_content_buffer.append(stripped_line)
|
|
|
|
save_chunk()
|
|
return chunks
|
|
|
|
|
|
def extract_text_from_pdf(pdf_path: str) -> str:
|
|
"""PDF에서 텍스트를 추출합니다."""
|
|
if not HAS_FITZ:
|
|
raise ImportError("PyMuPDF(fitz) is required for PDF processing")
|
|
text = ""
|
|
with fitz.open(pdf_path) as doc:
|
|
for page in doc:
|
|
text += page.get_text()
|
|
return text
|
|
|
|
|
|
def parse_legal_text(text: str) -> Dict[str, str]:
|
|
"""법률 텍스트(PDF 판례)를 파싱합니다."""
|
|
|
|
# (?!\S) : 바로 뒤에 "공백이 아닌 문자"가 오면 매칭하지 않음
|
|
# => '이유로', '주문에', '참조조문은' 같은 접미사가 붙은 줄을 헤더로 오인식 방지
|
|
patterns = {
|
|
'판시사항': r'^\s*판\s*시\s*사\s*항(?!\S)',
|
|
'판결요지': r'^\s*판\s*결\s*요\s*지(?!\S)',
|
|
'참조조문': r'^\s*참\s*조\s*조\s*문(?!\S)',
|
|
'참조판례': r'^\s*참\s*조\s*판\s*례(?!\S)',
|
|
'따름판례': r'^\s*따\s*름\s*판\s*례(?!\S)',
|
|
'전문': r'^\s*전\s*문(?!\S)',
|
|
'주문': r'^\s*주\s*문(?!\S)',
|
|
'청구취지 및 항소취지': r'^\s*청\s*구\s*취\s*지\s*및\s*항\s*소\s*취\s*지(?!\S)',
|
|
'청구취지': r'^\s*청\s*구\s*취\s*지(?!\S)',
|
|
'항소취지': r'^\s*항\s*소\s*취\s*지(?!\S)',
|
|
'이유': r'^\s*이\s*유(?!\S)',
|
|
}
|
|
|
|
found_sections: List[dict] = []
|
|
for key, pattern in patterns.items():
|
|
for m in re.finditer(pattern, text, re.MULTILINE):
|
|
found_sections.append({
|
|
'key': key,
|
|
'start': m.start(),
|
|
'end': m.end()
|
|
})
|
|
|
|
# 시작 위치가 같은 헤더가 여러 개 잡히면(예: '청구취지' vs '청구취지 및 항소취지')
|
|
# 더 긴(구체적인) 헤더를 우선으로 하나만 남김
|
|
found_sections.sort(key=lambda x: (x['start'], -(x['end'] - x['start'])))
|
|
deduped = []
|
|
last_start = None
|
|
for sec in found_sections:
|
|
if sec['start'] == last_start:
|
|
continue
|
|
deduped.append(sec)
|
|
last_start = sec['start']
|
|
found_sections = deduped
|
|
|
|
parsed_data: Dict[str, str] = {}
|
|
|
|
if found_sections:
|
|
first_section_start = found_sections[0]['start']
|
|
header_content = text[:first_section_start].strip()
|
|
if header_content:
|
|
parsed_data['판결정보'] = header_content
|
|
else:
|
|
parsed_data['판결정보'] = text.strip()
|
|
return parsed_data
|
|
|
|
for i, current in enumerate(found_sections):
|
|
key = current['key']
|
|
content_start = current['end']
|
|
content_end = found_sections[i + 1]['start'] if i < len(found_sections) - 1 else len(text)
|
|
content = text[content_start:content_end].strip()
|
|
|
|
if key in parsed_data and content:
|
|
parsed_data[key] += "\n" + content
|
|
else:
|
|
parsed_data[key] = content
|
|
|
|
return parsed_data
|
|
|
|
|
|
def create_weaviate_object(parsed_data: dict) -> dict:
|
|
"""파싱된 데이터에서 Weaviate 객체를 생성합니다."""
|
|
field_mapping = {
|
|
"raw_case_id": "판결정보",
|
|
"raw_case_holding_and_reasoning": "판시사항",
|
|
"raw_case_summary_of_decision": "판결요지",
|
|
"raw_case_reference_statutes": "참조조문",
|
|
"raw_case_reference_cases": "참조판례",
|
|
"raw_case_full_header": "전문",
|
|
"raw_case_ruling_order": "주문",
|
|
"raw_case_reasoning": "이유"
|
|
}
|
|
|
|
properties_to_insert = {}
|
|
for weaviate_key, parse_key in field_mapping.items():
|
|
value = parsed_data.get(parse_key)
|
|
if value and value.strip():
|
|
properties_to_insert[weaviate_key] = value.strip()
|
|
|
|
claim_appeal_value = parsed_data.get("청구취지 및 항소취지")
|
|
if not claim_appeal_value:
|
|
c = parsed_data.get("청구취지", "")
|
|
a = parsed_data.get("항소취지", "")
|
|
combined = (c + "\n" + a).strip()
|
|
if combined:
|
|
claim_appeal_value = combined
|
|
|
|
if claim_appeal_value and claim_appeal_value.strip():
|
|
properties_to_insert["raw_case_claim_and_appeal"] = claim_appeal_value.strip()
|
|
|
|
return properties_to_insert
|
|
|
|
|
|
def create_weaviate_objects_with_chunking(parsed_data: dict, max_tokens: int) -> list:
|
|
"""토큰 제한을 넘지 않도록 청킹하여 Weaviate 객체 리스트를 반환합니다."""
|
|
field_mapping = {
|
|
"raw_case_id": "판결정보",
|
|
"raw_case_holding_and_reasoning": "판시사항",
|
|
"raw_case_summary_of_decision": "판결요지",
|
|
"raw_case_reference_statutes": "참조조문",
|
|
"raw_case_reference_cases": "참조판례",
|
|
"raw_case_full_header": "전문",
|
|
"raw_case_ruling_order": "주문",
|
|
"raw_case_reasoning": "이유"
|
|
}
|
|
|
|
base_object = {}
|
|
|
|
claim_appeal = parsed_data.get("청구취지 및 항소취지")
|
|
if not claim_appeal:
|
|
c = parsed_data.get("청구취지", "")
|
|
a = parsed_data.get("항소취지", "")
|
|
combined = (c + "\n" + a).strip()
|
|
if combined:
|
|
claim_appeal = combined
|
|
|
|
if claim_appeal:
|
|
base_object["raw_case_claim_and_appeal"] = claim_appeal
|
|
|
|
for weaviate_key, parse_key in field_mapping.items():
|
|
val = parsed_data.get(parse_key)
|
|
if val and val.strip():
|
|
base_object[weaviate_key] = val.strip()
|
|
|
|
total_text_content = " ".join(base_object.values())
|
|
total_token_count = count_tokens(total_text_content)
|
|
|
|
if total_token_count <= max_tokens:
|
|
base_object["chunk_index"] = 0
|
|
return [base_object]
|
|
|
|
print(f" - 토큰 수({total_token_count})가 {max_tokens}를 초과하여 분할합니다.")
|
|
|
|
metadata_keys = ["raw_case_id", "raw_case_ruling_order", "raw_case_holding_and_reasoning"]
|
|
metadata_obj = {k: base_object[k] for k in metadata_keys if k in base_object}
|
|
|
|
metadata_tokens = count_tokens(" ".join(metadata_obj.values()))
|
|
available_tokens_per_chunk = max_tokens - metadata_tokens - 100
|
|
|
|
if available_tokens_per_chunk < 500:
|
|
full_text_chunks = split_text_by_tokens(total_text_content, max_tokens)
|
|
result_objects = []
|
|
for idx, text_chunk in enumerate(full_text_chunks):
|
|
result_objects.append({
|
|
"raw_case_id": base_object.get("raw_case_id", "Unknown"),
|
|
"raw_case_reasoning": text_chunk,
|
|
"chunk_index": idx
|
|
})
|
|
return result_objects
|
|
|
|
long_content_keys = [k for k in base_object.keys() if k not in metadata_keys]
|
|
|
|
long_text_combined = ""
|
|
for k in long_content_keys:
|
|
long_text_combined += f"[{k}] {base_object[k]}\n"
|
|
|
|
text_chunks = split_text_by_tokens(long_text_combined, available_tokens_per_chunk)
|
|
|
|
result_objects = []
|
|
for idx, text_chunk in enumerate(text_chunks):
|
|
new_obj = copy.deepcopy(metadata_obj)
|
|
new_obj["raw_case_reasoning"] = text_chunk
|
|
new_obj["chunk_index"] = idx
|
|
result_objects.append(new_obj)
|
|
|
|
return result_objects
|
|
|
|
|
|
# ============================================================
|
|
# 컬렉션 생성/관리 함수
|
|
# ============================================================
|
|
def create_collection(client, name: str, schema: list, multi_tenant: bool = False):
|
|
"""컬렉션을 생성합니다. 이미 존재하면 기존 컬렉션을 반환합니다."""
|
|
if client.collections.exists(name):
|
|
print(f" - 컬렉션 '{name}'이 이미 존재합니다.")
|
|
return client.collections.get(name)
|
|
|
|
config_kwargs = {
|
|
"name": name,
|
|
"vector_config": wvc.config.Configure.Vectors.text2vec_openai(base_url="http://qwen3-embedding-8b:8000", model="Qwen/Qwen3-Embedding-8B", dimensions=4096),
|
|
# "vector_config": wvc.config.Configure.Vectors.text2vec_openai(),
|
|
"generative_config": wvc.config.Configure.Generative.openai(),
|
|
"properties": schema,
|
|
}
|
|
|
|
if multi_tenant:
|
|
config_kwargs["multi_tenancy_config"] = wvc.config.Configure.multi_tenancy(
|
|
enabled=True, auto_tenant_creation=True
|
|
)
|
|
|
|
client.collections.create(**config_kwargs)
|
|
print(f" - 컬렉션 '{name}'을 생성했습니다.")
|
|
return client.collections.get(name)
|
|
|
|
|
|
# ============================================================
|
|
# 데이터 업로드 함수들
|
|
# ============================================================
|
|
def batch_insert(tenant, objects: list, batch_size: int = BATCH_SIZE):
|
|
"""객체를 배치 단위로 나눠서 업로드합니다."""
|
|
import time
|
|
total = len(objects)
|
|
for i in range(0, total, batch_size):
|
|
batch = objects[i:i + batch_size]
|
|
tenant.data.insert_many(objects=batch)
|
|
print(f" - 배치 {i // batch_size + 1}/{(total + batch_size - 1) // batch_size} 완료 ({len(batch)}개)")
|
|
time.sleep(0.5) # 배치 간 딜레이
|
|
|
|
|
|
def upload_check_list(client, file_paths: list):
|
|
"""Check_List 데이터를 업로드합니다."""
|
|
print("\n[1/5] Check_List 업로드 시작...")
|
|
collection = create_collection(client, "Check_List", SCHEMAS["Check_List"], multi_tenant=False)
|
|
|
|
uploaded_count = 0
|
|
for file_path in file_paths:
|
|
if not os.path.exists(file_path):
|
|
print(f" - 파일 없음: {file_path}")
|
|
continue
|
|
|
|
with open(file_path, "r", encoding="utf-8") as f:
|
|
content = f.read()
|
|
|
|
file_name = os.path.splitext(os.path.basename(file_path))[0]
|
|
obj = {"check_list_name": file_name, "check_list_content": content}
|
|
|
|
try:
|
|
collection.data.insert(obj)
|
|
uploaded_count += 1
|
|
print(f" - 업로드 완료: {file_name}")
|
|
except Exception as e:
|
|
print(f" - 업로드 실패: {file_name} - {e}")
|
|
|
|
print(f" => Check_List 업로드 완료: {uploaded_count}개")
|
|
|
|
|
|
def upload_analyzed_cases(client, tenant_paths: dict):
|
|
"""Analyzed_Cases 데이터를 업로드합니다."""
|
|
print("\n[2/5] Analyzed_Cases 업로드 시작...")
|
|
collection = create_collection(client, "Analyzed_Cases", SCHEMAS["Analyzed_Cases"], multi_tenant=True)
|
|
|
|
for tenant_name, path in tenant_paths.items():
|
|
if not os.path.exists(path):
|
|
print(f" - 경로 없음: {path}")
|
|
continue
|
|
|
|
legal_objects = []
|
|
for file_path in glob(f"{path}/*.txt"):
|
|
with open(file_path, "r", encoding="utf-8") as file:
|
|
content = file.read()
|
|
parsed_chunks = parse_analyzed_cases(content)
|
|
legal_objects.extend(parsed_chunks)
|
|
|
|
if legal_objects:
|
|
try:
|
|
# 테넌트가 없으면 생성
|
|
existing_tenants = collection.tenants.get()
|
|
if tenant_name not in existing_tenants:
|
|
collection.tenants.create(tenant_name)
|
|
|
|
tenant = collection.with_tenant(tenant_name)
|
|
batch_insert(tenant, legal_objects)
|
|
print(f" - 테넌트 '{tenant_name}' 업로드 완료: {len(legal_objects)}개 객체")
|
|
except Exception as e:
|
|
print(f" - 테넌트 '{tenant_name}' 업로드 실패: {e}")
|
|
|
|
print(" => Analyzed_Cases 업로드 완료")
|
|
|
|
|
|
def upload_law_and_rules(client, tenant_paths: dict):
|
|
"""Law_and_rules 데이터를 업로드합니다."""
|
|
print("\n[3/5] Law_and_rules 업로드 시작...")
|
|
collection = create_collection(client, "Law_and_rules", SCHEMAS["Law_and_rules"], multi_tenant=True)
|
|
|
|
for base_name, path in tenant_paths.items():
|
|
if not os.path.exists(path):
|
|
print(f" - 경로 없음: {path}")
|
|
continue
|
|
|
|
for file_path in glob(f"{path}/*.txt"):
|
|
with open(file_path, "r", encoding="utf-8") as file:
|
|
content = file.read()
|
|
parsed_chunks = parse_law_and_rule(content)
|
|
|
|
file_name = os.path.splitext(os.path.basename(file_path))[0]
|
|
|
|
if parsed_chunks:
|
|
try:
|
|
existing_tenants = collection.tenants.get()
|
|
if file_name not in existing_tenants:
|
|
collection.tenants.create(file_name)
|
|
|
|
tenant = collection.with_tenant(file_name)
|
|
batch_insert(tenant, parsed_chunks)
|
|
print(f" - 테넌트 '{file_name}' 업로드 완료: {len(parsed_chunks)}개 객체")
|
|
except Exception as e:
|
|
print(f" - 테넌트 '{file_name}' 업로드 실패: {e}")
|
|
|
|
print(" => Law_and_rules 업로드 완료")
|
|
|
|
|
|
def upload_legal_books(client, tenant_paths: dict, individual_files: list = None):
|
|
"""Legal_Books 데이터를 업로드합니다."""
|
|
print("\n[4/5] Legal_Books 업로드 시작...")
|
|
collection = create_collection(client, "Legal_Books", SCHEMAS["Legal_Books"], multi_tenant=True)
|
|
|
|
# 디렉토리 기반 업로드 (기존 방식)
|
|
for base_name, path in tenant_paths.items():
|
|
if not os.path.exists(path):
|
|
print(f" - 경로 없음: {path}")
|
|
continue
|
|
|
|
for file_path in glob(f"{path}/*.txt"):
|
|
with open(file_path, "r", encoding="utf-8") as file:
|
|
content = file.read()
|
|
|
|
# 본문과 각주 분리
|
|
if "---" in content:
|
|
main_content, footnote_content = content.split("---", 1)
|
|
footnote_list = footnote_content.split("\n\n")
|
|
footnote_list = [f for f in footnote_list if f.strip()]
|
|
footnotes = {}
|
|
for footnote in footnote_list:
|
|
if ":" in footnote:
|
|
key = footnote.split(":")[0]
|
|
footnotes[key] = footnote
|
|
else:
|
|
main_content = content
|
|
footnotes = None
|
|
|
|
parsed_chunks = parse_legal_document(main_content, footnotes=footnotes)
|
|
file_name = os.path.splitext(os.path.basename(file_path))[0]
|
|
|
|
if parsed_chunks:
|
|
try:
|
|
existing_tenants = collection.tenants.get()
|
|
if file_name not in existing_tenants:
|
|
collection.tenants.create(file_name)
|
|
|
|
tenant = collection.with_tenant(file_name)
|
|
batch_insert(tenant, parsed_chunks)
|
|
print(f" - 테넌트 '{file_name}' 업로드 완료: {len(parsed_chunks)}개 객체")
|
|
except Exception as e:
|
|
print(f" - 테넌트 '{file_name}' 업로드 실패: {e}")
|
|
|
|
# 개별 파일 업로드 (파일명이 테넌트가 됨)
|
|
if individual_files:
|
|
print(" - 개별 파일 업로드 시작...")
|
|
for file_path in individual_files:
|
|
if not os.path.exists(file_path):
|
|
print(f" - 파일 없음: {file_path}")
|
|
continue
|
|
|
|
with open(file_path, "r", encoding="utf-8") as file:
|
|
content = file.read()
|
|
|
|
# 본문과 각주 분리
|
|
if "---" in content:
|
|
main_content, footnote_content = content.split("---", 1)
|
|
footnote_list = footnote_content.split("\n\n")
|
|
footnote_list = [f for f in footnote_list if f.strip()]
|
|
footnotes = {}
|
|
for footnote in footnote_list:
|
|
if ":" in footnote:
|
|
key = footnote.split(":")[0]
|
|
footnotes[key] = footnote
|
|
else:
|
|
main_content = content
|
|
footnotes = None
|
|
|
|
parsed_chunks = parse_legal_document(main_content, footnotes=footnotes)
|
|
file_name = os.path.splitext(os.path.basename(file_path))[0]
|
|
|
|
if parsed_chunks:
|
|
try:
|
|
existing_tenants = collection.tenants.get()
|
|
if file_name not in existing_tenants:
|
|
collection.tenants.create(file_name)
|
|
|
|
tenant = collection.with_tenant(file_name)
|
|
batch_insert(tenant, parsed_chunks)
|
|
print(f" - 테넌트 '{file_name}' 업로드 완료: {len(parsed_chunks)}개 객체")
|
|
except Exception as e:
|
|
print(f" - 테넌트 '{file_name}' 업로드 실패: {e}")
|
|
|
|
print(" => Legal_Books 업로드 완료")
|
|
|
|
|
|
def upload_past_cases(client, tenant_paths: dict):
|
|
"""Past_Cases (PDF) 데이터를 업로드합니다."""
|
|
print("\n[5/5] Past_Cases 업로드 시작...")
|
|
|
|
if not HAS_FITZ:
|
|
print(" - PyMuPDF가 설치되지 않아 PDF 처리를 건너뜁니다.")
|
|
return
|
|
|
|
collection = create_collection(client, "Past_Cases", SCHEMAS["Past_Cases"], multi_tenant=True)
|
|
|
|
for tenant_name, path in tenant_paths.items():
|
|
if not os.path.exists(path):
|
|
print(f" - 경로 없음: {path}")
|
|
continue
|
|
|
|
legal_objects = []
|
|
for file_path in glob(f"{path}/*.pdf"):
|
|
try:
|
|
text = extract_text_from_pdf(file_path)
|
|
parsed_data = parse_legal_text(text)
|
|
|
|
if count_tokens(text) < MAX_TOKENS:
|
|
legal_object = create_weaviate_object(parsed_data)
|
|
legal_objects.append(legal_object)
|
|
else:
|
|
chunked_legal_objects = create_weaviate_objects_with_chunking(parsed_data, MAX_TOKENS)
|
|
legal_objects.extend(chunked_legal_objects)
|
|
|
|
print(f" - PDF 파싱 완료: {os.path.basename(file_path)}")
|
|
except Exception as e:
|
|
print(f" - PDF 파싱 실패: {os.path.basename(file_path)} - {e}")
|
|
|
|
if legal_objects:
|
|
try:
|
|
existing_tenants = collection.tenants.get()
|
|
if tenant_name not in existing_tenants:
|
|
collection.tenants.create(tenant_name)
|
|
|
|
tenant = collection.with_tenant(tenant_name)
|
|
batch_insert(tenant, legal_objects)
|
|
print(f" - 테넌트 '{tenant_name}' 업로드 완료: {len(legal_objects)}개 객체")
|
|
except Exception as e:
|
|
print(f" - 테넌트 '{tenant_name}' 업로드 실패: {e}")
|
|
|
|
print(" => Past_Cases 업로드 완료")
|
|
|
|
|
|
# ============================================================
|
|
# DB 초기화 함수
|
|
# ============================================================
|
|
def reset_database(client):
|
|
"""모든 컬렉션을 삭제하여 데이터베이스를 초기화합니다."""
|
|
print("\n[DB 초기화] 기존 컬렉션 삭제 중...")
|
|
|
|
collections_to_delete = ["Check_List", "Analyzed_Cases", "Law_and_rules", "Legal_Books", "Past_Cases"]
|
|
|
|
for name in collections_to_delete:
|
|
if client.collections.exists(name):
|
|
try:
|
|
client.collections.delete(name)
|
|
print(f" - 컬렉션 '{name}' 삭제 완료")
|
|
except Exception as e:
|
|
print(f" - 컬렉션 '{name}' 삭제 실패: {e}")
|
|
else:
|
|
print(f" - 컬렉션 '{name}' 존재하지 않음 (건너뜀)")
|
|
|
|
print(" => DB 초기화 완료\n")
|
|
|
|
|
|
# ============================================================
|
|
# 메인 함수
|
|
# ============================================================
|
|
def main():
|
|
"""모든 데이터를 Weaviate에 업로드합니다."""
|
|
print("=" * 60)
|
|
print("Weaviate 통합 데이터 업로드 시작")
|
|
print("=" * 60)
|
|
|
|
# Weaviate 클라이언트 연결
|
|
print(f"\nWeaviate 서버 연결 중... ({WEAVIATE_HOST}:{WEAVIATE_PORT})")
|
|
try:
|
|
client = weaviate.connect_to_local(
|
|
host=WEAVIATE_HOST,
|
|
port=WEAVIATE_PORT,
|
|
additional_config=wvc.init.AdditionalConfig(
|
|
timeout=wvc.init.Timeout(init=60, query=300, insert=600) # insert 10분
|
|
)
|
|
)
|
|
print("연결 성공!")
|
|
except Exception as e:
|
|
print(f"연결 실패: {e}")
|
|
return
|
|
|
|
try:
|
|
# 0. DB 초기화 (기존 컬렉션 삭제)
|
|
reset_database(client)
|
|
|
|
# 1. Check_List 업로드
|
|
# upload_check_list(client, DATA_PATHS["check_list"])
|
|
|
|
# 2. Analyzed_Cases 업로드
|
|
upload_analyzed_cases(client, DATA_PATHS["analyzed_cases"])
|
|
|
|
# 3. Law_and_rules 업로드
|
|
upload_law_and_rules(client, DATA_PATHS["law_and_rules"])
|
|
|
|
# 4. Legal_Books 업로드
|
|
upload_legal_books(client, DATA_PATHS["legal_books"], DATA_PATHS.get("legal_books_files"))
|
|
|
|
# 5. Past_Cases 업로드
|
|
upload_past_cases(client, DATA_PATHS["past_cases"])
|
|
|
|
print("\n" + "=" * 60)
|
|
print("모든 데이터 업로드 완료!")
|
|
print("=" * 60)
|
|
|
|
# 결과 요약
|
|
print("\n[결과 요약]")
|
|
for name in ["Check_List", "Analyzed_Cases", "Law_and_rules", "Legal_Books", "Past_Cases"]:
|
|
if client.collections.exists(name):
|
|
collection = client.collections.get(name)
|
|
print(f" - {name}: 존재함")
|
|
else:
|
|
print(f" - {name}: 없음")
|
|
|
|
finally:
|
|
client.close()
|
|
print("\nWeaviate 연결 종료")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
main()
|