Compare commits
8 Commits
51deee1493
...
736d644792
| Author | SHA1 | Date | |
|---|---|---|---|
| 736d644792 | |||
| 72bab25f3f | |||
| 8521aee7f9 | |||
| 783c633a15 | |||
| 15f79e69ad | |||
| d66e2826e1 | |||
| 1412b0ed0b | |||
| c745f0c070 |
36
.gitignore
vendored
Normal file
36
.gitignore
vendored
Normal file
@ -0,0 +1,36 @@
|
||||
# =========================================================
|
||||
# 运行时数据 —— 绝不入库
|
||||
# instance/ 里是 SQLite 生产库(已 36MB+)及其备份,
|
||||
# 一旦提交会让仓库体积失控,且会覆盖服务器上的真实数据。
|
||||
# =========================================================
|
||||
instance/
|
||||
*.db
|
||||
*.db-journal
|
||||
*.db-wal
|
||||
*.db-shm
|
||||
*.sqlite3
|
||||
|
||||
# =========================================================
|
||||
# 构建产物(可由源码重新生成,不入库)
|
||||
# =========================================================
|
||||
__pycache__/
|
||||
*.py[cod]
|
||||
*$py.class
|
||||
build/
|
||||
build_new/
|
||||
dist/
|
||||
dist_new/
|
||||
web_dist/
|
||||
*.exe
|
||||
|
||||
# =========================================================
|
||||
# 临时与备份
|
||||
# =========================================================
|
||||
*.bak
|
||||
*.bak-*
|
||||
*.log
|
||||
|
||||
# =========================================================
|
||||
# 编辑器
|
||||
# =========================================================
|
||||
.vscode/
|
||||
134
2_1banben/app.py
134
2_1banben/app.py
@ -1,4 +1,5 @@
|
||||
import os
|
||||
import io
|
||||
import sys
|
||||
import json
|
||||
import mimetypes
|
||||
@ -10,6 +11,41 @@ from flask import Flask, send_from_directory, jsonify
|
||||
from flask_cors import CORS
|
||||
from flask_apscheduler import APScheduler
|
||||
|
||||
# ==============================================================================
|
||||
# ✅ 0. 控制台编码兜底(必须最先执行)
|
||||
# ==============================================================================
|
||||
# 保留原始流引用,防止被 GC 回收时连带关闭底层 buffer
|
||||
_CONSOLE_ORIGINALS = []
|
||||
|
||||
|
||||
def _force_utf8_console():
|
||||
"""
|
||||
强制 stdout/stderr 以 UTF-8 输出,防止 Windows GBK 终端下 emoji 崩溃。
|
||||
|
||||
Windows 中文控制台默认 cp936(GBK),而本应用(app.py / services.core /
|
||||
crawler_106 / crawler_82)大量使用 emoji 打印日志。一旦 print 抛
|
||||
UnicodeEncodeError,异常会落进采集任务的 try 块,导致整个事务被 rollback
|
||||
并报"数据写入失败" —— 表现就是定时采集长期静默不落库。
|
||||
"""
|
||||
for name in ('stdout', 'stderr'):
|
||||
stream = getattr(sys, name, None)
|
||||
# PyInstaller --noconsole 或输出重定向时,stream 或其 buffer 可能不存在
|
||||
if stream is None or not hasattr(stream, 'buffer'):
|
||||
continue
|
||||
enc = (getattr(stream, 'encoding', None) or '').lower().replace('-', '')
|
||||
if enc == 'utf8':
|
||||
continue
|
||||
# 替换前先 flush:否则旧流缓冲区里尚未写出的内容会随旧对象一起丢掉
|
||||
try:
|
||||
stream.flush()
|
||||
except Exception:
|
||||
pass
|
||||
_CONSOLE_ORIGINALS.append(stream)
|
||||
setattr(sys, name, io.TextIOWrapper(stream.buffer, encoding='utf-8', errors='replace'))
|
||||
|
||||
|
||||
_force_utf8_console()
|
||||
|
||||
# ==============================================================================
|
||||
# ✅ 1. 核心模块引用
|
||||
# ==============================================================================
|
||||
@ -19,6 +55,8 @@ try:
|
||||
from models import Device, DeviceHistory
|
||||
# 引入核心爬虫调度
|
||||
from services.core import execute_monitor_task
|
||||
# 引入统一入库管道
|
||||
from services.db_ingest import ingest_device_data
|
||||
|
||||
try:
|
||||
from services.iot_api import sync_iot_data_service
|
||||
@ -80,14 +118,13 @@ mimetypes.add_type('text/css', '.css')
|
||||
# ==============================================================================
|
||||
def auto_monitor_job(app):
|
||||
"""
|
||||
[关键修复]
|
||||
1. 使用 app.app_context() 确保线程中有 Flask 上下文
|
||||
2. 使用 db.session.remove() 强制清理旧连接
|
||||
3. 使用 db.session.merge() 确保对象状态被正确追踪
|
||||
4. 增加详细日志,对比爬虫返回的数据与入库行为
|
||||
每天的定时采集任务。
|
||||
|
||||
入库细节全部收敛到 services.db_ingest.ingest_device_data,本函数只负责:
|
||||
建立应用上下文 -> 触发爬虫 -> 调用入库管道 -> 提交事务。
|
||||
"""
|
||||
with app.app_context():
|
||||
# A. 强制清理会话,确保线程获取的是全新的数据库连接
|
||||
# 强制清理会话,确保线程获取的是全新的数据库连接
|
||||
db.session.remove()
|
||||
|
||||
tz = pytz.timezone('Asia/Shanghai')
|
||||
@ -101,7 +138,6 @@ def auto_monitor_job(app):
|
||||
return
|
||||
|
||||
try:
|
||||
# B. 执行爬虫
|
||||
task_result = execute_monitor_task()
|
||||
|
||||
if not task_result:
|
||||
@ -111,96 +147,18 @@ def auto_monitor_job(app):
|
||||
scraped_list = task_result.get('device_list', [])
|
||||
print(f"📦 [数据获取] 爬取到 {len(scraped_list)} 条设备数据")
|
||||
|
||||
current_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
stats = {'updated': 0, 'history': 0}
|
||||
updated, history = ingest_device_data(scraped_list)
|
||||
|
||||
for item in scraped_list:
|
||||
d_name = item.get('name')
|
||||
if not d_name: continue
|
||||
|
||||
# --- 1. 数据解包与默认值处理 ---
|
||||
# 显式提取,防止 None 覆盖数据库现有的值(如果业务需要)
|
||||
# 这里假设爬虫返回 None 就是要写入 None,或者空字符串
|
||||
raw_status = item.get('status', '未知')
|
||||
raw_value = item.get('value', '')
|
||||
f_count = item.get('num_files', 0)
|
||||
|
||||
# 时间处理:必须有时间,否则用当前时间
|
||||
target_date = item.get('target_time')
|
||||
if not target_date:
|
||||
target_date = current_time
|
||||
|
||||
raw_json = item.get('raw_json', {})
|
||||
|
||||
# [调试日志] 仅打印第一条或特定的设备,防止刷屏,但能帮你确认数据是否为空
|
||||
# if '0025' in d_name:
|
||||
# print(f" >>> [写入前检查] {d_name}: Value='{raw_value}' | Files={f_count}")
|
||||
|
||||
# --- 2. 数据库操作 (使用 Merge 机制) ---
|
||||
# 先尝试查询
|
||||
device = Device.query.filter_by(name=d_name).first()
|
||||
|
||||
if not device:
|
||||
# 如果不存在,新建对象
|
||||
device = Device(name=d_name, source=item.get('source', '自动爬虫'), install_site="")
|
||||
db.session.add(device)
|
||||
db.session.flush() # 立即获取 ID
|
||||
|
||||
# 更新字段
|
||||
device.status = raw_status
|
||||
device.current_value = raw_value
|
||||
device.latest_time = target_date
|
||||
device.check_time = current_time
|
||||
device.file_count = f_count
|
||||
|
||||
# 计算 Offset
|
||||
try:
|
||||
device.offset = calculate_offset(target_date)
|
||||
except:
|
||||
device.offset = 0
|
||||
|
||||
# JSON 数据合并
|
||||
old_json = {}
|
||||
try:
|
||||
if device.json_data:
|
||||
old_json = json.loads(device.json_data)
|
||||
except:
|
||||
old_json = {}
|
||||
|
||||
if isinstance(raw_json, dict):
|
||||
old_json.update(raw_json)
|
||||
|
||||
device.json_data = json.dumps(old_json, ensure_ascii=False)
|
||||
|
||||
# [核心修复] 使用 merge 告诉 Session "这个对象归你管,请更新它"
|
||||
# 这能解决后台线程中 "DetachedInstanceError" 或更新丢失的问题
|
||||
db.session.merge(device)
|
||||
stats['updated'] += 1
|
||||
|
||||
# --- 3. 写入历史记录 ---
|
||||
history = DeviceHistory(
|
||||
device_id=device.id,
|
||||
status=raw_status,
|
||||
result_data=raw_value,
|
||||
data_time=target_date,
|
||||
file_count=f_count,
|
||||
json_data=device.json_data
|
||||
)
|
||||
db.session.add(history)
|
||||
stats['history'] += 1
|
||||
|
||||
# C. 提交事务
|
||||
db.session.commit()
|
||||
print(f"✅ [入库成功] 设备更新: {stats['updated']} | 历史追加: {stats['history']}")
|
||||
print(f"✅ [入库成功] 设备更新: {updated} | 历史追加: {history}")
|
||||
|
||||
except Exception as e:
|
||||
db.session.rollback()
|
||||
print(f"❌ [严重异常] 数据写入失败: {e}")
|
||||
# 打印堆栈以便排查
|
||||
import traceback
|
||||
traceback.print_exc()
|
||||
finally:
|
||||
# D. 再次清理 Session,防止内存泄漏或污染下一次任务
|
||||
# 再次清理 Session,防止内存泄漏或污染下一次任务
|
||||
db.session.remove()
|
||||
print(f"{'=' * 50}\n")
|
||||
|
||||
|
||||
44
2_1banben/device_monitor.spec
Normal file
44
2_1banben/device_monitor.spec
Normal file
@ -0,0 +1,44 @@
|
||||
# -*- mode: python ; coding: utf-8 -*-
|
||||
|
||||
|
||||
a = Analysis(
|
||||
['app.py'],
|
||||
pathex=[],
|
||||
binaries=[],
|
||||
datas=[('web_dist', 'web_dist')],
|
||||
hiddenimports=[],
|
||||
hookspath=[],
|
||||
hooksconfig={},
|
||||
runtime_hooks=[],
|
||||
excludes=[],
|
||||
noarchive=False,
|
||||
optimize=0,
|
||||
)
|
||||
pyz = PYZ(a.pure)
|
||||
|
||||
exe = EXE(
|
||||
pyz,
|
||||
a.scripts,
|
||||
[],
|
||||
exclude_binaries=True,
|
||||
name='device_monitor',
|
||||
debug=False,
|
||||
bootloader_ignore_signals=False,
|
||||
strip=False,
|
||||
upx=True,
|
||||
console=True,
|
||||
disable_windowed_traceback=False,
|
||||
argv_emulation=False,
|
||||
target_arch=None,
|
||||
codesign_identity=None,
|
||||
entitlements_file=None,
|
||||
)
|
||||
coll = COLLECT(
|
||||
exe,
|
||||
a.binaries,
|
||||
a.datas,
|
||||
strip=False,
|
||||
upx=True,
|
||||
upx_exclude=[],
|
||||
name='device_monitor',
|
||||
)
|
||||
21
2_1banben/fix_fake_time.py
Normal file
21
2_1banben/fix_fake_time.py
Normal file
@ -0,0 +1,21 @@
|
||||
from app import create_app
|
||||
from extensions import db
|
||||
from models import Device
|
||||
|
||||
app = create_app()
|
||||
|
||||
with app.app_context():
|
||||
# 查找状态异常但 offset 被判定为当天的假数据
|
||||
fake_devices = Device.query.filter(
|
||||
Device.status.in_(['离线', '异常', '已离线', 'offline', '未知']),
|
||||
Device.offset == '当天'
|
||||
).all()
|
||||
|
||||
count = 0
|
||||
for dev in fake_devices:
|
||||
dev.latest_time = None
|
||||
dev.offset = '从未同步'
|
||||
count += 1
|
||||
|
||||
db.session.commit()
|
||||
print(f"✅ 成功清洗了 {count} 台处于离线但伪造了时间的设备!")
|
||||
@ -1,6 +1,8 @@
|
||||
# models.py
|
||||
from datetime import datetime
|
||||
from extensions import db
|
||||
# services.time_utils 只依赖 re/datetime,不会与 extensions/models 形成循环导入
|
||||
from services.time_utils import staleness_of
|
||||
|
||||
|
||||
class Device(db.Model):
|
||||
@ -31,8 +33,20 @@ class Device(db.Model):
|
||||
is_whitelist = db.Column(db.Boolean, default=False)
|
||||
|
||||
def to_dict(self):
|
||||
# 统一状态映射逻辑
|
||||
api_status = 'offline' if self.status in ['离线', '异常', '已离线'] else 'online'
|
||||
# 状态归一:中英文各种写法都要覆盖。
|
||||
# 注意 'offline' 是 /add_device 创建手动设备时写入的值,原先的白名单
|
||||
# 漏了它,导致手动添加的设备被映射成 online、在面板上永远显示“在线”。
|
||||
raw_status = str(self.status or '').strip().lower()
|
||||
is_offline = raw_status in (
|
||||
'离线', '异常', '已离线', 'offline', 'unknown', '未知'
|
||||
)
|
||||
api_status = 'offline' if is_offline else 'online'
|
||||
|
||||
# 数据时效在读取时实时计算,不读 offset 列。
|
||||
# 该列是写入时刻的快照,采集一停摆就整体腐烂(库里停在 2026-02-06,
|
||||
# offset 却仍显示“当天”)。列保留但已废弃,见 services/db_ingest.py。
|
||||
stale = staleness_of(self.latest_time)
|
||||
|
||||
return {
|
||||
'id': self.id,
|
||||
'name': self.name,
|
||||
@ -46,8 +60,11 @@ class Device(db.Model):
|
||||
'is_maintaining': self.is_maintaining,
|
||||
'is_hidden': self.is_hidden,
|
||||
'is_whitelist': self.is_whitelist,
|
||||
'offset': self.offset,
|
||||
'file_count': self.file_count # ✅ 返回给前端
|
||||
# 时效三件套:前端不再自己算滞后天数,统一消费这三个字段
|
||||
'offset': stale['text'], # 中文描述,兼容旧调用方
|
||||
'stale_level': stale['level'], # unknown/ok/yesterday/lagging/severe
|
||||
'stale_days': stale['days'], # 滞后自然日数,unknown 时为 None
|
||||
'file_count': self.file_count # ✅ 返回给前端
|
||||
}
|
||||
|
||||
|
||||
|
||||
@ -15,6 +15,11 @@ try:
|
||||
except ImportError:
|
||||
execute_monitor_task = None
|
||||
|
||||
try:
|
||||
from services.db_ingest import ingest_device_data
|
||||
except ImportError:
|
||||
ingest_device_data = None
|
||||
|
||||
try:
|
||||
from services.iot_api import sync_iot_data_service
|
||||
except ImportError:
|
||||
@ -27,21 +32,11 @@ api_bp = Blueprint('api', __name__, url_prefix='/api')
|
||||
# 0. 核心算法区:数据质量分析与辅助函数
|
||||
# =========================================================
|
||||
|
||||
def calculate_offset(latest_time_str):
|
||||
"""
|
||||
计算时间滞后天数
|
||||
用于前端展示设备数据是否过时
|
||||
"""
|
||||
if not latest_time_str or latest_time_str == "N/A":
|
||||
return "从未同步"
|
||||
try:
|
||||
# 兼容处理 2026_01_13 和 2026-01-13 格式
|
||||
clean = str(latest_time_str).split()[0].replace('_', '-')
|
||||
target = datetime.strptime(clean, "%Y-%m-%d").date()
|
||||
diff = (datetime.now().date() - target).days
|
||||
return "当天" if diff == 0 else f"滞后 {diff} 天"
|
||||
except:
|
||||
return "时间解析失败"
|
||||
# calculate_offset 已下沉到 services/time_utils.py。
|
||||
# 原因:services.db_ingest 需要计算 offset,而本模块又要 import db_ingest,
|
||||
# 留在本模块会形成循环导入。这里 re-export 一次,保证
|
||||
# `from routes.api import calculate_offset`(app.py 在用)继续可用。
|
||||
from services.time_utils import calculate_offset # noqa: E402
|
||||
|
||||
|
||||
def check_data_quality(content_data, source_type, data_time_str=None):
|
||||
@ -383,77 +378,14 @@ def run_monitor():
|
||||
|
||||
try:
|
||||
# --- A. 执行爬虫并入库 ---
|
||||
if execute_monitor_task:
|
||||
# 入库细节统一走 services.db_ingest,与定时任务 auto_monitor_job 保持一致,
|
||||
# 避免同一条数据经手动/自动两个入口落库后字段不一致。
|
||||
if execute_monitor_task and ingest_device_data:
|
||||
task_result = execute_monitor_task()
|
||||
if task_result:
|
||||
scraped_list = task_result.get('device_list', [])
|
||||
current_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
count_crawler = 0
|
||||
for item in scraped_list:
|
||||
d_name = item.get('name')
|
||||
if not d_name: continue
|
||||
|
||||
d_raw = item.get('raw_json', {})
|
||||
source = item.get('source', '')
|
||||
target_time = item.get('target_time')
|
||||
|
||||
if '106' in str(source):
|
||||
try:
|
||||
path_str = d_raw.get('path', '')
|
||||
match = re.search(r'/Data/(\d{4}_\d{2}_\d{2})/\w+_(\d{2}_\d{2}_\d{2})\.csv', path_str)
|
||||
if match:
|
||||
date_part = match.group(1).replace('_', '-')
|
||||
time_part = match.group(2).replace('_', ':')
|
||||
target_time = f"{date_part} {time_part}"
|
||||
except:
|
||||
pass
|
||||
|
||||
device = Device.query.filter_by(name=d_name).first()
|
||||
if not device:
|
||||
device = Device(name=d_name, source=source, install_site="")
|
||||
db.session.add(device)
|
||||
db.session.flush()
|
||||
|
||||
if device.source == 'iot_card':
|
||||
device.source = source
|
||||
|
||||
device.status = item.get('status')
|
||||
device.current_value = item.get('value')
|
||||
device.latest_time = target_time
|
||||
device.check_time = current_time
|
||||
|
||||
# ✅ [核心修改] 获取爬虫返回的文件数量并保存
|
||||
f_count = item.get('num_files', 0)
|
||||
device.file_count = f_count
|
||||
|
||||
old_json = {}
|
||||
try:
|
||||
if device.json_data:
|
||||
old_json = json.loads(device.json_data)
|
||||
except:
|
||||
old_json = {}
|
||||
|
||||
new_json = d_raw if isinstance(d_raw, dict) else item.get('raw_json', {})
|
||||
if isinstance(new_json, dict):
|
||||
old_json.update(new_json)
|
||||
|
||||
device.json_data = json.dumps(old_json, ensure_ascii=False)
|
||||
device.offset = calculate_offset(device.latest_time)
|
||||
|
||||
# ✅ [核心修改] 写入历史记录时包含 file_count
|
||||
new_history = DeviceHistory(
|
||||
device_id=device.id,
|
||||
status=item.get('status'),
|
||||
result_data=item.get('value'),
|
||||
data_time=target_time,
|
||||
json_data=device.json_data,
|
||||
file_count=f_count # 确保历史数据也记录文件数
|
||||
)
|
||||
db.session.add(new_history)
|
||||
count_crawler += 1
|
||||
|
||||
msg_list.append(f"爬虫更新: {count_crawler}")
|
||||
updated, _ = ingest_device_data(scraped_list)
|
||||
msg_list.append(f"爬虫更新: {updated}")
|
||||
else:
|
||||
msg_list.append("爬虫无数据")
|
||||
|
||||
|
||||
@ -1,7 +1,9 @@
|
||||
# services/crawler_106.py
|
||||
import os
|
||||
import re
|
||||
import requests
|
||||
import logging
|
||||
import pytz
|
||||
from datetime import datetime
|
||||
from config import Config
|
||||
|
||||
@ -29,11 +31,18 @@ def get_106_dynamic_token(port):
|
||||
|
||||
def find_closest_item(items, is_date_level=True):
|
||||
"""
|
||||
在列表中找到与当前日期最接近的文件夹或文件
|
||||
在列表中找到【不晚于当前时间】的最新的文件夹或文件。
|
||||
|
||||
旧实现用 abs(now - item) 取"离现在最近",服务端只要存在未来日期的
|
||||
目录就会被选中;这里改成只保留已经发生过的,再取其中最新的一个。
|
||||
|
||||
返回 (diff, item, target_str);没有合法项时返回 None。
|
||||
"""
|
||||
if not items or not isinstance(items, list): return None
|
||||
today = datetime.now()
|
||||
scored_items = []
|
||||
|
||||
tz = pytz.timezone('Asia/Shanghai')
|
||||
now = datetime.now(tz)
|
||||
valid_items = []
|
||||
|
||||
for item in items:
|
||||
name_val = item.get('name', '')
|
||||
@ -43,23 +52,30 @@ def find_closest_item(items, is_date_level=True):
|
||||
|
||||
try:
|
||||
if is_date_level:
|
||||
# 解析文件夹日期格式: YYYY_MM_DD
|
||||
current_date = datetime.strptime(target_str, "%Y_%m_%d")
|
||||
# 目录名只有日期没有时刻,按当天 00:00 参与比较。
|
||||
# 注意不能用 23:59:59:那样"今天"的目录会大于当前时刻而被
|
||||
# 当成未来剔除,导致永远取到昨天的数据。
|
||||
item_dt = tz.localize(
|
||||
datetime.strptime(f"{target_str} 00:00:00", "%Y_%m_%d %H:%M:%S")
|
||||
)
|
||||
else:
|
||||
# 解析文件修改时间
|
||||
# 文件修改时间:带 'Z' 表示 UTC;不带时区的按东八区处理,
|
||||
# 否则会与 now 出现 naive/aware 比较异常而被静默丢弃。
|
||||
mod_str = item.get('modified', '')
|
||||
current_date = datetime.fromisoformat(mod_str.replace('Z', '+00:00'))
|
||||
item_dt = datetime.fromisoformat(mod_str.replace('Z', '+00:00'))
|
||||
if item_dt.tzinfo is None:
|
||||
item_dt = tz.localize(item_dt)
|
||||
|
||||
# 计算与当前时间的差距
|
||||
diff = abs((today - current_date.replace(tzinfo=None)).total_seconds())
|
||||
scored_items.append((diff, item, target_str))
|
||||
except:
|
||||
# 【关键修复】剔除未来时间,只保留已经发生过的
|
||||
if item_dt <= now:
|
||||
valid_items.append(((now - item_dt).total_seconds(), item, target_str))
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
if not scored_items: return None
|
||||
# 按时间差排序,取最小的
|
||||
scored_items.sort(key=lambda x: x[0])
|
||||
return scored_items[0]
|
||||
if not valid_items: return None
|
||||
# 按时间差升序排序(取最近的过去)
|
||||
valid_items.sort(key=lambda x: x[0])
|
||||
return valid_items[0]
|
||||
|
||||
|
||||
def run_106_logic():
|
||||
@ -89,15 +105,16 @@ def run_106_logic():
|
||||
if not (is_tower_underscore or is_tower_i): continue
|
||||
|
||||
# --- 构建基础数据包 ---
|
||||
# 默认使用标准当前时间作为兜底,防止后续步骤失败时时间为空
|
||||
current_standard_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
# target_time 默认 None:只有真正从文件路径解析出记录时间才会被覆盖。
|
||||
# 离线/异常/token 失败等分支一律留空,交给入库层"冻结"上一个有效时间,
|
||||
# 绝不用 datetime.now() 冒充数据时间。
|
||||
# status 默认 '正常':成功路径不会回头重设它,改这里等于改所有健康设备的显示。
|
||||
data_packet = {
|
||||
'source': '106网站',
|
||||
'name': name,
|
||||
'status': '正常',
|
||||
'value': '',
|
||||
'target_time': current_standard_time,
|
||||
'target_time': None,
|
||||
'raw_json': {},
|
||||
'temp_file': None,
|
||||
'num_files': 0
|
||||
@ -133,16 +150,13 @@ def run_106_logic():
|
||||
continue
|
||||
|
||||
# ==============================================================================
|
||||
# ✅ [核心修复] 时间格式标准化
|
||||
# 原逻辑: data_packet['target_time'] = best_date[2] (得到 "2026_02_08")
|
||||
# 新逻辑: 将 "2026_02_08" 转换为 "2026-02-08 HH:MM:SS"
|
||||
# ⚠️ 这里刻意不再给 target_time 赋值。
|
||||
# 旧逻辑是 formatted_date_part + datetime.now() 的时分秒,写进库的
|
||||
# 实际是"爬虫运行时刻",导致 calculate_offset 恒为"当天"。
|
||||
# 真实时间要等确定最终文件后,从文件路径里解析(见下方 [核心修复])。
|
||||
# ==============================================================================
|
||||
raw_folder_name = best_date[2] # 例如 "2026_02_08"
|
||||
formatted_date_part = raw_folder_name.replace('_', '-') # 变成 "2026-02-08"
|
||||
current_time_part = datetime.now().strftime("%H:%M:%S")
|
||||
|
||||
# 覆盖默认时间,确保数据库存入的是标准时间戳格式
|
||||
data_packet['target_time'] = f"{formatted_date_part} {current_time_part}"
|
||||
|
||||
date_path = f"{api_root}{raw_folder_name}/"
|
||||
|
||||
@ -165,6 +179,41 @@ def run_106_logic():
|
||||
file_item = best_file[1]
|
||||
full_path = file_item.get('path') or f"{date_path}{file_item.get('name')}"
|
||||
|
||||
# ==============================================================================
|
||||
# ✅ [核心修复] 从最终确定的文件路径里正则提取真正的记录时间
|
||||
# 原逻辑在 routes/api.py 的 run_monitor 里,只有手动触发才生效;
|
||||
# 下沉到爬虫层后,自动/手动两条路径拿到的是同一个真实时间。
|
||||
# ==============================================================================
|
||||
real_time = None
|
||||
# 放宽点:路径分隔符、data 大小写(TOWER-I 走小写分支)、
|
||||
# 是否以 / 开头、扩展名大小写都做兼容
|
||||
match = re.search(
|
||||
r'(?:^|[/\\])data/(\d{4}_\d{2}_\d{2})/[\w\-.]+_(\d{2}_\d{2}_\d{2})\.csv',
|
||||
full_path,
|
||||
re.IGNORECASE
|
||||
)
|
||||
if match:
|
||||
date_part = match.group(1).replace('_', '-')
|
||||
time_part = match.group(2).replace('_', ':')
|
||||
real_time = f"{date_part} {time_part}"
|
||||
|
||||
# 正则没命中(例如二进制 .db 文件)时退化到文件修改时间,
|
||||
# 绝不拿系统当前时间去伪造记录时间
|
||||
if not real_time:
|
||||
mod_str = file_item.get('modified', '')
|
||||
if mod_str:
|
||||
try:
|
||||
mod_dt = datetime.fromisoformat(mod_str.replace('Z', '+00:00'))
|
||||
if mod_dt.tzinfo is None:
|
||||
mod_dt = pytz.timezone('Asia/Shanghai').localize(mod_dt)
|
||||
else:
|
||||
mod_dt = mod_dt.astimezone(pytz.timezone('Asia/Shanghai'))
|
||||
real_time = mod_dt.strftime("%Y-%m-%d %H:%M:%S")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
data_packet['target_time'] = real_time
|
||||
|
||||
# --- 4. 下载/读取内容逻辑 ---
|
||||
if is_tower_i:
|
||||
# [二进制文件] 下载逻辑
|
||||
|
||||
@ -28,7 +28,9 @@ def run_82_logic():
|
||||
'name': str(sid),
|
||||
'status': '正常',
|
||||
'value': '',
|
||||
'target_time': datetime.now().strftime("%Y-%m-%d %H:%M:%S"),
|
||||
# 同 106:默认留空,只有真正拿到数据时间才覆盖,
|
||||
# 不用 datetime.now() 冒充数据时间
|
||||
'target_time': None,
|
||||
'raw_json': {},
|
||||
'temp_file': None
|
||||
}
|
||||
@ -42,7 +44,9 @@ def run_82_logic():
|
||||
|
||||
if data:
|
||||
d_list = data.get('date', [])
|
||||
latest = str(d_list[-1]) if d_list else "N/A"
|
||||
# 拿不到时间时必须是 None,不能是 "N/A" —— 字符串在入库层是真值,
|
||||
# 会把设备主表已经冻结的有效 latest_time 覆盖掉,令 offset 倒退成"从未同步"
|
||||
latest = str(d_list[-1]) if d_list else None
|
||||
data_packet['target_time'] = latest
|
||||
data_packet['value'] = f"Data Points: {len(d_list)}"
|
||||
data_packet['raw_json'] = data # 🔥 存完整JSON
|
||||
|
||||
119
2_1banben/services/db_ingest.py
Normal file
119
2_1banben/services/db_ingest.py
Normal file
@ -0,0 +1,119 @@
|
||||
# services/db_ingest.py
|
||||
"""
|
||||
统一的设备数据入库管道。
|
||||
|
||||
此前 app.py:auto_monitor_job(定时任务)和 routes/api.py:run_monitor(手动触发)
|
||||
各自维护了一份几乎相同但细节不一致的写入逻辑,导致同一条数据经不同入口落库后
|
||||
latest_time / offset / source / file_count 可能不同。这里合并为一份。
|
||||
|
||||
本函数只负责 add/flush,不 commit —— 事务边界由调用方掌握。
|
||||
"""
|
||||
import json
|
||||
from datetime import datetime
|
||||
|
||||
from extensions import db
|
||||
from models import Device, DeviceHistory
|
||||
|
||||
# 由本应用接口写入 Device.json_data、而非爬虫返回的键,覆盖时必须保留:
|
||||
# bound_iccid —— routes/api.py:/bind_device_card 写入的设备-流量卡手工绑定
|
||||
# is_whitelist —— routes/api.py:/toggle_whitelist 写入的白名单标记
|
||||
APP_OWNED_KEYS = ('bound_iccid', 'is_whitelist')
|
||||
|
||||
|
||||
def ingest_device_data(scraped_list):
|
||||
"""
|
||||
把爬虫返回的设备列表写入 Device(快照)和 DeviceHistory(历史)。
|
||||
|
||||
返回 (设备更新数, 历史追加数)。
|
||||
"""
|
||||
if not scraped_list:
|
||||
return 0, 0
|
||||
|
||||
current_time = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||||
stats = {'updated': 0, 'history': 0}
|
||||
|
||||
for item in scraped_list:
|
||||
d_name = item.get('name')
|
||||
if not d_name:
|
||||
continue
|
||||
|
||||
# --- 1. 数据解包 ---
|
||||
raw_status = item.get('status', '未知')
|
||||
raw_value = item.get('value', '')
|
||||
f_count = item.get('num_files', 0)
|
||||
source = item.get('source', '自动爬虫')
|
||||
|
||||
# target_time 由爬虫层解析,为真实记录时间;离线/异常/没抓到时为 None。
|
||||
# 注意:这里绝不回退到 current_time —— 那会把"采集时刻"冒充成"数据时刻"。
|
||||
target_date = item.get('target_time')
|
||||
|
||||
raw_json = item.get('raw_json', {})
|
||||
|
||||
# --- 2. 设备主表更新 ---
|
||||
device = Device.query.filter_by(name=d_name).first()
|
||||
|
||||
if not device:
|
||||
device = Device(name=d_name, source=source, install_site="")
|
||||
db.session.add(device)
|
||||
db.session.flush()
|
||||
elif device.source == 'iot_card':
|
||||
# IoT 卡与爬虫设备共用 devices 表且靠 name 关联,撞名时以爬虫来源为准
|
||||
device.source = source
|
||||
|
||||
device.status = raw_status
|
||||
device.current_value = raw_value
|
||||
device.check_time = current_time
|
||||
device.file_count = f_count
|
||||
|
||||
# ✅ [核心逻辑] 只有爬虫拿到了真实业务时间才更新设备主表的时间。
|
||||
# 拿不到(离线/异常)就保留上一次的有效值 —— 即"冻结"。
|
||||
# 这样一台断线多天的设备,offset 会停在"滞后 N 天"而不是被刷成"当天"。
|
||||
if target_date:
|
||||
device.latest_time = target_date
|
||||
|
||||
# 注意:不再写入 device.offset。
|
||||
# 该列是"写入时刻"的快照,采集一停摆就整体失真,现改为 Device.to_dict()
|
||||
# 读取时用 calculate_offset(latest_time) 实时计算。列保留但已废弃。
|
||||
|
||||
# --- 3. JSON 数据:覆盖为本次最新,切断无限膨胀 ---
|
||||
# 旧写法 old_json.update(raw_json) 会让主表 JSON 只增不减、永久累积。
|
||||
# 但 APP_OWNED_KEYS 里的键由本应用自己的接口写入(不属于爬虫数据),
|
||||
# 必须原样保留 —— 否则每次采集都会把用户的设备-流量卡绑定清掉。
|
||||
old_json = {}
|
||||
if device.json_data:
|
||||
try:
|
||||
old_json = json.loads(device.json_data)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if isinstance(raw_json, dict) and raw_json:
|
||||
new_json = dict(raw_json)
|
||||
for key in APP_OWNED_KEYS:
|
||||
if old_json.get(key) is not None:
|
||||
new_json[key] = old_json[key]
|
||||
device.json_data = json.dumps(new_json, ensure_ascii=False)
|
||||
# raw_json 为空(离线/异常)时保持原有 json_data 不变,避免单次爬虫失误清空状态
|
||||
|
||||
db.session.merge(device)
|
||||
stats['updated'] += 1
|
||||
|
||||
# --- 4. 历史表写入 ---
|
||||
# 历史记录的是"事件",所以没有数据时间时用本次观测时间,
|
||||
# 与主表的"冻结"策略不同:主表回答"数据到什么时候",历史回答"什么时候采过"。
|
||||
history_time = target_date if target_date else current_time
|
||||
# [核心修复] 只存本次抓取的切片,不再复制主表那份越滚越大的累积 JSON。
|
||||
# 旧写法 json_data=device.json_data 会让每一条历史都完整复制一遍全量数据,
|
||||
# 历史表体积随采集次数线性膨胀。
|
||||
incremental_json = json.dumps(raw_json, ensure_ascii=False) if raw_json else "{}"
|
||||
history = DeviceHistory(
|
||||
device_id=device.id,
|
||||
status=raw_status,
|
||||
result_data=raw_value,
|
||||
data_time=history_time,
|
||||
file_count=f_count,
|
||||
json_data=incremental_json
|
||||
)
|
||||
db.session.add(history)
|
||||
stats['history'] += 1
|
||||
|
||||
return stats['updated'], stats['history']
|
||||
97
2_1banben/services/time_utils.py
Normal file
97
2_1banben/services/time_utils.py
Normal file
@ -0,0 +1,97 @@
|
||||
# services/time_utils.py
|
||||
"""
|
||||
数据时效判定的唯一实现。
|
||||
|
||||
此前「滞后几天」这件事在后端(calculate_offset)和前端(Dashboard.vue 里
|
||||
自算 diffDays/diffHours)各算了一套,阈值和粒度都不同,跨日必然打架:
|
||||
昨天 23:00 的数据在次日 01:00 查看,后端说「滞后 1 天」,前端标「昨日数据」。
|
||||
现在统一由这里判定,前端只消费 Device.to_dict() 暴露的
|
||||
offset / stale_level / stale_days 三个字段做渲染。
|
||||
|
||||
刻意不依赖 flask / db / models,避免被 services.db_ingest 和 routes.api
|
||||
互相引用时产生循环导入。
|
||||
"""
|
||||
import re
|
||||
from datetime import datetime, date
|
||||
|
||||
# 视为「从未同步」的占位值
|
||||
_NEVER_SYNCED = ('', 'N/A', 'NA', 'NONE', 'NULL', 'NAN')
|
||||
|
||||
# 从任意时间字符串里抠出 YYYY-MM-DD。分隔符兼容 - / _ ,
|
||||
# 因此 ISO 的 2026-09-15T08:30:00Z 和 2026_09_15 都能命中。
|
||||
# 滞后判定按自然日,时分秒不影响结果,所以不做完整的时间解析。
|
||||
_DATE_RE = re.compile(r'(\d{4})[-/_](\d{1,2})[-/_](\d{1,2})')
|
||||
|
||||
# 超过这个自然日数算「严重滞后」
|
||||
SEVERE_THRESHOLD_DAYS = 7
|
||||
|
||||
# 判定级别
|
||||
LEVEL_UNKNOWN = 'unknown' # 从未同步 / 时间无法解析
|
||||
LEVEL_OK = 'ok' # 当天
|
||||
LEVEL_YESTERDAY = 'yesterday' # 昨天
|
||||
LEVEL_LAGGING = 'lagging' # 滞后 2~7 天
|
||||
LEVEL_SEVERE = 'severe' # 滞后 7 天以上
|
||||
|
||||
|
||||
def parse_record_date(value):
|
||||
"""
|
||||
把各种格式的记录时间解析成 date,失败返回 None。
|
||||
|
||||
兼容:2026-09-15 08:30:00 / 2026_09_15 / 2026-09-15T08:30:00Z /
|
||||
2026/09/15 08:30 / 2026-09-15 08:30 / 2026-09-15
|
||||
"""
|
||||
if value is None:
|
||||
return None
|
||||
|
||||
text = str(value).strip()
|
||||
if not text or text.upper() in _NEVER_SYNCED:
|
||||
return None
|
||||
|
||||
match = _DATE_RE.search(text)
|
||||
if not match:
|
||||
return None
|
||||
|
||||
try:
|
||||
return date(int(match.group(1)), int(match.group(2)), int(match.group(3)))
|
||||
except ValueError:
|
||||
# 例如 2026-13-45 这种越界日期
|
||||
return None
|
||||
|
||||
|
||||
def staleness_of(value):
|
||||
"""
|
||||
统一的滞后判定。
|
||||
|
||||
返回 {'level': str, 'days': int|None, 'text': str}
|
||||
level —— unknown / ok / yesterday / lagging / severe
|
||||
days —— 滞后自然日数,unknown 时为 None
|
||||
text —— 中文描述,可直接展示
|
||||
"""
|
||||
if value is None or str(value).strip().upper() in _NEVER_SYNCED:
|
||||
return {'level': LEVEL_UNKNOWN, 'days': None, 'text': '从未同步'}
|
||||
|
||||
record_date = parse_record_date(value)
|
||||
if record_date is None:
|
||||
return {'level': LEVEL_UNKNOWN, 'days': None, 'text': '时间解析失败'}
|
||||
|
||||
days = (datetime.now().date() - record_date).days
|
||||
if days < 0:
|
||||
# 数据时间在未来(设备/服务端时钟超前),按当天处理,不显示负数
|
||||
days = 0
|
||||
|
||||
if days == 0:
|
||||
return {'level': LEVEL_OK, 'days': 0, 'text': '当天'}
|
||||
if days == 1:
|
||||
return {'level': LEVEL_YESTERDAY, 'days': 1, 'text': '滞后 1 天'}
|
||||
if days <= SEVERE_THRESHOLD_DAYS:
|
||||
return {'level': LEVEL_LAGGING, 'days': days, 'text': f'滞后 {days} 天'}
|
||||
return {'level': LEVEL_SEVERE, 'days': days, 'text': f'滞后 {days} 天'}
|
||||
|
||||
|
||||
def calculate_offset(latest_time_str):
|
||||
"""
|
||||
(兼容保留) 返回滞后天数的中文描述文本。
|
||||
|
||||
新代码请直接用 staleness_of(),需要分级时不要再从这段文本里反解。
|
||||
"""
|
||||
return staleness_of(latest_time_str)['text']
|
||||
@ -1,8 +1,10 @@
|
||||
import { createRouter, createWebHistory } from 'vue-router'
|
||||
import { ElMessage } from 'element-plus'
|
||||
|
||||
// 1. 引入登录页面(建议新建 views/Login.vue)
|
||||
import Login from '../views/Login.vue'
|
||||
// 1. 引入登录页面
|
||||
// 注意:实际文件名是全小写的 login.vue。Windows 文件系统不区分大小写,所以本地
|
||||
// 构建能过;但在 Linux/Docker/CI 上会直接报 Could not resolve,必须严格匹配。
|
||||
import Login from '../views/login.vue'
|
||||
// 2. 首页组件
|
||||
import Dashboard from '../views/Dashboard.vue'
|
||||
|
||||
|
||||
@ -341,8 +341,13 @@ const fetchData = async () => {
|
||||
const isOrphanIoT = (item.source === 'iot_card')
|
||||
const isWhitelist = !!item.is_whitelist
|
||||
|
||||
// === 1. 智能时间解析与格式化 (增强版) ===
|
||||
let diffDays = 0, diffHours = 0, isToday = false, validTime = false
|
||||
// === 1. 时间显示格式化 ===
|
||||
// 滞后判定已统一由后端提供(stale_level / stale_days / offset),
|
||||
// 这里只负责把 latest_time 渲染成统一格式,不再本地计算天数。
|
||||
// 判定与渲染分家的原因见 services/time_utils.py 顶部说明。
|
||||
const staleLevel = item.stale_level || 'unknown'
|
||||
const staleDays = (typeof item.stale_days === 'number') ? item.stale_days : null
|
||||
const validTime = staleLevel !== 'unknown'
|
||||
let timeStr = item.latest_time
|
||||
|
||||
// 默认显示原始值,稍后如果解析成功则覆盖它
|
||||
@ -377,16 +382,8 @@ const fetchData = async () => {
|
||||
d = new Date(cleanStr);
|
||||
}
|
||||
|
||||
// C. 如果解析成功,强制重新生成统一的显示字符串
|
||||
// C. 解析成功则强制重新生成统一的显示格式 YYYY-MM-DD HH:mm:ss
|
||||
if (d && !isNaN(d.getTime())) {
|
||||
validTime = true
|
||||
isToday = d.toDateString() === now.toDateString()
|
||||
|
||||
const diff = now - d
|
||||
diffHours = (diff > 0 ? diff : 0) / (1000 * 3600)
|
||||
diffDays = diffHours / 24
|
||||
|
||||
// 🌟 核心修改点:生成标准显示格式 YYYY-MM-DD HH:mm:ss 🌟
|
||||
const y = d.getFullYear()
|
||||
const m = String(d.getMonth() + 1).padStart(2, '0')
|
||||
const dd = String(d.getDate()).padStart(2, '0')
|
||||
@ -429,29 +426,33 @@ const fetchData = async () => {
|
||||
}
|
||||
}
|
||||
|
||||
// 4. 状态判定
|
||||
// 4. 状态判定(分级由后端给出,这里只做映射,不再本地判天数)
|
||||
let statusColor = '#67C23A', statusLabel = '正常', statusType = 'normal', statusLabelColor = '#fff'
|
||||
let statusReason = ''
|
||||
let sortWeight = diffHours
|
||||
// 排序权重:问题设备排前面。staleDays * 24 与原 diffHours 同为小时量纲,
|
||||
// 排序结果与原实现一致(同为“当天”的设备不再按小时细分,顺序退化为稳定序)
|
||||
let sortWeight = (staleDays === null) ? 80000000 : staleDays * 24
|
||||
|
||||
if (item.is_maintaining) {
|
||||
statusColor = '#409EFF'; statusLabel = '维修中'; statusType = 'maintenance';
|
||||
sortWeight = Number.MAX_SAFE_INTEGER;
|
||||
} else if (!validTime || item.status === 'offline') {
|
||||
statusLabel = '离线'; statusColor = '#F56C6C'; statusType = 'error';
|
||||
statusReason = validTime ? '设备离线' : '暂无数据';
|
||||
// 文案与后端 offset 统一:拿不到可解析的数据时间就说“从未同步”,
|
||||
// 不再另造一套“暂无数据”的说法
|
||||
statusReason = (item.status === 'offline') ? '设备离线' : (item.offset || '从未同步');
|
||||
sortWeight = 80000000;
|
||||
} else if (diffDays > 7) {
|
||||
} else if (staleLevel === 'severe') {
|
||||
statusLabel = '严重滞后'; statusColor = '#F56C6C'; statusType = 'error';
|
||||
statusReason = `滞后 ${Math.floor(diffDays)} 天`;
|
||||
} else if (diffHours > 24) {
|
||||
statusReason = item.offset || `滞后 ${staleDays} 天`;
|
||||
} else if (staleLevel === 'lagging') {
|
||||
statusLabel = '滞后'; statusColor = '#E6A23C'; statusType = 'warning';
|
||||
statusReason = `滞后 ${Math.floor(diffDays)} 天`;
|
||||
statusReason = item.offset || `滞后 ${staleDays} 天`;
|
||||
} else if (expireWarning) {
|
||||
statusLabel = '即将过期'; statusColor = '#E6A23C'; statusType = 'warning';
|
||||
statusReason = `即将过期`;
|
||||
sortWeight = 400;
|
||||
} else if (!isToday) {
|
||||
} else if (staleLevel === 'yesterday') {
|
||||
statusLabel = '昨日数据'; statusColor = '#FAC858'; statusType = 'slight-warning'; statusLabelColor = '#333';
|
||||
statusReason = '非今日数据';
|
||||
} else {
|
||||
@ -465,7 +466,7 @@ const fetchData = async () => {
|
||||
isOrphanIoT,
|
||||
isBound,
|
||||
isWhitelist,
|
||||
diffDays, diffHours, sortWeight, isToday,
|
||||
sortWeight, staleLevel, staleDays,
|
||||
statusColor, statusLabel, statusType, statusLabelColor, statusReason,
|
||||
isEditingSite: false, tempSite: '',
|
||||
data_quality: item.data_quality || 'ok',
|
||||
@ -534,7 +535,22 @@ const handleAddDeviceSubmit = async () => { if (!newDeviceForm.name) return; awa
|
||||
const handleDeviceClick = (row) => { if (!row.is_hidden && dataMonitorRef.value) dataMonitorRef.value.open(row) }
|
||||
const openLogCenter = (row) => { if (maintenanceLogsRef.value) maintenanceLogsRef.value.open(row ? { deviceName: row.name } : null) }
|
||||
const openIoTBinder = () => { if (iotBinderRef.value) iotBinderRef.value.open() }
|
||||
const runManualMonitor = async () => { runningTask.value=true; await axios.post(`${API_BASE}/api/run_monitor`); setTimeout(()=>fetchData(), 3000); setTimeout(()=>runningTask.value=false, 1000) }
|
||||
const runManualMonitor = async () => {
|
||||
if (runningTask.value) return; // 防止重复点击
|
||||
runningTask.value = true;
|
||||
try {
|
||||
// 后端 execute_monitor_task() 是同步阻塞的,105 台设备逐个登录+下载要跑几分钟,
|
||||
// 必须给足超时,否则请求会先断开、爬虫却还在后台跑
|
||||
await axios.post(`${API_BASE}/api/run_monitor`, {}, { timeout: 600000 });
|
||||
ElMessage.success('手动采集完成!');
|
||||
await fetchData(); // 等后端真正跑完再取数,不再用固定 setTimeout 瞎猜
|
||||
} catch (error) {
|
||||
console.error("手动采集失败:", error);
|
||||
ElMessage.error('采集失败或超时:' + (error.response?.data?.message || error.message));
|
||||
} finally {
|
||||
runningTask.value = false; // 无论成功失败都关掉 loading
|
||||
}
|
||||
}
|
||||
const handleEditSite = (row) => { row.tempSite = row.install_site; row.isEditingSite = true; nextTick(() => document.querySelector('.site-input-inner input')?.focus()) }
|
||||
const saveSite = async (row) => { if(!row.isEditingSite)return; row.isEditingSite=false; await axios.post(`${API_BASE}/api/update_site`, {name:row.name, site:row.tempSite}); row.install_site=row.tempSite }
|
||||
const handleMaintenanceBeforeChange = (row) => { return new Promise(r => { axios.post(`${API_BASE}/api/toggle_maintenance`, {name:row.name, is_maintaining:!row.is_maintaining}).then(() => {row.is_maintaining=!row.is_maintaining; fetchData(); r(true)}).catch(()=>r(false)) }) }
|
||||
|
||||
Reference in New Issue
Block a user