Files

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()