Repository navigation
Expand file tree
/
Copy pathmain.py
More file actions
144 lines (109 loc) · 4.78 KB
/
Copy pathmain.py
File metadata and controls
144 lines (109 loc) · 4.78 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
import datetime
import threading
from contextlib import asynccontextmanager
import uvicorn
from apscheduler.schedulers.background import BackgroundScheduler
from fastapi import FastAPI, BackgroundTasks
from Ai_Tags import aiTags_isNull, aiTags_change, extract_tags_in_batch, update_ai_tags_to_db
from RapidOCR import rapid_ocr, filter_data
# 全局互斥锁:防止定时巡检与手动触发同时打标同一批商品
_tagging_lock = threading.Lock()
# 模块级持有调度器,避免被 GC 回收,同时防重复启动
_scheduler = BackgroundScheduler()
def _process_image_map(image_map: dict) -> list:
"""
公共打标流水线:OCR 识别 -> 规则过滤 -> LLM 批量打标 -> 写库。
定时巡检和单商品手动触发共用,保证两处逻辑一致。
"""
imageocr_list = []
for item_pid, urls in image_map.items():
combined_texts = []
for url in urls:
orc = rapid_ocr(url)
data = filter_data(orc)
if data:
combined_texts.append(data)
final_product_text = "\n".join(combined_texts)
if not final_product_text:
# OCR 失败或无有效文本的商品直接跳过,不浪费大模型调用
print(f"⚠️ 商品 {item_pid} 没有识别到有效文本,跳过。")
continue
imageocr_list.append({
"id": item_pid,
"ocr_info": final_product_text
})
if not imageocr_list:
print("⚠️ 没有可用的 OCR 文本,跳过本轮打标。")
return []
ai_tags_list = extract_tags_in_batch(imageocr_list)
update_ai_tags_to_db(ai_tags_list)
return ai_tags_list
def process_tags_task(pid: int):
"""
后台执行的具体打标任务(耗时操作)
"""
with _tagging_lock:
try:
image_map = aiTags_change(pid)
if not image_map:
print(f"✅ 商品 {pid} 当前没有需要打标的图片。")
return
print("👀 Map数据抽样预览:", image_map)
ai_tags_list = _process_image_map(image_map)
print("👀 Ai数据返回结果预览:", ai_tags_list)
print(f"✨ 商品 {pid} 本轮打标任务圆满完成!")
except Exception as e:
print(f"❌ 商品 {pid} 任务执行过程中发生异常: {e}")
def run_tagging_job():
"""
具体的打标任务逻辑,被抽取成一个独立的方法,方便定时器调用
"""
print(f"\n[{datetime.datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] 🚀 开始执行定时 AI 打标任务...")
with _tagging_lock:
try:
# 获取空标签的商品图片数据(每轮限量,防止一次把上下文打爆)
image_map = aiTags_isNull()
if not image_map:
print("✅ 当前没有需要打标的空标签商品,休息一下。")
return
print("👀 Map数据抽样预览:", image_map)
ai_tags_list = _process_image_map(image_map)
print("👀 Ai数据返回结果预览:", ai_tags_list)
print("✨ 本轮打标任务圆满完成!")
except Exception as e:
print(f"❌ 本轮任务执行过程中发生异常: {e}")
def start_scheduler():
"""启动 2 小时一次的全局巡检定时任务"""
if _scheduler.running:
return
_scheduler.add_job(run_tagging_job, 'interval', hours=2)
_scheduler.start()
print("⏰ 定时任务已随 Web 服务成功启动!每 2 小时自动执行一次全局巡检...")
# 启动时先异步跑一次,避免阻塞服务启动(首次全量打标可能耗时数分钟)
threading.Thread(target=run_tagging_job, daemon=True).start()
@asynccontextmanager
async def lifespan(app: FastAPI):
start_scheduler()
yield
app = FastAPI(title="HxMall AI 打标微服务", lifespan=lifespan)
@app.post("/api/ai/update_AiTags")
async def trigger_tag_update(pid: int, background_tasks: BackgroundTasks):
"""
触发单个商品 AI 标签更新的接口
"""
# 1. 参数校验等前置操作(可选)
if pid <= 0:
return {"code": 400, "message": "非法的商品 ID"}
# 2. ⚡️ 将耗时的任务丢入后台线程池!
# 就像你往 MQ (RabbitMQ/Kafka) 里发了一条消息一样
background_tasks.add_task(process_tags_task, pid)
# 3. 毫秒级立刻响应调用方
return {
"code": 200,
"message": f"商品 {pid} 的图片变更已收到,AI 正在后台静默刷新标签..."
}
if __name__ == '__main__':
# 🚀 真正启动 Web 服务器的地方!
print("🚀 正在启动 HxMall AI 标签微服务...")
# 端口 8100:与 AI 客服服务(8000)错开,避免同机部署端口冲突
uvicorn.run(app, host="0.0.0.0", port=8100)