120 lines
5.0 KiB
Python
120 lines
5.0 KiB
Python
|
|
# 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']
|