Commit 70774657 by 管志勇

注释改为英文,防止乱码

parent 1e617112
......@@ -8,7 +8,7 @@ import re
import os
import time
# 配置日志
# Configure logging
log_path = os.path.join(os.path.dirname(__file__), 'sync_log_common.log')
logging.basicConfig(
filename=log_path,
......@@ -18,16 +18,16 @@ logging.basicConfig(
)
# 定义替换print语句的函数
# Define a function to replace print statements with logging
def log_message(message):
logging.info(message)
print(message)
# 配置参数(需根据实际环境修改)
DIFY_API_BASE_URL = 'http://122.112.204.51/v1/datasets'
DATASET_ID = 'a6f4574f-9fac-4b08-ad2f-673f8533e33d'
API_KEY = 'dataset-YmAQDZbuEq5c4Q2KhSV7afyY'
# Configuration parameters (modify according to your environment)
DIFY_API_BASE_URL = 'http://192.168.141.145/v1/datasets'
DATASET_ID = 'ad4b7b26-e12b-4e8f-9d40-e927d57baa08'
API_KEY = 'dataset-eETdI20aVtQwr9w89Fnsu9uJ'
DB_CONFIG = {
'host': 'localhost',
'user': 'root',
......@@ -36,21 +36,21 @@ DB_CONFIG = {
'charset': 'utf8mb4',
'cursorclass': pymysql.cursors.DictCursor,
}
RERANK_MAX_LENGTH = 510 # 分段最大长度(包括alarm_name)
OVERLAP_LENGTH = 50 # 分段重叠长度
RERANK_MAX_LENGTH = 510 # Max segment length (including alarm_name)
OVERLAP_LENGTH = 0 # Segment overlap length
SYNC_STATE_FILE = os.path.join(os.path.dirname(__file__), 'sync_state_common.txt')
DOCUMENT_READY_TIMEOUT = 10 # 创建就绪超时时间(秒)
POLLING_INTERVAL = 1 # 轮询间隔(秒)
DOCUMENT_READY_TIMEOUT = 10 # Document ready timeout (seconds)
POLLING_INTERVAL = 1 # Polling interval (seconds)
# 增强的时间戳管理
# Enhanced timestamp management
def read_sync_state():
"""初始化同步时间戳策略"""
"""Initialize sync timestamp strategy"""
try:
with open(SYNC_STATE_FILE, 'r') as f:
content = f.read().strip()
if not content:
raise ValueError("触发全量同步初始化")
raise ValueError("Trigger full sync initialization")
return datetime.strptime(content, '%Y-%m-%d %H:%M:%S')
except (FileNotFoundError, ValueError):
initial_time = datetime.now() - timedelta(days=365 * 30)
......@@ -60,14 +60,14 @@ def read_sync_state():
def save_sync_state(timestamp):
"""原子时间戳存储"""
"""Atomic timestamp storage"""
with open(SYNC_STATE_FILE, 'w') as f:
f.write(timestamp.strftime('%Y-%m-%d %H:%M:%S'))
# 数据同步核心模块
# Core data synchronization module
def fetch_records(sync_type, last_sync):
"""动态SQL生成器"""
"""Dynamic SQL generator"""
conn = pymysql.connect(**DB_CONFIG)
try:
with conn.cursor() as cursor:
......@@ -90,85 +90,124 @@ def fetch_records(sync_type, last_sync):
def clean_solution_html(solution_html):
"""清洗solution_html,保留img标签,其他标签替换为空"""
"""Clean solution_html: keep img tags, remove all other HTML tags"""
return re.sub(r'&nbsp;|<(?!img\b)[^>]*>', '', solution_html)
def split_content(alarm_name, content):
"""将内容分段,支持重叠长度控制,每段以alarm_name开头并添加换行,总长度不超过RERANK_MAX_LENGTH"""
def split_content(alarm_name, content, labels=None):
"""Split content into segments with overlap control. Each segment starts with alarm_name followed by a newline. Total length does not exceed RERANK_MAX_LENGTH."""
segments = []
start = 0
content_length = len(content)
# 处理空报警名称的情况
# Handle empty alarm name
if not alarm_name:
alarm_name = ""
# 计算alarm_name占用的长度(包含换行符)
header_length = len(alarm_name) + 1 # +1 是换行符的长度
label_str = ",".join(labels) if labels else "无"
# 计算内容部分的最大长度
# Context template
context_prefix = f"{alarm_name}\n标签:{label_str}\n"
# Calculate the length occupied by the context prefix
header_length = len(context_prefix)
# Calculate the max length for the content portion
max_content_length = RERANK_MAX_LENGTH - header_length
# 确保重叠长度不超过最大内容长度
# Ensure overlap length does not exceed max content length
effective_overlap = min(OVERLAP_LENGTH, max_content_length)
# 如果报警名称太长,直接返回报警名称作为一个分段
# If context prefix is too long, return it as a single segment
if header_length > RERANK_MAX_LENGTH:
segments.append(alarm_name[:RERANK_MAX_LENGTH])
segments.append(context_prefix[:RERANK_MAX_LENGTH])
return segments
while start < content_length:
# 计算本次分段的结束位置
# Calculate the end position of the current segment
end = start + max_content_length
# 如果是最后一段,直接截取到末尾
# If this is the last segment, take content to the end
if end >= content_length:
segment_content = content[start:]
# 构造以alarm_name开头的分段
formatted_segment = f"{alarm_name}\n{segment_content}"
# Build segment prefixed with context
formatted_segment = f"{context_prefix}{segment_content}"
segments.append(formatted_segment)
break
# 截取当前分段
# Extract current segment
segment_content = content[start:end]
# 构造以alarm_name开头的分段
formatted_segment = f"{alarm_name}\n{segment_content}"
# Build segment prefixed with context
formatted_segment = f"{context_prefix}{segment_content}"
segments.append(formatted_segment)
# 计算下一段的起始位置(考虑重叠)
# Calculate next segment start position (considering overlap)
next_start = end - effective_overlap
# 防止重叠后起始位置回退
# Prevent start position from going backwards after overlap
if next_start <= start:
next_start = start + 1 # 至少前进1个字符,避免无限循环
next_start = start + 1 # Advance at least 1 char to avoid infinite loop
start = next_start
# 过滤空分段
# Filter out empty segments
return [s for s in segments if s.strip()]
def get_labels_by_knowledge_id(knowledge_id):
"""Get label list by knowledge_id"""
conn = pymysql.connect(**DB_CONFIG)
try:
with conn.cursor() as cursor:
sql = """
SELECT ul.title
FROM api_knowledge_open_label kol
JOIN api_userlabel ul ON kol.userlabel_id = ul.id
WHERE kol.knowledge_id = %s
"""
cursor.execute(sql, (knowledge_id,))
results = cursor.fetchall()
return [row['title'] for row in results]
finally:
conn.close()
def get_document_id_by_name(alarm_name):
"""通过文档名称查询文档ID"""
"""Query document ID by document name"""
headers = {
'Authorization': f'Bearer {API_KEY}',
'Content-Type': 'application/json'
}
url = f"{DIFY_API_BASE_URL}/{DATASET_ID}/documents?keyword={alarm_name}"
try:
response = requests.get(url, headers=headers)
log_message(f"[get_document_id_by_name] GET {url} - Status: {response.status_code}")
response.raise_for_status()
try:
data = response.json()
except requests.exceptions.JSONDecodeError as e:
log_message(f"[ERROR] JSON decode error in get_document_id_by_name: {e}")
log_message(f"[ERROR] Status Code: {response.status_code}")
log_message(f"[ERROR] Response Headers: {dict(response.headers)}")
log_message(f"[ERROR] Response Content: {response.text[:1000]}")
raise
# 遍历结果,查找精确匹配的文档名称
# Iterate results to find exact match by document name
for doc in data.get('data', []):
if doc.get('name') == alarm_name:
return doc.get('id')
return None
except requests.exceptions.RequestException as e:
log_message(f"[ERROR] Request failed in get_document_id_by_name: {e}")
if hasattr(e, 'response') and e.response is not None:
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
def create_document(alarm_name):
"""创建文档"""
"""Create a document"""
headers = {
'Authorization': f'Bearer {API_KEY}',
'Content-Type': 'application/json'
......@@ -178,16 +217,46 @@ def create_document(alarm_name):
"name": alarm_name,
"text": '',
"indexing_technique": "high_quality",
"process_rule": {"mode": "automatic"},
"process_rule": {
"mode": "custom",
"rules": {
"pre_processing_rules": [
{"id": "remove_extra_spaces", "enabled": True},
{"id": "remove_urls_emails", "enabled": False}
],
"segmentation": {
"separator": "\n",
"max_tokens": 1024
}
}
},
}
try:
response = requests.post(create_url, headers=headers, json=payload)
log_message(f"[create_document] POST {create_url} - Status: {response.status_code}")
response.raise_for_status()
return response.json()['document']['id']
try:
result = response.json()
except requests.exceptions.JSONDecodeError as e:
log_message(f"[ERROR] JSON decode error in create_document: {e}")
log_message(f"[ERROR] Status Code: {response.status_code}")
log_message(f"[ERROR] Response Headers: {dict(response.headers)}")
log_message(f"[ERROR] Response Content: {response.text[:1000]}")
raise
return result['document']['id']
except requests.exceptions.RequestException as e:
log_message(f"[ERROR] Request failed in create_document: {e}")
if hasattr(e, 'response') and e.response is not None:
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
def is_document_ready(document_id, segments, alarm_name):
"""通过尝试创建分段判断文档是否就绪,并返回测试分段ID"""
def is_document_ready(document_id, segments, alarm_name, labels=None):
"""Check if document is ready by attempting to create a segment, and return the test segment ID"""
headers = {
'Authorization': f'Bearer {API_KEY}',
'Content-Type': 'application/json'
......@@ -200,31 +269,55 @@ def is_document_ready(document_id, segments, alarm_name):
return False, None
test_segment = valid_segments[0]
# Merge labels and alarm_name as keywords
keywords = labels if labels else []
if alarm_name not in keywords:
keywords.insert(0, alarm_name)
payload = {
"segments": [{"content": test_segment, "keywords": [alarm_name]}]
"segments": [{"content": test_segment, "keywords": keywords}]
}
try:
response = requests.post(create_segment_url, headers=headers, json=payload)
log_message(f"[is_document_ready] POST {create_segment_url} - Status: {response.status_code}")
response.raise_for_status()
segment_id = response.json()['data'][0]['id']
try:
result = response.json()
except requests.exceptions.JSONDecodeError as e:
log_message(f"[ERROR] JSON decode error in is_document_ready: {e}")
log_message(f"[ERROR] Status Code: {response.status_code}")
log_message(f"[ERROR] Response Headers: {dict(response.headers)}")
log_message(f"[ERROR] Response Content: {response.text[:1000]}")
raise
segment_id = result['data'][0]['id']
log_message(f"Created test segment ID: {segment_id} for document {document_id}")
return True, segment_id
except requests.exceptions.HTTPError as e:
if e.response.status_code == 404:
log_message(f"Document {document_id} not ready yet - 404 error")
return False, None
log_message(f"[ERROR] HTTP error in is_document_ready: {e}")
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
except requests.exceptions.RequestException as e:
log_message(f"[ERROR] Request failed in is_document_ready: {e}")
if hasattr(e, 'response') and e.response is not None:
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
def wait_for_document_ready(document_id, segments, alarm_name):
"""等待文档就绪并返回测试分段ID"""
def wait_for_document_ready(document_id, segments, alarm_name, labels=None):
"""Wait for document to be ready and return test segment ID"""
start_time = time.time()
log_message(f"Waiting for document {document_id} to be ready...")
while time.time() - start_time < DOCUMENT_READY_TIMEOUT:
try:
ready, segment_id = is_document_ready(document_id, segments, alarm_name)
ready, segment_id = is_document_ready(document_id, segments, alarm_name, labels)
if ready:
test_segment_id = segment_id
log_message(f"Document {document_id} is ready, test segment ID: {test_segment_id}")
......@@ -239,8 +332,8 @@ def wait_for_document_ready(document_id, segments, alarm_name):
return False, None
def create_segments(document_id, segments, alarm_name):
"""创建文档分段,排除测试用的第一个分段"""
def create_segments(document_id, segments, alarm_name, labels=None):
"""Create document segments, excluding the first test segment"""
headers = {
'Authorization': f'Bearer {API_KEY}',
'Content-Type': 'application/json'
......@@ -252,49 +345,104 @@ def create_segments(document_id, segments, alarm_name):
log_message(f"No valid segments to create for document {document_id}")
return {"message": "No valid segments provided"}
# 等待文档就绪并获取测试分段ID
ready, test_segment_id = wait_for_document_ready(document_id, valid_segments, alarm_name)
# Wait for document to be ready and get test segment ID
ready, test_segment_id = wait_for_document_ready(document_id, valid_segments, alarm_name, labels)
if not ready:
raise Exception(f"Document {document_id} not ready after waiting")
# 插入所有分段(排除测试用的第一个分段)
# Insert all segments (excluding the first test segment)
segments_to_create = valid_segments[1:]
if not segments_to_create:
log_message(f"Only 1 segment available, no new segments to create for document {document_id}")
return {"message": "No new segments to create"}
# Merge labels and alarm_name as keywords
keywords = labels if labels else []
if alarm_name not in keywords:
keywords.insert(0, alarm_name)
payload = {
"segments": [{"content": s, "keywords": [alarm_name]} for s in segments_to_create]
"segments": [{"content": s, "keywords": keywords} for s in segments_to_create]
}
# Log API request parameters
log_message(f"[create_segments] POST {create_segment_url}")
log_message(f"[create_segments] Creating {len(segments_to_create)} segments")
try:
response = requests.post(create_segment_url, headers=headers, json=payload)
log_message(f"[create_segments] Status: {response.status_code}")
response.raise_for_status()
try:
result = response.json()
except requests.exceptions.JSONDecodeError as e:
log_message(f"[ERROR] JSON decode error in create_segments: {e}")
log_message(f"[ERROR] Status Code: {response.status_code}")
log_message(f"[ERROR] Response Headers: {dict(response.headers)}")
log_message(f"[ERROR] Response Content: {response.text[:1000]}")
raise
log_message(f"Created {len(segments_to_create)} segments for document {document_id}")
return response.json()
return result
except requests.exceptions.RequestException as e:
log_message(f"[ERROR] Request failed in create_segments: {e}")
if hasattr(e, 'response') and e.response is not None:
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
def delete_document(document_id):
"""删除文档"""
"""Delete a document"""
headers = {
'Authorization': f'Bearer {API_KEY}',
'Content-Type': 'application/json'
}
delete_url = f"{DIFY_API_BASE_URL}/{DATASET_ID}/documents/{document_id}"
try:
response = requests.delete(delete_url, headers=headers)
log_message(f"[delete_document] DELETE {delete_url} - Status: {response.status_code}")
response.raise_for_status()
# Dify API returns 204 No Content on successful deletion (no response body)
if response.status_code == 204:
log_message(f"Deleted document with ID: {document_id}")
return {"result": "success", "message": "Document deleted"}
# Handle other success codes with potential JSON response
if response.text.strip():
try:
return response.json()
except requests.exceptions.JSONDecodeError:
log_message(f"[WARNING] Non-JSON response: {response.text[:500]}")
return {"result": "success", "message": "Document deleted"}
return {"result": "success", "message": "Document deleted"}
except requests.exceptions.RequestException as e:
log_message(f"[ERROR] Request failed in delete_document: {e}")
if hasattr(e, 'response') and e.response is not None:
log_message(f"[ERROR] Status Code: {e.response.status_code}")
log_message(f"[ERROR] Response Content: {e.response.text[:1000]}")
raise
def sync_record(record):
"""同步单条记录到Dify"""
"""Sync a single record to Dify"""
alarm_name = record['alarm_name']
knowledge_id = record['id']
solution_html = clean_solution_html(record['solution_html'])
segments = split_content(alarm_name, solution_html)
log_message(f"Processing record: {alarm_name}")
# 处理删除逻辑
# Get knowledge base labels
labels = get_labels_by_knowledge_id(knowledge_id)
if labels:
log_message(f"Labels for '{alarm_name}': {labels}")
segments = split_content(alarm_name, solution_html, labels)
# Handle deletion logic
if record['is_delete'] == 1:
document_id = get_document_id_by_name(alarm_name)
if document_id:
......@@ -304,7 +452,7 @@ def sync_record(record):
log_message(f"Document not found for deletion: {alarm_name}")
return "not_found"
# 处理更新或新增逻辑
# Handle update or create logic
document_id = get_document_id_by_name(alarm_name)
if document_id:
delete_document(document_id)
......@@ -315,16 +463,16 @@ def sync_record(record):
log_message(f"Skipping document creation for '{alarm_name}' - no valid content")
return "skipped"
# 创建新文档和分段
# Create new document and segments
dify_id = create_document(alarm_name)
create_segments(dify_id, valid_segments, alarm_name)
create_segments(dify_id, valid_segments, alarm_name, labels)
log_message(f"Created new document: {alarm_name}, ID: {dify_id} with {len(valid_segments) - 1} segments")
return "recreated"
# 可视化控制模块
# Progress visualization module
def print_progress(idx, total, prefix=""):
"""终端进度可视化"""
"""Terminal progress bar visualization"""
bar_length = 40
filled = int(bar_length * idx / total)
bar = '#' * filled + '-' * (bar_length - filled)
......@@ -334,10 +482,10 @@ def print_progress(idx, total, prefix=""):
def main():
"""主控制流程"""
"""Main control flow"""
parser = argparse.ArgumentParser(description='Sync data with Dify API')
parser.add_argument('--sync-type', choices=['full', 'incremental'], default=None,
help='指定同步类型:全量或增量')
help='Specify sync type: full or incremental')
args = parser.parse_args()
last_sync = read_sync_state()
if args.sync_type:
......@@ -370,8 +518,13 @@ def main():
max_update_datetime = current_update
print_progress(idx, total, prefix=f"SYNCING")
except Exception as e:
error_log.append(f"ID:{record['alarm_name']} ERROR:{str(e)}")
log_message(f"ERROR: {record['alarm_name'][:15]}... - {str(e)}")
error_msg = f"ID:{record['id']} NAME:{record['alarm_name']}"
error_log.append(f"{error_msg} ERROR:{str(e)}")
log_message(f"[ERROR] {error_msg}")
log_message(f"[ERROR] Exception Type: {type(e).__name__}")
log_message(f"[ERROR] Exception Details: {str(e)}")
import traceback
log_message(f"[ERROR] Traceback:\n{traceback.format_exc()}")
log_message("\n" + "=" * 50)
save_sync_state(datetime.now())
log_message(
......
Markdown is supported
0% or
You are about to add 0 people to the discussion. Proceed with caution.
Finish editing this message first!
Please register or to comment