Загрузка данных
"""
Этап 1: Агентный пайплайн для ОЧИСТКИ данных.
Исправления по итогам ревью:
- убран `break`, из-за которого обрабатывалась только одна строка
- добавлен импорт `base64` (использовался в _encode_image, но не был импортирован)
- убран неиспользуемый импорт `json`
- добавлено периодическое сохранение результатов (SAVE_EVERY) на случай сбоя посреди батча
- добавлен resume: строки со статусом "ok" повторно не отправляются на сервер
- ошибки сети/парсинга пишутся в отдельный STATUS_COLUMN, а не поверх очищенного текста
- добавлены повторные попытки (retry) с задержкой при сбоях запроса
- температура снижена до 0 — задача очистки текста детерминированная, вариативность не нужна
- после удаления <think>...</think> добавлен .strip()
- поправлена нумерация пунктов в system prompt (была "3." дважды, пропущен "4.")
- перед стартом батча проверяется, что модель реально доступна на сервере
"""
import base64
import re
import os
import time
from typing import Tuple
import requests
import pandas as pd
class LLMClient:
"""Мультимодальный клиент для отправки текста (и, опционально, изображений) на LLM-сервер."""
def __init__(self, base_url: str, model_name: str, api_key: str = "none"):
self.base_url = base_url.rstrip("/")
self.server_url = self.base_url + "/v1/chat/completions"
self.models_url = self.base_url + "/v1/models"
self.model_name = model_name
self.headers = {
"Content-Type": "application/json",
"Authorization": f"Bearer {api_key}",
}
def get_models(self) -> list:
"""Получает список доступных моделей с сервера."""
try:
response = requests.get(self.models_url, headers=self.headers, timeout=30)
response.raise_for_status()
data = response.json()
if "data" in data:
return [model["id"] for model in data["data"]]
return []
except requests.exceptions.RequestException as e:
print(f"Ошибка при получении списка моделей: {e}")
return []
except (KeyError, TypeError):
print(f"Ошибка парсинга списка моделей. Ответ: {response.text}")
return []
def _encode_image(self, image_path: str) -> str:
"""Кодирует локальное изображение в Base64 (для будущей мультимодальной отправки)."""
with open(image_path, "rb") as image_file:
return base64.b64encode(image_file.read()).decode("utf-8")
def send_multimodal_prompt(
self,
system_role: str,
user_text: str,
temperature: float = 0.0,
max_retries: int = 3,
retry_delay: float = 5.0,
) -> Tuple[str, bool]:
"""
Отправляет текст на модель.
Возвращает (текст, success). При сбое после всех попыток
success=False, а текст содержит описание последней ошибки.
"""
user_content = [{"type": "text", "text": user_text}]
payload = {
"model": self.model_name,
"messages": [
{"role": "system", "content": system_role},
{"role": "user", "content": user_content},
],
"temperature": temperature,
}
last_error = "Неизвестная ошибка"
for attempt in range(1, max_retries + 1):
try:
response = requests.post(
self.server_url, headers=self.headers, json=payload, timeout=120
)
response.raise_for_status()
return response.json()["choices"][0]["message"]["content"], True
except requests.exceptions.RequestException as e:
last_error = f"Ошибка соединения с VL-моделью: {e}"
except (KeyError, IndexError):
last_error = f"Ошибка парсинга ответа. Сервер вернул: {response.text}"
if attempt < max_retries:
print(f" Попытка {attempt}/{max_retries} неудачна, повтор через {retry_delay} с...")
time.sleep(retry_delay)
return last_error, False
BASE_URL = "http://vllm-server:8000"
MODEL_NAME = "/models/Qwen3-14B-AWQ"
API_KEY = os.getenv("VLLM_API_KEY_QWEN_14B", "none")
client = LLMClient(base_url=BASE_URL, model_name=MODEL_NAME, api_key=API_KEY)
system_prompt = """
<instructions>
- ВСЕГДА следуй <answering_rules> и <self_reflection>
<self_reflection>
1. Ты эксперт в области очистки данных
2. Действуй согласно роли. Ты ДОЛЖЕН сочетать глубокие знания и ясное мышление
3. Не пытайся вникать в суть сообщения и решать проблему
4. Ты должен ТОЛЬКО очистить сообщение
</self_reflection>
<answering_rules>
1. ИСПОЛЬЗУЙ язык сообщения ПОЛЬЗОВАТЕЛЯ
2. Твой ответ очень важен для меня
3. Первой строкой выведи должность работника из подписи, если она есть
4. Удали из входящего сообщения:
- Разметку Markdown
- Изображения
- Ссылки
- URL-адреса
- Имена и фамилии
- Номера телефонов
- Различные артефакты
- Заявление о конфиденциальности
- Лишние пробелы
- Подпись
- Ненужные переносы строк
- Другой мусор
5. НИКОГДА не удаляй другой текст сообщения
6. Используй <example> для оценки <input> и <output> текста
</answering_rules>
<example>
<input>
Добрый день, проблема следующая- после согласования в ОБД скважины не появляются в ГеоНОВА. Последний раз подгружали планшет в феврали с тем же названием, все было в порядке. Грузилось 3 скважины 1м проектом. Скважины без месторождения, только ЛУ. С уважением, Лудникова Нина Викторовна Эксперт Петрофизика 16_Утренний, Штормовой, Зап. берег Восточно-Тамбейский Раб.тел. +7 (3452)684-696 (вн.22-649) [cid:image001.png@01DCF442.25A4F9E0] [Created via e-mail received from: nina.ludnikova@novatek.ru]
</input>
<output>
Проблема следующая после согласования в ОБД скважины не появляются в ГеоНОВА. Последний раз подгружали планшет в феврали с тем же названием, все было в порядке. Грузилось 3 скважины 1м проектом. Скважины без месторождения, только ЛУ. Эксперт Петрофизика
</output>
</example>
</instructions>
"""
EXCEL_INPUT = "data.xlsx" # Имя исходного файла
EXCEL_OUTPUT = "cleaned_tasks.xlsx" # Имя файла для сохранения результатов
TARGET_COLUMN = "Письмо пользователя" # Название целевого столбца
RESULT_COLUMN = "Обработанный ответ" # Столбец с очищенным текстом
STATUS_COLUMN = "Статус обработки" # "ok" / "error: ..." — не путаем данные с ошибками
SAVE_EVERY = 20 # периодическое сохранение на случай сбоя посреди батча
TEMPERATURE = 0.0 # задача детерминированная — вариативность не нужна
def strip_think_tags(text: str) -> str:
"""Убирает блок рассуждений reasoning-модели и обрезает лишние пробелы вокруг."""
cleaned = re.sub(r"<think>.*?</think>", "", text, flags=re.DOTALL)
return cleaned.strip()
def main():
try:
df = pd.read_excel(EXCEL_INPUT)
print(f"Успешно загружено строк: {len(df)}")
except Exception as e:
print(f"Ошибка при чтении Excel-файла: {e}")
return
if TARGET_COLUMN not in df.columns:
print(f"Ошибка: Столбец '{TARGET_COLUMN}' не найден в файле!")
return
if RESULT_COLUMN not in df.columns:
df[RESULT_COLUMN] = None
if STATUS_COLUMN not in df.columns:
df[STATUS_COLUMN] = None
# Проверяем, что модель реально доступна на сервере, до старта батча
available_models = client.get_models()
if available_models and MODEL_NAME not in available_models:
print(f"Внимание: модель '{MODEL_NAME}' не найдена среди доступных: {available_models}")
processed_since_save = 0
for index, row in df.iterrows():
# Resume: строки, уже успешно обработанные ранее, повторно не отправляем
if df.at[index, STATUS_COLUMN] == "ok":
continue
raw_text = row[TARGET_COLUMN]
if pd.isna(raw_text) or str(raw_text).strip() == "":
print(f"Строка {index}: пропущена (пустое значение)")
continue
print(f"Обработка строки {index + 1}/{len(df)}...")
llm_response, success = client.send_multimodal_prompt(
system_role=system_prompt,
user_text=str(raw_text),
temperature=TEMPERATURE,
)
if success:
df.at[index, RESULT_COLUMN] = strip_think_tags(llm_response)
df.at[index, STATUS_COLUMN] = "ok"
print(df.at[index, RESULT_COLUMN])
else:
df.at[index, RESULT_COLUMN] = None
df.at[index, STATUS_COLUMN] = f"error: {llm_response}"
print(f" Строка {index}: ошибка — {llm_response}")
processed_since_save += 1
if processed_since_save >= SAVE_EVERY:
try:
df.to_excel(EXCEL_OUTPUT, index=False)
print(f" Промежуточное сохранение ({index + 1}/{len(df)} строк)...")
except Exception as e:
print(f"Ошибка при промежуточном сохранении: {e}")
processed_since_save = 0
try:
df.to_excel(EXCEL_OUTPUT, index=False)
n_errors = int(df[STATUS_COLUMN].astype(str).str.startswith("error").sum())
print(f"Обработка завершена! Результаты сохранены в {EXCEL_OUTPUT}")
if n_errors:
print(f"Внимание: {n_errors} строк(и) завершились с ошибкой (см. столбец '{STATUS_COLUMN}')")
except Exception as e:
print(f"Ошибка при сохранении Excel-файла: {e}")
if __name__ == "__main__":
main()