Files
sbnews/docker/crawler/scripts/docker_task_runner.py

441 lines
15 KiB
Python

#!/usr/bin/env python3
from __future__ import annotations
import argparse
import os
import shlex
import socket
import subprocess
import sys
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Callable, Iterable, Optional
import pymysql
import yaml
APP_DIR = Path(__file__).resolve().parents[1]
if str(APP_DIR) not in sys.path:
sys.path.insert(0, str(APP_DIR))
from main import get_action_timeout # noqa: E402
from src.core.config import load_config # noqa: E402
STATUS_RUNNING = 0
STATUS_SUCCESS = 1
STATUS_FAILED = 2
STATUS_SKIPPED = 3
OUTPUT_LIMIT = 5000
@dataclass
class CrawlerTask:
task_key: str
task_name: str
action: str
cron_expression: str
active: bool = True
timeout: Optional[int] = None
@classmethod
def from_dict(cls, row: dict) -> "CrawlerTask":
action = str(row.get("action") or row.get("params") or "").strip()
return cls(
task_key=str(row.get("key") or row.get("task_key") or action).strip(),
task_name=str(row.get("name") or row.get("task_name") or action).strip(),
action=action,
cron_expression=str(row.get("cron") or row.get("cron_expression") or "* * * * *").strip(),
active=bool(row.get("active", True)),
timeout=int(row["timeout"]) if row.get("timeout") not in (None, "") else None,
)
def resolved_timeout(self) -> int:
return int(self.timeout or get_action_timeout(self.action))
@dataclass
class CommandResult:
exit_code: int
output: str
elapsed: float
timed_out: bool = False
@dataclass
class TaskRunResult:
log_id: int
status: int
output: str
elapsed: float
error_message: str = ""
class TaskLock:
def __init__(self, action: str, lock_dir: Path):
safe_action = "".join(ch if ch.isalnum() or ch in ("-", "_") else "_" for ch in str(action or "all"))
self.lock_dir = Path(lock_dir)
self.path = self.lock_dir / f"{safe_action}.lock"
self._fp = None
self._locked = False
def acquire(self) -> bool:
self.lock_dir.mkdir(parents=True, exist_ok=True)
self._fp = open(self.path, "a+b")
try:
self._fp.seek(0)
if os.name == "nt":
import msvcrt
msvcrt.locking(self._fp.fileno(), msvcrt.LK_NBLCK, 1)
else:
import fcntl
fcntl.flock(self._fp.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB)
self._locked = True
self._fp.seek(0)
self._fp.truncate()
self._fp.write(f"pid={os.getpid()} started_at={int(time.time())}\n".encode("utf-8"))
self._fp.flush()
return True
except (BlockingIOError, OSError):
self.release()
return False
def release(self) -> None:
if not self._fp:
return
try:
if self._locked:
if os.name == "nt":
import msvcrt
self._fp.seek(0)
msvcrt.locking(self._fp.fileno(), msvcrt.LK_UNLCK, 1)
else:
import fcntl
fcntl.flock(self._fp.fileno(), fcntl.LOCK_UN)
finally:
self._locked = False
self._fp.close()
self._fp = None
class TaskLogStore:
def __init__(self):
cfg = load_config().database
self.prefix = cfg.prefix
self.conn = pymysql.connect(
host=cfg.host,
port=cfg.port,
user=cfg.username,
password=cfg.password,
database=cfg.database,
charset=cfg.charset,
autocommit=True,
)
self.ensure_table()
@property
def table(self) -> str:
return f"`{self.prefix}crawler_task_log`"
def close(self) -> None:
self.conn.close()
def ensure_table(self) -> None:
sql = f"""
CREATE TABLE IF NOT EXISTS {self.table} (
`id` BIGINT UNSIGNED NOT NULL AUTO_INCREMENT,
`task_key` VARCHAR(100) NOT NULL DEFAULT '' COMMENT '任务标识',
`task_name` VARCHAR(100) NOT NULL DEFAULT '' COMMENT '任务名称',
`action` VARCHAR(64) NOT NULL DEFAULT '' COMMENT '爬虫动作',
`cron_expression` VARCHAR(64) NOT NULL DEFAULT '' COMMENT 'cron规则',
`source` VARCHAR(32) NOT NULL DEFAULT 'docker' COMMENT '来源',
`container_id` VARCHAR(128) NOT NULL DEFAULT '' COMMENT '容器ID',
`status` TINYINT UNSIGNED NOT NULL DEFAULT 0 COMMENT '0执行中 1成功 2失败 3跳过',
`trigger_type` TINYINT UNSIGNED NOT NULL DEFAULT 1 COMMENT '1自动 2手动',
`output` TEXT NULL COMMENT '执行输出',
`error_message` TEXT NULL COMMENT '错误信息',
`elapsed` DECIMAL(10,2) NOT NULL DEFAULT 0.00 COMMENT '耗时秒',
`started_at` INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '开始时间',
`finished_at` INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '结束时间',
`create_time` INT UNSIGNED NOT NULL DEFAULT 0,
`update_time` INT UNSIGNED NOT NULL DEFAULT 0,
PRIMARY KEY (`id`),
KEY `idx_action_status` (`action`, `status`),
KEY `idx_create_time` (`create_time`),
KEY `idx_task_key` (`task_key`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COMMENT='Docker爬虫任务执行日志'
"""
with self.conn.cursor() as cur:
cur.execute(sql)
def create_running(self, task: CrawlerTask) -> int:
now = int(time.time())
with self.conn.cursor() as cur:
cur.execute(
f"""
INSERT INTO {self.table}
(`task_key`, `task_name`, `action`, `cron_expression`, `source`, `container_id`,
`status`, `trigger_type`, `output`, `error_message`, `elapsed`,
`started_at`, `finished_at`, `create_time`, `update_time`)
VALUES (%s, %s, %s, %s, 'docker', %s, %s, 1, '', '', 0, %s, 0, %s, %s)
""",
(
task.task_key,
task.task_name,
task.action,
task.cron_expression,
socket.gethostname(),
STATUS_RUNNING,
now,
now,
now,
),
)
return int(cur.lastrowid)
def update_progress(self, log_id: int, output: str, elapsed: float) -> None:
with self.conn.cursor() as cur:
cur.execute(
f"UPDATE {self.table} SET output=%s, elapsed=%s, update_time=%s WHERE id=%s",
(output[-OUTPUT_LIMIT:], round(elapsed, 2), int(time.time()), log_id),
)
def finish(self, log_id: int, status: int, output: str, elapsed: float, error_message: str = "") -> None:
now = int(time.time())
with self.conn.cursor() as cur:
cur.execute(
f"""
UPDATE {self.table}
SET status=%s, output=%s, error_message=%s, elapsed=%s, finished_at=%s, update_time=%s
WHERE id=%s
""",
(status, output[-OUTPUT_LIMIT:], error_message[:2000], round(elapsed, 2), now, now, log_id),
)
def load_tasks(path: Path) -> list[CrawlerTask]:
raw = yaml.safe_load(Path(path).read_text(encoding="utf-8")) or {}
rows = raw if isinstance(raw, list) else raw.get("tasks", [])
return [CrawlerTask.from_dict(row) for row in rows]
def render_crontab(tasks: Iterable[CrawlerTask], python_bin: str = "/usr/local/bin/python", app_dir: str = "/app") -> str:
lines = [
"SHELL=/bin/bash",
"BASH_ENV=/etc/sport-era-crawler/env.sh",
"PATH=/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin",
"",
]
for task in tasks:
if not task.active:
continue
args = [
python_bin,
"scripts/docker_task_runner.py",
"run",
task.action,
"--task-key",
task.task_key,
"--task-name",
task.task_name,
"--cron-expression",
task.cron_expression,
]
if task.timeout:
args.extend(["--timeout", str(task.timeout)])
cmd = " ".join(shlex.quote(str(arg)) for arg in args)
lines.append(f"{task.cron_expression} root cd {shlex.quote(app_dir)} && {cmd}")
lines.append("")
return "\n".join(lines)
def run_main_action(action: str, timeout: int, output_callback: Callable[[str], None]) -> CommandResult:
start = time.monotonic()
env = os.environ.copy()
env["PYTHONUNBUFFERED"] = "1"
env["DQD_RUNNING_UNDER_TASK_RUNNER"] = "1"
proc = subprocess.Popen(
[sys.executable, "main.py", action],
cwd=str(APP_DIR),
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
encoding="utf-8",
errors="replace",
env=env,
)
output_parts: list[str] = []
timed_out = False
if os.name == "nt":
try:
stdout, _ = proc.communicate(timeout=timeout)
stdout = stdout or ""
output_parts.append(stdout)
if stdout:
output_callback(stdout)
except subprocess.TimeoutExpired:
timed_out = True
proc.kill()
stdout, _ = proc.communicate()
stdout = stdout or ""
output_parts.append(stdout)
if stdout:
output_callback(stdout)
else:
import selectors
assert proc.stdout is not None
selector = selectors.DefaultSelector()
selector.register(proc.stdout, selectors.EVENT_READ)
try:
while proc.poll() is None:
remaining = timeout - (time.monotonic() - start)
if remaining <= 0:
timed_out = True
proc.kill()
break
for key, _ in selector.select(timeout=min(1.0, remaining)):
line = key.fileobj.readline()
if line:
output_parts.append(line)
output_callback(line)
rest = proc.stdout.read()
if rest:
output_parts.append(rest)
output_callback(rest)
finally:
selector.close()
elapsed = round(time.monotonic() - start, 2)
return CommandResult(
exit_code=proc.returncode if proc.returncode is not None else -1,
output="".join(output_parts),
elapsed=elapsed,
timed_out=timed_out,
)
def run_task(
task: CrawlerTask,
store=None,
lock_dir: Optional[Path] = None,
command_runner: Callable[[str, int, Callable[[str], None]], CommandResult] = run_main_action,
flush_interval: float = 5.0,
) -> TaskRunResult:
if store is None:
store = TaskLogStore()
if hasattr(store, "ensure_table"):
store.ensure_table()
log_id = int(store.create_running(task))
start = time.monotonic()
lock = TaskLock(task.action, Path(lock_dir or os.environ.get("DQD_LOCK_DIR") or APP_DIR / "data" / "locks"))
output_parts: list[str] = []
last_flush = 0.0
def append_output(text: str) -> None:
nonlocal last_flush
output_parts.append(text)
now = time.monotonic()
elapsed = round(now - start, 2)
if now - last_flush >= flush_interval:
store.update_progress(log_id, "".join(output_parts), elapsed)
last_flush = now
if not lock.acquire():
output = f"{task.action} 已有任务执行中,跳过本次"
store.finish(log_id, STATUS_SKIPPED, output, 0, "")
return TaskRunResult(log_id, STATUS_SKIPPED, output, 0, "")
try:
result = command_runner(task.action, task.resolved_timeout(), append_output)
output = result.output or "".join(output_parts) or "执行完成"
if result.timed_out:
status = STATUS_FAILED
error = f"{task.action} 执行超时,超过 {task.resolved_timeout()} 秒"
elif result.exit_code == 0:
status = STATUS_SUCCESS
error = ""
else:
status = STATUS_FAILED
error = f"{task.action} 执行失败,退出码 {result.exit_code}"
store.finish(log_id, status, output, result.elapsed, error)
return TaskRunResult(log_id, status, output, result.elapsed, error)
except Exception as exc:
elapsed = round(time.monotonic() - start, 2)
output = "".join(output_parts) or str(exc)
store.finish(log_id, STATUS_FAILED, output, elapsed, str(exc))
return TaskRunResult(log_id, STATUS_FAILED, output, elapsed, str(exc))
finally:
lock.release()
def build_task_from_args(args) -> CrawlerTask:
if args.tasks:
for task in load_tasks(Path(args.tasks)):
if task.action == args.action:
if args.timeout:
task.timeout = args.timeout
return task
return CrawlerTask(
task_key=args.task_key or args.action,
task_name=args.task_name or args.action,
action=args.action,
cron_expression=args.cron_expression or "",
active=True,
timeout=args.timeout,
)
def parse_args(argv: Optional[list[str]] = None):
parser = argparse.ArgumentParser(description="Sport Era Docker crawler task runner")
sub = parser.add_subparsers(dest="command", required=True)
run = sub.add_parser("run")
run.add_argument("action")
run.add_argument("--task-key", default="")
run.add_argument("--task-name", default="")
run.add_argument("--cron-expression", default="")
run.add_argument("--timeout", type=int, default=None)
run.add_argument("--tasks", default="")
render = sub.add_parser("render-cron")
render.add_argument("--tasks", required=True)
render.add_argument("--output", required=True)
render.add_argument("--python-bin", default="/usr/local/bin/python")
render.add_argument("--app-dir", default="/app")
return parser.parse_args(argv)
def main(argv: Optional[list[str]] = None) -> int:
args = parse_args(argv)
if args.command == "render-cron":
tasks = load_tasks(Path(args.tasks))
content = render_crontab(tasks, python_bin=args.python_bin, app_dir=args.app_dir)
output = Path(args.output)
output.parent.mkdir(parents=True, exist_ok=True)
output.write_text(content, encoding="utf-8")
output.chmod(0o644)
return 0
if args.command == "run":
store = TaskLogStore()
try:
result = run_task(build_task_from_args(args), store=store)
return 0 if result.status in (STATUS_SUCCESS, STATUS_SKIPPED) else 1
finally:
store.close()
return 1
if __name__ == "__main__":
raise SystemExit(main())