Browse Source

采集线程防掉线、网络抖动防掉线

master
han\hanst 4 weeks ago
parent
commit
9cef5ac880
  1. 73
      client-file-collector/client_file_collector.py
  2. BIN
      client-file-collector/dist/QMSFileCollector.exe

73
client-file-collector/client_file_collector.py

@ -13,6 +13,8 @@ from tkinter import Tk, Label, Entry, Button, StringVar, END, DISABLED, NORMAL,
from tkinter.scrolledtext import ScrolledText from tkinter.scrolledtext import ScrolledText
import requests import requests
from requests.adapters import HTTPAdapter
from urllib3.util.retry import Retry
def get_runtime_dir(): def get_runtime_dir():
@ -138,6 +140,7 @@ class CollectorApp:
self.active_rows = [] self.active_rows = []
self.upload_state = {} self.upload_state = {}
self.notice_ts = {} self.notice_ts = {}
self.http_session = self._create_http_session()
self.base_url_var = StringVar() self.base_url_var = StringVar()
self.poll_seconds_var = StringVar(value="10") self.poll_seconds_var = StringVar(value="10")
@ -153,9 +156,27 @@ class CollectorApp:
}) })
self._build_ui() self._build_ui()
self.root.protocol("WM_DELETE_WINDOW", self._on_window_close)
self._init_windows_startup() self._init_windows_startup()
self._load_config() self._load_config()
self.root.after(200, self._flush_logs) self.root.after(200, self._flush_logs)
self.root.after(3000, self._watch_worker)
def _create_http_session(self):
retry = Retry(
total=3,
connect=3,
read=3,
backoff_factor=1.0,
status_forcelist=[500, 502, 503, 504],
allowed_methods=frozenset(["POST"]),
raise_on_status=False,
)
adapter = HTTPAdapter(max_retries=retry, pool_connections=8, pool_maxsize=8)
session = requests.Session()
session.mount("http://", adapter)
session.mount("https://", adapter)
return session
def _init_windows_startup(self): def _init_windows_startup(self):
changed, message = ensure_windows_startup() changed, message = ensure_windows_startup()
@ -176,6 +197,7 @@ class CollectorApp:
Entry(self.root, textvariable=self.poll_seconds_var, width=10).place(x=85, y=58) Entry(self.root, textvariable=self.poll_seconds_var, width=10).place(x=85, y=58)
Button(self.root, text="开始同步", command=self.start_collect).place(x=180, y=54) Button(self.root, text="开始同步", command=self.start_collect).place(x=180, y=54)
Button(self.root, text="停止同步", command=self.stop_collect).place(x=265, y=54) Button(self.root, text="停止同步", command=self.stop_collect).place(x=265, y=54)
Button(self.root, text="退出程序", command=self.exit_app).place(x=1080, y=54)
Label(self.root, text="配置说明:最多3行,每行需填写 site + buNo + equipmentNo + 本地目录;备份目录可选(为空则自动使用本地目录_bak)").place(x=360, y=58) Label(self.root, text="配置说明:最多3行,每行需填写 site + buNo + equipmentNo + 本地目录;备份目录可选(为空则自动使用本地目录_bak)").place(x=360, y=58)
@ -360,10 +382,13 @@ class CollectorApp:
return rows return rows
def start_collect(self, auto_mode=False): def start_collect(self, auto_mode=False):
if self.running:
if self.running and self.worker_thread is not None and self.worker_thread.is_alive():
if not auto_mode: if not auto_mode:
self.log("采集任务已在运行中。") self.log("采集任务已在运行中。")
return return
if self.running and (self.worker_thread is None or not self.worker_thread.is_alive()):
self.log("检测到采集线程已退出,正在自动重启。")
self.running = False
if not self._base_url(): if not self._base_url():
if auto_mode: if auto_mode:
@ -413,13 +438,39 @@ class CollectorApp:
self.stop_event.set() self.stop_event.set()
self.log("正在停止采集任务...") self.log("正在停止采集任务...")
def _on_window_close(self):
self.root.iconify()
self._log_with_interval("window_hidden", "配置窗口已最小化到任务栏,采集仍在后台运行。", 5)
def exit_app(self):
self.running = False
self.stop_event.set()
try:
self.http_session.close()
except Exception:
pass
self.root.destroy()
def _watch_worker(self):
try:
if self.running and (self.worker_thread is None or not self.worker_thread.is_alive()):
self.log("检测到采集线程中断,正在自动恢复...")
self.running = False
self.start_collect(auto_mode=True)
finally:
self.root.after(3000, self._watch_worker)
def _worker_loop(self, stop_event): def _worker_loop(self, stop_event):
while not stop_event.is_set(): while not stop_event.is_set():
try:
rows = list(self.active_rows) rows = list(self.active_rows)
for row in rows: for row in rows:
if stop_event.is_set(): if stop_event.is_set():
break break
try:
self._sync_one_row(row, stop_event) self._sync_one_row(row, stop_event)
except Exception as row_error:
self.log("%d行同步异常,已跳过并继续: %s" % (row.get("row_index", 0), row_error))
try: try:
poll_seconds = int(self.poll_seconds_var.get().strip() or "10") poll_seconds = int(self.poll_seconds_var.get().strip() or "10")
@ -430,6 +481,12 @@ class CollectorApp:
if stop_event.is_set(): if stop_event.is_set():
break break
time.sleep(1) time.sleep(1)
except Exception as e:
self.log("采集线程异常,将自动继续: %s" % e)
for _ in range(5):
if stop_event.is_set():
break
time.sleep(1)
if self.stop_event is stop_event: if self.stop_event is stop_event:
self.running = False self.running = False
self.worker_thread = None self.worker_thread = None
@ -446,17 +503,21 @@ class CollectorApp:
return return
files = [] files = []
try:
for name in os.listdir(source_path): for name in os.listdir(source_path):
file_path = Path(source_path) / name file_path = Path(source_path) / name
if file_path.is_file(): if file_path.is_file():
files.append(file_path) files.append(file_path)
except Exception as e:
self._log_with_interval(path_key, "%d行读取目录失败: %s, 原因: %s" % (row_index, source_path, e), 30)
return
self._cleanup_row_upload_state(row_index, files) self._cleanup_row_upload_state(row_index, files)
if not files: if not files:
self._log_with_interval(empty_key, "%d行目录为空,等待新文件..." % row_index, 60) self._log_with_interval(empty_key, "%d行目录为空,等待新文件..." % row_index, 60)
return return
files.sort(key=lambda p: p.stat().st_mtime)
files.sort(key=self._safe_mtime)
for file_path in files: for file_path in files:
if stop_event.is_set(): if stop_event.is_set():
return return
@ -497,6 +558,12 @@ class CollectorApp:
self.log("读取文件状态失败: %s, 原因: %s" % (file_path, e)) self.log("读取文件状态失败: %s, 原因: %s" % (file_path, e))
return None return None
def _safe_mtime(self, file_path):
try:
return file_path.stat().st_mtime
except Exception:
return 0
def _get_backup_dir(self, row): def _get_backup_dir(self, row):
backup_path = str(row.get("backup_path", "")).strip() backup_path = str(row.get("backup_path", "")).strip()
if backup_path: if backup_path:
@ -542,7 +609,7 @@ class CollectorApp:
try: try:
with open(file_path, "rb") as fp: with open(file_path, "rb") as fp:
files = {"file": (file_path.name, fp, "application/octet-stream")} files = {"file": (file_path.name, fp, "application/octet-stream")}
resp = requests.post(url, data=data, files=files, timeout=180)
resp = self.http_session.post(url, data=data, files=files, timeout=180)
resp.raise_for_status() resp.raise_for_status()
result = resp.json() result = resp.json()

BIN
client-file-collector/dist/QMSFileCollector.exe

Loading…
Cancel
Save