Files
ReinLoopTest/ReinLoop/core/data_collector.py
2026-08-03 11:16:49 +08:00

234 lines
9.2 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# data_collector.py
"""数据采集器:管理 Episode 数据记录与上上传。
纯业务逻辑,无 UI 依赖。通过回调与 UI 层通信。
"""
import io
import json
import pickle
import datetime
import threading
import requests
from api import device_post, 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:
result = device_post({
"type": "uploadDataFile",
"fileName": filename,
"folder": folder,
}, timeout=30)
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 = []