fix(scheduler): 调度日志补行缓冲,让后台任务的执行结果可见

问题:gunicorn 的 stdout 是管道,Python 默认对它做**块缓冲** ——
print() 的内容会攒在缓冲区里,docker logs 里看不到。

实测 2026-09-23 17:30 的日报任务:8 个 worker 里只刷出 2 条
"本轮跳过",连抢到锁那个 worker 的成败都没出现,完全无从判断
那封日报到底发了没有。事后排查时只能靠猜。

这不是偶发 —— 定时任务跑在后台线程,出问题本来就不容易发现,
再叠加输出缓冲就等于"瞎跑"。补齐后同一条路径 8 条日志全部刷出,
7 个 worker 跳过、1 个执行,一目了然。

改动:
  · run.py 顶部 sys.stdout.reconfigure(line_buffering=True),全局生效;
  · 调度任务的 print 一律再加 flush=True 作为双保险 —— 这几行输出是
    出问题时唯一的线索,不能依赖缓冲策略。
This commit is contained in:
yueli
2026-09-23 17:41:27 +08:00
parent 446c064836
commit 4ed8abcc6d

View File

@ -1,6 +1,16 @@
# inventory-backend/run.py # inventory-backend/run.py
import sys
from app import create_app from app import create_app
# ★ stdout 行缓冲。
# gunicorn 的 stdout 是管道,Python 默认对它做**块缓冲** —— print() 的内容
# 会攒在缓冲区里,docker logs 里看不到。实测 17:30 的日报任务:8 个 worker
# 里只刷出 2 条日志,连抢到锁那个 worker 的成败都没出现,完全无从排查。
# 调度任务是后台线程,出问题本来就不容易发现,再叠加缓冲就等于瞎跑。
if hasattr(sys.stdout, 'reconfigure'):
sys.stdout.reconfigure(line_buffering=True)
app = create_app() app = create_app()
# ========================================================= # =========================================================
@ -24,6 +34,8 @@ import pytz
beijing_tz = pytz.timezone('Asia/Shanghai') beijing_tz = pytz.timezone('Asia/Shanghai')
# 调度日志一律再加一层 flush=True(理由见文件头部的行缓冲说明)——
# 定时任务跑在后台线程,出问题时唯一的线索就是这几行输出,不能丢。
def _run_warning_job(): def _run_warning_job():
"""库存预警扫描与邮件发送(每天 9:30 北京时间)""" """库存预警扫描与邮件发送(每天 9:30 北京时间)"""
with app.app_context(): with app.app_context():
@ -31,14 +43,14 @@ def _run_warning_job():
with advisory_lock(LOCK_INVENTORY_WARNING) as acquired: with advisory_lock(LOCK_INVENTORY_WARNING) as acquired:
if not acquired: if not acquired:
# 正常路径:另一个 worker 正在执行同一任务 # 正常路径:另一个 worker 正在执行同一任务
print("[Scheduler] 库存预警:另一 worker 正在执行,本轮跳过") print("[Scheduler] 库存预警:另一 worker 正在执行,本轮跳过", flush=True)
return return
try: try:
from app.services.inventory_task import InventoryWarningService from app.services.inventory_task import InventoryWarningService
result = InventoryWarningService.check_and_send_warning_emails() result = InventoryWarningService.check_and_send_warning_emails()
print(f"[Scheduler] 库存预警扫描完成: red={result['red_count']}, yellow={result['yellow_count']}") print(f"[Scheduler] 库存预警扫描完成: red={result['red_count']}, yellow={result['yellow_count']}", flush=True)
except Exception as e: except Exception as e:
print(f"[Scheduler] 库存预警任务失败: {e}") print(f"[Scheduler] 库存预警任务失败: {e}", flush=True)
def _run_daily_report_job(): def _run_daily_report_job():
@ -47,14 +59,14 @@ def _run_daily_report_job():
from app.utils.job_lock import advisory_lock, LOCK_DAILY_REPORT from app.utils.job_lock import advisory_lock, LOCK_DAILY_REPORT
with advisory_lock(LOCK_DAILY_REPORT) as acquired: with advisory_lock(LOCK_DAILY_REPORT) as acquired:
if not acquired: if not acquired:
print("[Scheduler] 系统日报:另一 worker 正在执行,本轮跳过") print("[Scheduler] 系统日报:另一 worker 正在执行,本轮跳过", flush=True)
return return
try: try:
from app.services.daily_report_service import DailyReportService from app.services.daily_report_service import DailyReportService
r = DailyReportService.send_daily_report() r = DailyReportService.send_daily_report()
print(f"[Scheduler] 系统日报已发送: {r['subject']} -> {r['recipients']}") print(f"[Scheduler] 系统日报已发送: {r['subject']} -> {r['recipients']}", flush=True)
except Exception as e: except Exception as e:
print(f"[Scheduler] 系统日报任务失败: {e}") print(f"[Scheduler] 系统日报任务失败: {e}", flush=True)
scheduler = BackgroundScheduler(timezone=beijing_tz) scheduler = BackgroundScheduler(timezone=beijing_tz)
@ -73,7 +85,7 @@ scheduler.add_job(
replace_existing=True replace_existing=True
) )
scheduler.start() scheduler.start()
print("✅ 定时任务已启动:库存预警 9:30 / MOM系统日报 17:30(北京时间)") print("✅ 定时任务已启动:库存预警 9:30 / MOM系统日报 17:30(北京时间)", flush=True)
if __name__ == '__main__': if __name__ == '__main__':
# ================================================= # =================================================