v1.0.178: classify → subprocess.Popen (classify_worker.py). Изолированный процесс — свои коннекты к БД. HTTP-сервер продолжает отвечать.
This commit is contained in:
@@ -0,0 +1,33 @@
|
|||||||
|
"""
|
||||||
|
classify_worker.py — Фоновый процесс классификации (Фаза async classify).
|
||||||
|
|
||||||
|
Запускается через subprocess.Popen из convert_server.py.
|
||||||
|
Принимает batch_id как аргумент командной строки.
|
||||||
|
Не зависит от HTTP-сервера — свои коннекты к БД, своя память.
|
||||||
|
|
||||||
|
ЗАЧЕМ: threading.Thread внутри HTTP-процесса делит коннекты к БД
|
||||||
|
с HTTP-потоком. При 50+ файлах ThreadPoolExecutor(4) + HTTP-поток
|
||||||
|
исчерпывают пул коннектов → сервер не отвечает → 502.
|
||||||
|
|
||||||
|
subprocess.Popen создаёт отдельный питон-процесс со своей памятью
|
||||||
|
и своими коннектами. HTTP-сервер продолжает отвечать на batch-progress.
|
||||||
|
Процесс живёт пока classify не завершится, потом умирает.
|
||||||
|
"""
|
||||||
|
|
||||||
|
import sys, os
|
||||||
|
sys.path.insert(0, os.path.dirname(__file__))
|
||||||
|
|
||||||
|
from services.classify import classify_batch
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
if len(sys.argv) < 2:
|
||||||
|
print("Usage: python3 classify_worker.py <batch_id>", file=sys.stderr)
|
||||||
|
sys.exit(1)
|
||||||
|
|
||||||
|
batch_id = sys.argv[1]
|
||||||
|
try:
|
||||||
|
result = classify_batch(batch_id)
|
||||||
|
print(f"classify_worker DONE: {result}")
|
||||||
|
except Exception as e:
|
||||||
|
print(f"classify_worker FAILED: {e}", file=sys.stderr)
|
||||||
|
sys.exit(1)
|
||||||
@@ -3,7 +3,7 @@
|
|||||||
from http.server import HTTPServer, BaseHTTPRequestHandler
|
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||||||
from socketserver import ThreadingMixIn, TCPServer
|
from socketserver import ThreadingMixIn, TCPServer
|
||||||
from urllib.parse import urlparse, parse_qs
|
from urllib.parse import urlparse, parse_qs
|
||||||
import json, re, os
|
import json, re, os, sys
|
||||||
|
|
||||||
# ── DB auto-seed ──────────────────────────────────────────────────────────
|
# ── DB auto-seed ──────────────────────────────────────────────────────────
|
||||||
from db import prompts as db_prompts
|
from db import prompts as db_prompts
|
||||||
@@ -234,24 +234,37 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
# ── POST /classify-batch ─────────────────────────────────────────────
|
# ── POST /classify-batch ─────────────────────────────────────────────
|
||||||
|
|
||||||
def _handle_classify_batch(self):
|
def _handle_classify_batch(self):
|
||||||
|
"""
|
||||||
|
POST /api/classify-batch — запустить классификацию в фоне.
|
||||||
|
|
||||||
|
ЗАЧЕМ subprocess вместо threading.Thread:
|
||||||
|
ThreadPoolExecutor(4) + HTTP-поток делят коннекты к PostgreSQL.
|
||||||
|
При 50+ файлах пул исчерпывается → сервер не отвечает → 502.
|
||||||
|
Отдельный процесс — свои коннекты, своя память.
|
||||||
|
HTTP-сервер продолжает отвечать на batch-progress.
|
||||||
|
"""
|
||||||
length = int(self.headers.get("Content-Length", 0))
|
length = int(self.headers.get("Content-Length", 0))
|
||||||
body = json.loads(self.rfile.read(length))
|
body = json.loads(self.rfile.read(length))
|
||||||
batch_id = body.get("batch_id")
|
batch_id = body.get("batch_id")
|
||||||
if not batch_id:
|
if not batch_id:
|
||||||
self._json({"ok": False, "error": "batch_id required"}, 400)
|
self._json({"ok": False, "error": "batch_id required"}, 400)
|
||||||
return
|
return
|
||||||
from services.classify import classify_batch
|
|
||||||
from db import documents as db_docs
|
from db import documents as db_docs
|
||||||
import threading
|
import subprocess, os
|
||||||
|
|
||||||
# Посчитать сколько файлов — ответить сразу, классифицировать в фоне
|
# Посчитать сколько файлов — ответить сразу, classify в отдельном процессе
|
||||||
pending = db_docs.list_pending(batch_id)
|
pending = db_docs.list_pending(batch_id)
|
||||||
total = len(pending)
|
total = len(pending)
|
||||||
if total == 0:
|
if total == 0:
|
||||||
self._json({"ok": False, "error": "no pending documents"}, 400)
|
self._json({"ok": False, "error": "no pending documents"}, 400)
|
||||||
return
|
return
|
||||||
|
|
||||||
threading.Thread(target=classify_batch, args=(batch_id,), daemon=True).start()
|
# Запустить classify_worker.py в отдельном процессе
|
||||||
|
worker_path = os.path.join(os.path.dirname(__file__), "classify_worker.py")
|
||||||
|
subprocess.Popen(
|
||||||
|
[sys.executable, worker_path, batch_id],
|
||||||
|
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL,
|
||||||
|
)
|
||||||
self._json({"ok": True, "total": total}, 202)
|
self._json({"ok": True, "total": total}, 202)
|
||||||
|
|
||||||
# ── POST /apply-groups ───────────────────────────────────────────────
|
# ── POST /apply-groups ───────────────────────────────────────────────
|
||||||
|
|||||||
Reference in New Issue
Block a user