235 lines
9.3 KiB
Python
235 lines
9.3 KiB
Python
# data_collector.py
|
||
"""数据采集器:管理 Episode 数据记录与上上传。
|
||
|
||
纯业务逻辑,无 UI 依赖。通过回调与 UI 层通信。
|
||
"""
|
||
|
||
import io
|
||
import json
|
||
import pickle
|
||
import datetime
|
||
import threading
|
||
import requests
|
||
|
||
from api import base_url, data_record_url, the_folder
|
||
|
||
|
||
class DataCollector:
|
||
"""管理控制过程中的 Episode 数据采集与保存"""
|
||
|
||
def __init__(self):
|
||
self.episode_data_raw = [] # 所有已完成的 Episode
|
||
self.current_episode = None # 当前正在记录的 Episode
|
||
self.last_target_record = None
|
||
self._on_log = None
|
||
self._on_upload_complete = None
|
||
|
||
def set_log_callback(self, callback):
|
||
"""设置日志回调"""
|
||
self._on_log = callback
|
||
|
||
def set_upload_complete_callback(self, callback):
|
||
"""设置控制数据上传完成回调。
|
||
|
||
callback(success, manifest, error) 会在后台上传线程中调用。成功时
|
||
manifest 是已上传的清单字典;失败时 error 为可展示的错误信息。
|
||
"""
|
||
self._on_upload_complete = callback
|
||
|
||
def log(self, message):
|
||
if self._on_log:
|
||
self._on_log(message)
|
||
|
||
def _notify_upload_complete(self, success, manifest=None, error=None):
|
||
if self._on_upload_complete:
|
||
self._on_upload_complete(success, manifest, error)
|
||
|
||
def _upload_to_server(self, data_bytes: bytes, filename: str, folder: str) -> bool:
|
||
"""向 ReinLoop 云服务器申请上传地址并上传控制数据。"""
|
||
try:
|
||
resp = requests.post(data_record_url, json={
|
||
"type": "uploadDataFile",
|
||
"fileName": filename,
|
||
"folder": folder,
|
||
}, timeout=30)
|
||
result = resp.json()
|
||
except Exception as e:
|
||
self.log(f"向云服务器申请上传地址异常: {e}")
|
||
return False
|
||
|
||
if not result.get("success"):
|
||
self.log(f"申请服务器上传地址失败: {result.get('errMsg', result)}")
|
||
return False
|
||
|
||
meta = result.get("uploadMetadata")
|
||
if not meta or "url" not in meta:
|
||
self.log("云服务器未返回有效的上传地址")
|
||
return False
|
||
|
||
try:
|
||
files = {"file": (filename, io.BytesIO(data_bytes))}
|
||
upload_resp = requests.post(meta["url"], files=files, timeout=60)
|
||
|
||
if upload_resp.status_code in [200, 204]:
|
||
return True
|
||
else:
|
||
self.log(f"云服务器上传失败,状态码: {upload_resp.status_code}")
|
||
return False
|
||
except Exception as e:
|
||
self.log(f"云服务器上传异常: {e}")
|
||
return False
|
||
|
||
def reset(self):
|
||
"""重置所有采集状态(控制启动时调用)"""
|
||
self.episode_data_raw = []
|
||
self.current_episode = None
|
||
self.last_target_record = None
|
||
|
||
def record_step(self, cycle_count: int, current_pressure: float,
|
||
target_pressure: float, valve_opening: float,
|
||
kp: float, ki: float, kd: float,
|
||
q_in: float, v: float):
|
||
"""记录一个控制周期的数据点
|
||
|
||
Args:
|
||
cycle_count: 控制周期计数
|
||
current_pressure: 当前压力
|
||
target_pressure: 目标压力
|
||
valve_opening: 阀门开度
|
||
kp, ki, kd: PID 参数
|
||
q_in: 流量
|
||
v: 容积
|
||
"""
|
||
# 目标压力变化时自动切分 Episode
|
||
if self.current_episode is None or target_pressure != self.last_target_record:
|
||
if self.current_episode is not None:
|
||
self.episode_data_raw.append(self.current_episode)
|
||
self.log(f"Episode 结束,已记录 {len(self.current_episode['pressures'])} 个点")
|
||
|
||
self.current_episode = {
|
||
'pid': [float(kp), float(ki), float(kd)],
|
||
'target_pressure': target_pressure,
|
||
'Q_in': q_in,
|
||
'V': v,
|
||
'steps': [],
|
||
'pressures': [],
|
||
'errors': [],
|
||
'valves': []
|
||
}
|
||
self.last_target_record = target_pressure
|
||
|
||
# 记录当前步数据
|
||
error = -(target_pressure - current_pressure)
|
||
self.current_episode['steps'].append(cycle_count)
|
||
self.current_episode['pressures'].append(current_pressure)
|
||
self.current_episode['errors'].append(error)
|
||
self.current_episode['valves'].append(float(valve_opening))
|
||
|
||
def finalize_and_upload(self, flow: float, vol: float):
|
||
"""停止控制时:闭合最后一个 Episode,分片上传到云存储。
|
||
|
||
单文件超过 5MB 时自动拆分为多个分片,
|
||
同时上传一个 manifest.json 记录所有分片信息。
|
||
|
||
Args:
|
||
flow: 流量值 (用于文件名/路径)
|
||
vol: 容积值 (用于文件名/路径)
|
||
"""
|
||
# 闭合最后一个 Episode
|
||
if self.current_episode and len(self.current_episode['pressures']) > 0:
|
||
self.episode_data_raw.append(self.current_episode)
|
||
self.current_episode = None
|
||
|
||
if not self.episode_data_raw:
|
||
return
|
||
|
||
# 上传在线程中继续执行,因此必须持有本轮数据快照。否则 finally
|
||
# 清空缓存后,异步线程生成的 manifest 会错误地显示 0 个 Episode。
|
||
episodes = list(self.episode_data_raw)
|
||
|
||
try:
|
||
timestamp = datetime.datetime.now().strftime('%Y%m%d_%H%M%S_%f')
|
||
base_folder = f"{the_folder}/data_record/data_{flow}SLM_{vol}L"
|
||
|
||
# 拆成每片尽量不超过 5MB 的 episode 分组
|
||
MAX_CHUNK_BYTES = 5 * 1024 * 1024 # 5MB
|
||
|
||
chunks = [] # [(chunk_index, episodes_subset)]
|
||
current_chunk = []
|
||
for ep in episodes:
|
||
current_chunk.append(ep)
|
||
if pickle.dumps(current_chunk).__len__() >= MAX_CHUNK_BYTES:
|
||
# 当前片已满,回退一个 episode 后保存
|
||
current_chunk.pop()
|
||
# 单个 Episode 也可能超过 5 MB;此时仍上传该 Episode,
|
||
# 而不是产生一个无内容的空分片。
|
||
if current_chunk:
|
||
chunks.append(current_chunk)
|
||
current_chunk = [ep]
|
||
if current_chunk:
|
||
chunks.append(current_chunk)
|
||
|
||
total_chunks = len(chunks)
|
||
self.log(f"控制数据共 {len(episodes)} 个 Episode,"
|
||
f"拆为 {total_chunks} 个分片上传")
|
||
|
||
def upload_all():
|
||
part_files = []
|
||
part_metadata = []
|
||
for idx, chunk_eps in enumerate(chunks):
|
||
data_bytes = pickle.dumps(chunk_eps)
|
||
size_kb = len(data_bytes) / 1024
|
||
part_filename = f'episode_raw_data_{timestamp}_part{idx + 1}of{total_chunks}.pkl'
|
||
self.log(f" 上传分片 {idx + 1}/{total_chunks} ({size_kb:.0f} KB)...")
|
||
if self._upload_to_server(data_bytes, part_filename, base_folder):
|
||
part_files.append(part_filename)
|
||
part_metadata.append({
|
||
"file_name": part_filename,
|
||
"episode_count": len(chunk_eps),
|
||
"size_bytes": len(data_bytes),
|
||
})
|
||
else:
|
||
self.log(f" 分片 {idx + 1} 上传失败")
|
||
|
||
# 上传 manifest
|
||
manifest = {
|
||
"schema_version": 1,
|
||
"data_type": "control_episode",
|
||
"run_id": timestamp,
|
||
"timestamp": timestamp,
|
||
"folder": base_folder,
|
||
"total_chunks": total_chunks,
|
||
"uploaded_chunks": len(part_files),
|
||
"part_files": part_files,
|
||
"parts": part_metadata,
|
||
"total_episodes": len(episodes),
|
||
"flow": flow,
|
||
"volume": vol,
|
||
}
|
||
manifest_str = json.dumps(manifest, indent=2, ensure_ascii=False)
|
||
manifest_bytes = manifest_str.encode('utf-8')
|
||
manifest_filename = f'episode_raw_data_{timestamp}_manifest.json'
|
||
manifest_uploaded = self._upload_to_server(
|
||
manifest_bytes, manifest_filename, base_folder
|
||
)
|
||
|
||
if len(part_files) == total_chunks and manifest_uploaded:
|
||
self.log(f"控制数据上传成功 ({total_chunks} 个分片)")
|
||
self._notify_upload_complete(True, manifest, None)
|
||
else:
|
||
error = (
|
||
"控制数据清单上传失败"
|
||
if not manifest_uploaded
|
||
else f"控制数据部分上传失败 ({len(part_files)}/{total_chunks})"
|
||
)
|
||
self.log(error)
|
||
self._notify_upload_complete(False, manifest, error)
|
||
|
||
threading.Thread(target=upload_all, daemon=True).start()
|
||
|
||
except Exception as e:
|
||
self.log(f"保存收集数据时发生错误: {e}")
|
||
self._notify_upload_complete(False, None, str(e))
|
||
finally:
|
||
self.episode_data_raw = []
|