Просмотр исходного кода

feat: add vertical video query API

zhangliang 1 день назад
Родитель
Сommit
18f956e7da

+ 2 - 0
.gitignore

@@ -59,3 +59,5 @@ target/
 
 .idea/
 *.DS_Store
+.env
+/run/

+ 72 - 3
README.md

@@ -44,6 +44,76 @@
 - ✅ 视频封装为标准 `VideoItem`,统一推送到 MQ
 - ✅ 任务执行成功后再确认 ACK,保证一致性
 - ✅ 完善的配置管理(验证、健康检查、文档生成、命令行工具)
+- ✅ 异步高并发垂直 Spider 视频查询 API
+
+### 垂直 Spider API
+
+完整的接口参数、筛选规则、去重分页、响应、错误码、日志和部署说明见
+[`docs/chui_zhi_video_api.md`](docs/chui_zhi_video_api.md)。
+
+API 默认只监听 `127.0.0.1:8888`,并强制配置 Token。跨服务器访问时应绑定内网/VPN
+地址或使用网关反向代理,不建议直接暴露到公网。
+
+```bash
+export CHUI_ZHI_API_TOKEN='replace-with-a-strong-token'
+sh run_api.sh prod
+```
+
+- `POST /api/v1/crawler/videos/query`
+- `GET /health` 和 `GET /ready` 仅用于本机健康检查,Nginx 示例默认禁止公网访问
+- 调用方必须通过 `X-API-Key` 传递 `CHUI_ZHI_API_TOKEN`
+- `start_time`、`end_time` 推荐传13位毫秒时间戳,API 会统一转换为东八区数据库查询时间
+- 未传时间时默认按 `create_time` 查询最近3天;传入 `keywords` 时按 `video_title` 模糊匹配
+- 响应通过 `has_more` 和 `next_cursor` 表示是否存在下一页;调用方下一页原样回传 `cursor`
+- 每次API调用都会通过有界队列批量向阿里云SLS上报URL、请求参数、SQL结果长度、查询耗时、请求总耗时、成功状态、失败阶段、异常类型、HTTP状态码、request_id和错误消息;异常时 `message` 包含脱敏后的原始堆栈;本地日志仅记录
+  URL、请求参数、状态码,失败时额外记录错误原因
+- 可通过 `API_HOST`、`API_PORT`、`API_MAX_CONCURRENT_REQUESTS`、`DB_POOL_SIZE`、`API_QUERY_TIMEOUT`、`API_QUEUE_TIMEOUT`、`API_LOG_QUEUE_SIZE`、`API_MAX_LIMIT` 调整运行参数
+- 当前只提供垂直视频查询接口,路由在 `api/app.py` 中直接注册,参数校验、SQL 构造和查询逻辑集中在
+  `api/chui_zhi/videos.py`
+
+筛选条件使用 `field/operator/value` 结构;没有筛选条件时传空数组或省略 `filters`:
+
+```json
+{
+  "filters": [
+    {"field": "like_cnt", "operator": ">", "value": 3},
+    {"field": "play_cnt", "operator": ">=", "value": 1000},
+    {"field": "duration", "operator": "between", "value": [60, 300]},
+    {"field": "publish_time", "operator": ">=", "value": 1785772800000}
+  ]
+}
+```
+
+支持 `>`、`>=`、`=`、`<`、`<=`、`between`、`in` 和 `not_in`。API 会校验字段、操作符和值,
+并使用参数化 SQL。
+
+第一页不传 `cursor`。如果 `has_more=true`,下一页将响应里的 `next_cursor` 原样传回:
+
+```json
+{
+  "cursor": {"id": 6815000}
+}
+```
+
+API在筛选范围内按 `out_video_id` 分组并保留最大 `id` 的记录,然后只回表查询当前页完整字段;
+空 `out_video_id` 按各自 `id` 保留。游标作用于分组后的最大 `id`,避免重复记录跨页再次出现。
+
+公网部署时,API 仍保持 `API_HOST=127.0.0.1` 和 `API_PORT=8888`,由 Nginx 对外提供 HTTPS:
+
+```bash
+# 1. 配置 API 环境变量
+API_HOST=127.0.0.1
+API_PORT=8888
+CHUI_ZHI_API_TOKEN=replace-with-a-strong-token
+
+# 2. 使用现有部署脚本同时启动主爬虫和API
+bash deploy.sh
+
+# 3. 首次部署时安装 Nginx 配置(先替换域名、证书路径和可选 IP 白名单)
+cp deploy/nginx/chui_zhi_api.conf.example /etc/nginx/conf.d/chui_zhi_api.conf
+nginx -t
+systemctl reload nginx
+```
 
 ---
 
@@ -248,8 +318,8 @@ pip install -r requirements.txt
 # 启动服务
 python main.py
 
-# 或使用部署脚本
-sh deploy.sh
+# 推荐:同时部署主爬虫和垂直视频API,并执行API就绪检查
+bash deploy.sh
 ```
 
 ## 🧰 常用操作
@@ -309,4 +379,3 @@ sh deploy.sh
    ```
    https://w42nne6hzg.feishu.cn/sheets/U5dXsSlPOhiNNCtEfgqcm1iYnpf?sheet=K0gA9Y
    ```
-   

+ 1 - 0
api/__init__.py

@@ -0,0 +1 @@
+"""AutoScraperX HTTP API package."""

+ 475 - 0
api/app.py

@@ -0,0 +1,475 @@
+import asyncio
+import hmac
+import json
+import re
+import time
+import traceback
+import uuid
+from typing import Any
+
+from aiohttp import web
+from pydantic import ValidationError
+
+from api.chui_zhi.videos import (
+    MYSQL_KEY,
+    QUERY_METRICS_KEY,
+    QUERY_SEMAPHORE_KEY,
+    BusinessValidationError,
+    DatabaseQueryError,
+    ServiceBusyError,
+    api_response,
+    query_videos,
+)
+from config import settings
+from core.base.async_mysql_client import AsyncMySQLClient
+from core.utils.log.logger_manager import LoggerManager
+
+
+LOGGER_KEY = web.AppKey('logger', object)
+ALIYUN_LOGGER_KEY = web.AppKey('aliyun_logger', object)
+CLOUD_LOG_QUEUE_KEY = web.AppKey('cloud_log_queue', asyncio.Queue)
+CLOUD_LOG_WORKER_KEY = web.AppKey('cloud_log_worker', asyncio.Task)
+ERROR_REASON_KEY = web.AppKey('error_reason', str)
+ERROR_TRACEBACK_KEY = web.AppKey('error_traceback', str)
+ERROR_TYPE_KEY = web.AppKey('error_type', str)
+FAILURE_STAGE_KEY = web.AppKey('failure_stage', str)
+REQUEST_ID_KEY = web.AppKey('request_id', str)
+REQUEST_PARAMS_KEY = web.AppKey('request_params', dict)
+REQUEST_STARTED_AT_KEY = web.AppKey('request_started_at', float)
+API_PATH = '/api/v1/crawler/videos/query'
+HEALTH_PATH = '/health'
+READY_PATH = '/ready'
+AUTH_EXEMPT_PATHS = frozenset({HEALTH_PATH, READY_PATH})
+MAX_CLOUD_LOG_VALUE_LENGTH = 64 * 1024
+MAX_LOG_BATCH_SIZE = 50
+MAX_LOG_RETRIES = 3
+SENSITIVE_FIELD_PARTS = ('authorization', 'api_key', 'apikey', 'token', 'password', 'secret')
+
+
+def print_api_routes(app: web.Application) -> None:
+    print('\n已注册 API:', flush=True)
+    for route in app.router.routes():
+        path = route.resource.canonical
+        handler_name = getattr(route.handler, '__name__', route.handler.__class__.__name__)
+        print(f'  {route.method:<6} {path:<36} -> {handler_name}', flush=True)
+    print('', flush=True)
+
+
+def parse_json_text(text: str):
+    """尽量将日志内容还原为JSON;非JSON内容保留原字符串。"""
+    if not text:
+        return None
+    try:
+        return json.loads(text)
+    except (TypeError, ValueError):
+        return text
+
+
+def format_validation_error(exc: ValidationError) -> str:
+    """把Pydantic错误转换为调用方可直接定位的简洁参数信息。"""
+    messages = []
+    for error in exc.errors(include_url=False):
+        location = '.'.join(str(item) for item in error.get('loc', ())) or 'body'
+        error_type = error.get('type', '')
+        input_value = str(error.get('input', ''))[:100]
+        context = error.get('ctx') or {}
+        if error_type == 'extra_forbidden':
+            message = f'不支持的参数: {location}'
+        elif error_type == 'literal_error':
+            message = f'{location}不支持值 {input_value},允许值为{context.get("expected", "白名单值")}'
+        elif error_type == 'missing':
+            message = f'{location}不能为空'
+        elif error_type == 'list_too_long':
+            message = f'{location}最多允许{context.get("max_length")}项'
+        elif error_type == 'list_too_short':
+            message = f'{location}至少需要{context.get("min_length")}项'
+        elif error_type == 'less_than_equal':
+            message = f'{location}不能大于{context.get("le")}'
+        elif error_type == 'greater_than_equal':
+            message = f'{location}不能小于{context.get("ge")}'
+        elif error_type == 'value_error':
+            detail = error.get('msg', '').removeprefix('Value error, ')
+            message = f'{location}: {detail}'
+        else:
+            message = f'{location}: {error.get("msg", "参数格式错误")}'
+        messages.append(message)
+    return '参数校验失败: ' + '; '.join(messages)
+
+
+def redact_value(value: Any):
+    """递归脱敏凭证字段,避免未来扩展请求参数时意外泄露。"""
+    if isinstance(value, dict):
+        return {
+            key: '***' if any(part in str(key).lower() for part in SENSITIVE_FIELD_PARTS)
+            else redact_value(item)
+            for key, item in value.items()
+        }
+    if isinstance(value, list):
+        return [redact_value(item) for item in value]
+    return value
+
+
+def redact_text(value: str) -> str:
+    """保留异常堆栈,同时移除配置中已知的密钥内容。"""
+    result = value
+    for secret in (
+        settings.CHUI_ZHI_API_TOKEN,
+        settings.DB_PASSWORD,
+        settings.ALIYUN_ACCESS_KEY_ID,
+        settings.ALIYUN_ACCESS_KEY_SECRET,
+    ):
+        if secret:
+            result = result.replace(str(secret), '***')
+    return result
+
+
+def limit_cloud_log_value(value):
+    """限制单个日志字段大小,避免请求参数超过SLS单条日志限制。"""
+    value = redact_value(value)
+    text = json.dumps(value, ensure_ascii=False, default=str)
+    if len(text) <= MAX_CLOUD_LOG_VALUE_LENGTH:
+        return value
+    return {
+        'truncated': True,
+        'original_length': len(text),
+        'content': text[:MAX_CLOUD_LOG_VALUE_LENGTH],
+    }
+
+
+def get_request_url(request: web.Request) -> str:
+    scheme = request.headers.get('X-Forwarded-Proto', request.scheme).split(',', 1)[0].strip()
+    host = request.headers.get('X-Forwarded-Host', request.host).split(',', 1)[0].strip()
+    return f'{scheme}://{host}{request.rel_url}'
+
+
+async def send_cloud_log_batch(app: web.Application, events: list[dict]) -> None:
+    await asyncio.wait_for(
+        asyncio.to_thread(app[ALIYUN_LOGGER_KEY].logging_batch, events),
+        timeout=settings.API_LOG_FLUSH_TIMEOUT,
+    )
+
+
+async def cloud_log_worker(app: web.Application) -> None:
+    """后台批量上报SLS,日志服务异常不阻塞业务请求。"""
+    queue = app[CLOUD_LOG_QUEUE_KEY]
+    logger = app[LOGGER_KEY]
+    while True:
+        event = await queue.get()
+        if event is None:
+            queue.task_done()
+            return
+
+        events = [event]
+        while len(events) < MAX_LOG_BATCH_SIZE:
+            try:
+                next_event = queue.get_nowait()
+            except asyncio.QueueEmpty:
+                break
+            if next_event is None:
+                queue.task_done()
+                break
+            events.append(next_event)
+
+        try:
+            for attempt in range(MAX_LOG_RETRIES):
+                try:
+                    await send_cloud_log_batch(app, events)
+                    break
+                except Exception:
+                    if attempt + 1 >= MAX_LOG_RETRIES:
+                        logger.exception(f'阿里云API日志批量上报失败: count={len(events)}')
+                    else:
+                        await asyncio.sleep(0.5 * (2 ** attempt))
+        finally:
+            for _ in events:
+                queue.task_done()
+
+
+async def enqueue_cloud_log(app: web.Application, event: dict) -> None:
+    queue = app.get(CLOUD_LOG_QUEUE_KEY)
+    if queue is None:
+        # 单元测试或嵌入运行未启动cleanup context时,仍可验证日志行为。
+        try:
+            await send_cloud_log_batch(app, [event])
+        except Exception:
+            app[LOGGER_KEY].exception('阿里云API日志上报失败')
+        return
+    try:
+        queue.put_nowait(event)
+    except asyncio.QueueFull:
+        app[LOGGER_KEY].error(
+            f'阿里云API日志队列已满,丢弃日志: request_id={event.get("trace_id", "")}'
+        )
+
+
+async def report_api_request(
+    request: web.Request,
+    response: web.StreamResponse | None,
+    exception: Exception | None = None,
+) -> None:
+    """本地日志同步落地,SLS日志仅入队,不增加接口响应延迟。"""
+    logger = request.app[LOGGER_KEY]
+    status_code = response.status if response is not None else getattr(exception, 'status', 500)
+    request_duration_ms = round(
+        (time.perf_counter() - request.get(REQUEST_STARTED_AT_KEY, time.perf_counter())) * 1000,
+        2,
+    )
+
+    request_params = request.get(REQUEST_PARAMS_KEY, {'query': dict(request.query), 'body': None})
+    safe_request_params = limit_cloud_log_value(request_params)
+    request_id = request[REQUEST_ID_KEY]
+    error_message = request.get(ERROR_TRACEBACK_KEY, '') or request.get(ERROR_REASON_KEY, '')
+    if not error_message and exception is not None:
+        error_message = f'{type(exception).__name__}: {exception}'
+    error_message = redact_text(error_message)
+    error_type = request.get(ERROR_TYPE_KEY, '')
+    if not error_type and exception is not None:
+        error_type = type(exception).__name__
+    failure_stage = request.get(FAILURE_STAGE_KEY, '')
+
+    request_url = get_request_url(request)
+    local_message = (
+        f'API请求 request_id={request_id} method={request.method} url={request_url} '
+        f'params={json.dumps(safe_request_params, ensure_ascii=False, default=str)} '
+        f'status={status_code} duration_ms={request_duration_ms}'
+    )
+    if error_message:
+        logger.error(f'{local_message} error={error_message}')
+    else:
+        logger.info(local_message)
+
+    query_metrics = request.get(QUERY_METRICS_KEY, {})
+    request_body = request_params.get('body') if isinstance(request_params, dict) else None
+    request_body = request_body if isinstance(request_body, dict) else {}
+    cloud_data = {
+        'url': request_url,
+        'path': request.path,
+        'method': request.method,
+        'request_id': request_id,
+        'request_params': safe_request_params,
+        'platforms': request_body.get('platforms'),
+        'status_code': status_code,
+        'success': status_code < 400,
+        'request_duration_ms': request_duration_ms,
+        'failure_stage': failure_stage,
+        'error_type': error_type,
+        'message': error_message,
+        'query_result_count': query_metrics.get('result_count'),
+        'query_duration_ms': query_metrics.get('duration_ms'),
+        'query_success': query_metrics.get('success'),
+    }
+    await enqueue_cloud_log(request.app, {
+        'code': '2000' if status_code < 400 else '9000',
+        'message': 'API请求成功' if status_code < 400 else error_message,
+        'data': cloud_data,
+        'trace_id': request_id,
+    })
+
+
+@web.middleware
+async def request_context_middleware(request: web.Request, handler):
+    """建立请求上下文,确保成功、失败和未鉴权请求都有访问日志。"""
+    incoming_request_id = request.headers.get('X-Request-ID', '').strip()
+    if not re.fullmatch(r'[A-Za-z0-9._-]{1,128}', incoming_request_id):
+        incoming_request_id = uuid.uuid4().hex
+    request[REQUEST_ID_KEY] = incoming_request_id
+    request[REQUEST_STARTED_AT_KEY] = time.perf_counter()
+    request[REQUEST_PARAMS_KEY] = {
+        'query': dict(request.query),
+        'body': None,
+    }
+    response = None
+    exception = None
+    try:
+        response = await handler(request)
+        response.headers['X-Request-ID'] = request[REQUEST_ID_KEY]
+        return response
+    except Exception as exc:
+        exception = exc
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        raise
+    finally:
+        try:
+            await report_api_request(request, response, exception)
+        except Exception:
+            # 日志异常不能覆盖API原始响应或异常。
+            request.app[LOGGER_KEY].exception('API访问日志记录失败')
+
+
+@web.middleware
+async def request_body_middleware(request: web.Request, handler):
+    """鉴权通过后再读取请求体,避免未授权请求消耗JSON解析资源。"""
+    raw_body = await request.read()
+    request[REQUEST_PARAMS_KEY]['body'] = parse_json_text(
+        raw_body.decode('utf-8', errors='replace')
+    )
+    return await handler(request)
+
+
+@web.middleware
+async def error_middleware(request: web.Request, handler):
+    try:
+        return await handler(request)
+    except json.JSONDecodeError as exc:
+        request[ERROR_REASON_KEY] = str(exc)
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'request_validation'
+        return api_response(code=400, msg=f'请求JSON格式错误: {exc.msg}')
+    except ValidationError as exc:
+        message = format_validation_error(exc)
+        request[ERROR_REASON_KEY] = message
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'request_validation'
+        return api_response(code=400, msg=message, data=exc.errors(include_url=False))
+    except BusinessValidationError as exc:
+        request[ERROR_REASON_KEY] = str(exc)
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'business_validation'
+        return api_response(code=422, msg=str(exc))
+    except ServiceBusyError as exc:
+        request[ERROR_REASON_KEY] = str(exc) or 'server busy'
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'concurrency_limit'
+        return api_response(code=503, msg='server busy')
+    except asyncio.TimeoutError as exc:
+        request[ERROR_REASON_KEY] = str(exc) or 'query timeout'
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'database_query'
+        return api_response(code=504, msg='query timeout')
+    except DatabaseQueryError as exc:
+        request[ERROR_REASON_KEY] = str(exc)
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'database_query'
+        request.app[LOGGER_KEY].exception(str(exc))
+        return api_response(code=500, msg='database query failed')
+    except web.HTTPException:
+        raise
+    except Exception as exc:
+        request[ERROR_REASON_KEY] = f'{type(exc).__name__}: {exc}'
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'internal'
+        request.app[LOGGER_KEY].exception(f'API request failed: {exc}')
+        return api_response(code=500, msg='internal server error')
+
+
+@web.middleware
+async def auth_middleware(request: web.Request, handler):
+    if request.path in AUTH_EXEMPT_PATHS:
+        return await handler(request)
+    supplied_token = request.headers.get('X-API-Key', '')
+    if not hmac.compare_digest(supplied_token, settings.CHUI_ZHI_API_TOKEN):
+        request[ERROR_REASON_KEY] = 'unauthorized'
+        request[ERROR_TYPE_KEY] = 'AuthenticationError'
+        request[FAILURE_STAGE_KEY] = 'authentication'
+        return api_response(code=401, msg='unauthorized')
+    return await handler(request)
+
+
+async def health(_request: web.Request) -> web.Response:
+    return api_response({'status': 'ok'})
+
+
+async def ready(request: web.Request) -> web.Response:
+    try:
+        async with asyncio.timeout(min(3, settings.API_QUERY_TIMEOUT)):
+            row = await request.app[MYSQL_KEY].fetch_one('SELECT 1 AS ok')
+        if not row or row.get('ok') != 1:
+            raise DatabaseQueryError('database readiness check failed')
+        return api_response({'status': 'ready'})
+    except Exception as exc:
+        request[ERROR_REASON_KEY] = f'{type(exc).__name__}: {exc}'
+        request[ERROR_TRACEBACK_KEY] = traceback.format_exc()
+        request[ERROR_TYPE_KEY] = type(exc).__name__
+        request[FAILURE_STAGE_KEY] = 'readiness'
+        return api_response(code=503, msg='not ready')
+
+
+async def app_context(app: web.Application):
+    if not settings.CHUI_ZHI_API_TOKEN:
+        raise RuntimeError('CHUI_ZHI_API_TOKEN未配置,拒绝启动垂直spider API')
+
+    logger = LoggerManager.get_logger(platform='chui_zhi', mode='api')
+    aliyun_logger = LoggerManager.get_aliyun_logger(platform='chui_zhi', mode='api')
+    mysql = AsyncMySQLClient(
+        host=settings.DB_HOST,
+        port=settings.DB_PORT,
+        user=settings.DB_USER,
+        password=settings.DB_PASSWORD,
+        db=settings.DB_NAME,
+        charset=settings.DB_CHARSET,
+        minsize=min(5, settings.DB_POOL_SIZE),
+        maxsize=settings.DB_POOL_SIZE,
+        pool_recycle=settings.DB_POOL_RECYCLE,
+        logger=logger,
+        # API中间件统一上报异常,避免数据库客户端同步调用SLS。
+        aliyun_logr=None,
+    )
+    await mysql.init_pool()
+    query_concurrency = min(settings.API_MAX_CONCURRENT_REQUESTS, settings.DB_POOL_SIZE)
+    app[MYSQL_KEY] = mysql
+    app[QUERY_SEMAPHORE_KEY] = asyncio.Semaphore(query_concurrency)
+    app[LOGGER_KEY] = logger
+    app[ALIYUN_LOGGER_KEY] = aliyun_logger
+    app[CLOUD_LOG_QUEUE_KEY] = asyncio.Queue(maxsize=settings.API_LOG_QUEUE_SIZE)
+    app[CLOUD_LOG_WORKER_KEY] = asyncio.create_task(cloud_log_worker(app))
+    logger.info(
+        f'垂直spider API启动: db_pool={settings.DB_POOL_SIZE}, '
+        f'query_concurrency={query_concurrency}'
+    )
+    print_api_routes(app)
+    yield
+
+    queue = app[CLOUD_LOG_QUEUE_KEY]
+    try:
+        await asyncio.wait_for(queue.join(), timeout=settings.API_LOG_FLUSH_TIMEOUT)
+    except asyncio.TimeoutError:
+        logger.error(f'API关闭时日志队列未完全清空: remaining={queue.qsize()}')
+        app[CLOUD_LOG_WORKER_KEY].cancel()
+        try:
+            await app[CLOUD_LOG_WORKER_KEY]
+        except asyncio.CancelledError:
+            pass
+    else:
+        await queue.put(None)
+        await app[CLOUD_LOG_WORKER_KEY]
+    await mysql.close()
+    logger.info('垂直spider API已关闭')
+
+
+def create_app() -> web.Application:
+    app = web.Application(
+        middlewares=[
+            request_context_middleware,
+            error_middleware,
+            auth_middleware,
+            request_body_middleware,
+        ],
+        client_max_size=1024 * 1024,
+    )
+    app.cleanup_ctx.append(app_context)
+    app.add_routes([
+        web.post(API_PATH, query_videos, name='chui_zhi_videos'),
+        web.get(HEALTH_PATH, health, name='health'),
+        web.get(READY_PATH, ready, name='ready'),
+    ])
+    return app
+
+
+def main():
+    web.run_app(
+        create_app(),
+        host=settings.API_HOST,
+        port=settings.API_PORT,
+        access_log=None,
+    )
+
+
+if __name__ == '__main__':
+    main()

+ 0 - 0
api/chui_zhi/__init__.py


+ 373 - 0
api/chui_zhi/videos.py

@@ -0,0 +1,373 @@
+import asyncio
+import json
+import math
+import time
+from datetime import datetime, timedelta, timezone
+from typing import Any, List, Literal, Optional, Tuple
+
+from aiohttp import web
+from pydantic import BaseModel, ConfigDict, Field, field_validator, model_validator
+
+from config import settings
+from core.base.async_mysql_client import AsyncMySQLClient
+
+
+MYSQL_KEY = web.AppKey('mysql', AsyncMySQLClient)
+QUERY_SEMAPHORE_KEY = web.AppKey('query_semaphore', asyncio.Semaphore)
+QUERY_METRICS_KEY = web.AppKey('query_metrics', dict)
+CHINA_TIMEZONE = timezone(timedelta(hours=8))
+
+FilterField = Literal[
+    'play_cnt',
+    'like_cnt',
+    'share_cnt',
+    'collection_cnt',
+    'comment_cnt',
+    'duration',
+    'publish_time',
+]
+FilterOperator = Literal['>', '>=', '=', '<', '<=', 'between', 'in', 'not_in']
+SqlFragment = Tuple[str, List[Any]]
+SUPPORTED_PLATFORMS = frozenset({'xiaoniangao', 'xiaoniangaotuijianliu'})
+MAX_FILTER_SET_VALUES = 100
+
+
+class ServiceBusyError(Exception):
+    """API并发已满,暂时无法执行查询。"""
+
+
+class BusinessValidationError(Exception):
+    """请求格式正确,但业务筛选值无法执行。"""
+
+
+class DatabaseQueryError(Exception):
+    """数据库查询执行失败。"""
+
+
+# ==================== 请求参数 ====================
+
+class FilterCondition(BaseModel):
+    """单个结构化筛选条件,字段和操作符均通过白名单限制。"""
+
+    model_config = ConfigDict(extra='forbid')
+
+    field: FilterField
+    operator: FilterOperator
+    value: Any
+
+    @model_validator(mode='after')
+    def validate_collection_size(self):
+        if isinstance(self.value, list) and len(self.value) > MAX_FILTER_SET_VALUES:
+            raise ValueError(f'筛选值最多允许{MAX_FILTER_SET_VALUES}项')
+        return self
+
+
+class PageCursor(BaseModel):
+    """稳定翻页游标,对应上一页最后一条数据的自增主键。"""
+
+    model_config = ConfigDict(extra='forbid')
+
+    id: int = Field(gt=0)
+
+
+class VideoQueryRequest(BaseModel):
+    model_config = ConfigDict(extra='forbid')
+
+    platforms: List[str] = Field(
+        default_factory=lambda: ['xiaoniangao', 'xiaoniangaotuijianliu'],
+        min_length=1,
+        max_length=20,
+    )
+    start_time: Optional[datetime] = None
+    end_time: Optional[datetime] = None
+    keywords: List[str] = Field(default_factory=list, max_length=20)
+    filter_match_mode: Literal[1, 2] = 2  # 1=OR,2=AND
+    filters: List[FilterCondition] = Field(default_factory=list, max_length=50)
+    limit: int = Field(default=500, ge=1, le=settings.API_MAX_LIMIT)
+    cursor: Optional[PageCursor] = None
+
+    @field_validator('start_time', 'end_time', mode='before')
+    @classmethod
+    def normalize_query_time(cls, value):
+        """毫秒时间戳统一转换为东八区的数据库查询时间,同时兼容日期字符串。"""
+        is_number = isinstance(value, (int, float)) and not isinstance(value, bool)
+        is_digit_string = isinstance(value, str) and value.strip().isdigit()
+        if is_number or is_digit_string:
+            try:
+                timestamp = float(value)
+                if not math.isfinite(timestamp):
+                    raise ValueError('时间戳必须是有限数字')
+                if abs(timestamp) >= 10_000_000_000:
+                    timestamp /= 1000
+                return datetime.fromtimestamp(timestamp, tz=CHINA_TIMEZONE).replace(tzinfo=None)
+            except (OverflowError, OSError, ValueError) as exc:
+                raise ValueError(f'时间格式错误: {value}') from exc
+        return value
+
+    @field_validator('start_time', 'end_time', mode='after')
+    @classmethod
+    def normalize_timezone(cls, value):
+        if value is not None and value.tzinfo is not None:
+            return value.astimezone(CHINA_TIMEZONE).replace(tzinfo=None)
+        return value
+
+    @field_validator('platforms')
+    @classmethod
+    def validate_platforms(cls, values: List[str]) -> List[str]:
+        normalized = [value.strip() for value in values if value and value.strip()]
+        if not normalized:
+            raise ValueError('platforms不能为空')
+        if any(len(value) > 50 for value in normalized):
+            raise ValueError('platform长度不能超过50')
+        normalized = list(dict.fromkeys(normalized))
+        unsupported = sorted(set(normalized) - SUPPORTED_PLATFORMS)
+        if unsupported:
+            raise ValueError(f'不支持的平台: {", ".join(unsupported)}')
+        return normalized
+
+    @field_validator('keywords')
+    @classmethod
+    def validate_keywords(cls, values: List[str]) -> List[str]:
+        normalized = [value.strip() for value in values if value and value.strip()]
+        if values and not normalized:
+            raise ValueError('keywords不能全部为空')
+        if any(len(value) > 200 for value in normalized):
+            raise ValueError('keyword长度不能超过200')
+        return normalized
+
+    @model_validator(mode='after')
+    def fill_and_validate_time_range(self):
+        now = datetime.now(CHINA_TIMEZONE).replace(tzinfo=None)
+        if self.start_time is None and self.end_time is None:
+            self.end_time = now
+            self.start_time = self.end_time - timedelta(days=3)
+        elif self.start_time is None:
+            self.start_time = self.end_time - timedelta(days=3)
+        elif self.end_time is None:
+            self.end_time = now
+        if self.start_time >= self.end_time:
+            raise ValueError('start_time必须早于end_time')
+        return self
+
+
+# ==================== 功能实现 ====================
+
+SELECT_COLUMNS = (
+    'video_id', 'user_id', 'out_user_id', 'platform', 'strategy',
+    'out_video_id', 'video_title', 'cover_url', 'video_url', 'duration',
+    'publish_time', 'play_cnt', 'like_cnt', 'share_cnt', 'collection_cnt',
+    'comment_cnt', 'width', 'height', 'id', 'create_time',
+)
+COMPARISON_OPERATORS = frozenset({'>', '>=', '=', '<', '<='})
+
+
+def normalize_filter_value(field: FilterField, value: Any) -> Any:
+    """将请求值转换为数据库可比较的数字或日期字符串。"""
+    if field == 'publish_time':
+        if isinstance(value, datetime):
+            if value.tzinfo is not None:
+                value = value.astimezone(CHINA_TIMEZONE).replace(tzinfo=None)
+            return value.strftime('%Y-%m-%d %H:%M:%S')
+        if isinstance(value, (int, float)) or (isinstance(value, str) and value.strip().isdigit()):
+            timestamp = float(value)
+            if timestamp >= 10_000_000_000:
+                timestamp /= 1000
+            return datetime.fromtimestamp(timestamp, tz=CHINA_TIMEZONE).strftime('%Y-%m-%d %H:%M:%S')
+        try:
+            parsed = datetime.fromisoformat(str(value).strip())
+            if parsed.tzinfo is not None:
+                parsed = parsed.astimezone(CHINA_TIMEZONE).replace(tzinfo=None)
+            return parsed.strftime('%Y-%m-%d %H:%M:%S')
+        except ValueError as exc:
+            raise BusinessValidationError(f'publish_time筛选值格式错误: {value}') from exc
+    if isinstance(value, bool):
+        raise BusinessValidationError(f'{field}筛选值必须是数字: {value}')
+    try:
+        number = float(value)
+    except (TypeError, ValueError) as exc:
+        raise BusinessValidationError(f'{field}筛选值必须是数字: {value}') from exc
+    if not math.isfinite(number):
+        raise BusinessValidationError(f'{field}筛选值必须是有限数字: {value}')
+    return int(number) if number.is_integer() else number
+
+
+def qualified_column(field: str, table_alias: Optional[str] = None) -> str:
+    return f'`{table_alias}`.`{field}`' if table_alias else f'`{field}`'
+
+
+def compile_filter_condition(
+    condition: FilterCondition,
+    table_alias: Optional[str] = None,
+) -> SqlFragment:
+    """将结构化筛选条件转换成参数化SQL。"""
+    field, operator, value = condition.field, condition.operator, condition.value
+    column = qualified_column(field, table_alias)
+
+    if operator in COMPARISON_OPERATORS:
+        if value in (None, '') or isinstance(value, list):
+            raise BusinessValidationError(f'{field}的{operator}操作需要单个有效值')
+        return f'{column} {operator} %s', [normalize_filter_value(field, value)]
+
+    if not isinstance(value, list):
+        raise BusinessValidationError(f'{field}的{operator}操作需要数组值')
+    if operator == 'between':
+        if len(value) != 2:
+            raise BusinessValidationError(f'{field}的between操作必须传两个值')
+        normalized_values = [normalize_filter_value(field, item) for item in value]
+        if normalized_values[0] > normalized_values[1]:
+            raise BusinessValidationError(f'{field}的between起始值不能大于结束值')
+        return f'{column} BETWEEN %s AND %s', normalized_values
+    if operator in ('in', 'not_in'):
+        if not value:
+            raise BusinessValidationError(f'{field}的{operator}操作不能为空')
+        placeholders = ', '.join(['%s'] * len(value))
+        sql_operator = 'IN' if operator == 'in' else 'NOT IN'
+        return f'{column} {sql_operator} ({placeholders})', [
+            normalize_filter_value(field, item) for item in value
+        ]
+    raise BusinessValidationError(f'不支持的筛选操作符: {operator}')
+
+
+def build_filter_scope(
+    params: VideoQueryRequest,
+    table_alias: Optional[str] = None,
+) -> SqlFragment:
+    """生成可复用的筛选范围,供主查询和去重子查询保持一致。"""
+    column = lambda field: qualified_column(field, table_alias)
+    platform_placeholders = ', '.join(['%s'] * len(params.platforms))
+    clauses = [
+        f'{column("platform")} IN ({platform_placeholders})',
+        f'{column("create_time")} >= %s',
+        f'{column("create_time")} < %s',
+        f"{column('video_url')} <> ''",
+    ]
+    sql_params: List[Any] = [*params.platforms, params.start_time, params.end_time]
+
+    if params.keywords:
+        keyword_clause = f'{column("video_title")} LIKE %s'
+        clauses.append(f"({' OR '.join([keyword_clause] * len(params.keywords))})")
+        sql_params.extend([f'%{keyword}%' for keyword in params.keywords])
+
+    filter_clauses, filter_params = [], []
+    for condition in params.filters:
+        clause, values = compile_filter_condition(condition, table_alias)
+        filter_clauses.append(clause)
+        filter_params.extend(values)
+    if filter_clauses:
+        joiner = ' OR ' if params.filter_match_mode == 1 else ' AND '
+        clauses.append(f"({joiner.join(filter_clauses)})")
+        sql_params.extend(filter_params)
+    return ' AND '.join(clauses), sql_params
+
+
+def build_query(params: VideoQueryRequest) -> SqlFragment:
+    """先按out_video_id取最大id,再按主键回表读取当前页完整数据。"""
+    filter_scope, sql_params = build_filter_scope(params, 'source')
+    columns = ', '.join(f'`cv`.`{column}`' for column in SELECT_COLUMNS)
+    page_size = min(params.limit, settings.API_MAX_LIMIT)
+    having_clause = ''
+    if params.cursor:
+        # 游标必须作用在分组结果上,不能放入WHERE,否则旧重复记录会在下一页再次出现。
+        having_clause = 'HAVING MAX(`source`.`id`) < %s'
+        sql_params.append(params.cursor.id)
+    sql = f'''SELECT {columns}
+              FROM `crawler_video` AS `cv`
+              INNER JOIN (
+                  SELECT MAX(`source`.`id`) AS `selected_id`
+                  FROM `crawler_video` AS `source`
+                  WHERE {filter_scope}
+                  GROUP BY
+                      CASE WHEN `source`.`out_video_id` = '' THEN `source`.`id` ELSE 0 END,
+                      `source`.`out_video_id`
+                  {having_clause}
+                  ORDER BY `selected_id` DESC
+                  LIMIT %s
+              ) AS `dedup` ON `dedup`.`selected_id` = `cv`.`id`
+              ORDER BY `cv`.`id` DESC'''
+    # 多取1条判断下一页,避免为每次请求额外执行COUNT(*)。
+    sql_params.append(page_size + 1)
+    return sql, sql_params
+
+
+# ==================== 接口处理与返回 ====================
+
+def api_response(data=None, code: int = 0, msg: str = '') -> web.Response:
+    response = web.json_response(
+        {'code': code, 'msg': msg, 'data': data},
+        status=200 if code == 0 else code,
+        dumps=lambda value: json.dumps(value, ensure_ascii=False, default=str),
+    )
+    response.headers['Cache-Control'] = 'no-store'
+    return response
+
+
+def timestamp_ms(value: datetime) -> int:
+    if value.tzinfo is None:
+        value = value.replace(tzinfo=CHINA_TIMEZONE)
+    else:
+        value = value.astimezone(CHINA_TIMEZONE)
+    return int(value.timestamp() * 1000)
+
+
+async def execute_query(request: web.Request, sql: str, sql_params: List[Any]) -> List[dict]:
+    """在并发限制和查询超时保护下执行数据库查询。"""
+    semaphore = request.app[QUERY_SEMAPHORE_KEY]
+    try:
+        await asyncio.wait_for(semaphore.acquire(), timeout=settings.API_QUEUE_TIMEOUT)
+    except asyncio.TimeoutError as exc:
+        raise ServiceBusyError from exc
+
+    query_started_at = time.perf_counter()
+    try:
+        async with asyncio.timeout(settings.API_QUERY_TIMEOUT):
+            rows = await request.app[MYSQL_KEY].fetch_all(sql, sql_params)
+        request[QUERY_METRICS_KEY] = {
+            'result_count': len(rows),
+            'duration_ms': round((time.perf_counter() - query_started_at) * 1000, 2),
+            'success': True,
+        }
+        return rows
+    except asyncio.TimeoutError:
+        request[QUERY_METRICS_KEY] = {
+            'result_count': None,
+            'duration_ms': round((time.perf_counter() - query_started_at) * 1000, 2),
+            'success': False,
+        }
+        raise
+    except Exception as exc:
+        request[QUERY_METRICS_KEY] = {
+            'result_count': None,
+            'duration_ms': round((time.perf_counter() - query_started_at) * 1000, 2),
+            'success': False,
+        }
+        raise DatabaseQueryError('database query failed') from exc
+    finally:
+        semaphore.release()
+
+
+async def query_videos(request: web.Request) -> web.Response:
+    """查询一页垂直视频,并返回下一页信息。"""
+    # 1. 校验请求参数,并补齐默认三天时间范围。
+    params = VideoQueryRequest.model_validate(await request.json())
+
+    # 2. 生成参数化SQL并查询数据库。
+    sql, sql_params = build_query(params)
+    rows = await execute_query(request, sql, sql_params)
+
+    # 3. 截取当前页,并用最后一条数据的自增主键生成稳定游标。
+    page_size = min(params.limit, settings.API_MAX_LIMIT)
+    has_more = len(rows) > page_size
+    rows = rows[:page_size]
+    next_cursor = None
+    if has_more and rows:
+        last_row = rows[-1]
+        next_cursor = {'id': last_row['id']}
+    response_data = {
+        'data': rows,
+        'count': len(rows),
+        'has_more': has_more,
+        'next_cursor': next_cursor,
+        'start_time': timestamp_ms(params.start_time),
+        'end_time': timestamp_ms(params.end_time),
+    }
+    return api_response(response_data)

+ 29 - 3
config/base.py

@@ -35,6 +35,17 @@ class Settings(BaseSettings):
     DB_POOL_SIZE: int = 20
     DB_POOL_RECYCLE: int = 3600
 
+    # 垂直 spider 查询 API
+    API_HOST: str = "127.0.0.1"
+    API_PORT: int = 8888
+    API_MAX_LIMIT: int = 1000
+    API_MAX_CONCURRENT_REQUESTS: int = 20
+    API_QUERY_TIMEOUT: int = 30
+    API_QUEUE_TIMEOUT: float = 2.0
+    API_LOG_QUEUE_SIZE: int = 1000
+    API_LOG_FLUSH_TIMEOUT: float = 5.0
+    CHUI_ZHI_API_TOKEN: str = ""
+
     # 阿里云RocketMQ配置
     ROCKETMQ_ENDPOINT: str = Field(..., validation_alias="ROCKETMQ_ENDPOINT")
     ROCKETMQ_ACCESS_KEY_ID: str = Field(..., validation_alias="ROCKETMQ_ACCESS_KEY_ID")
@@ -67,20 +78,35 @@ class Settings(BaseSettings):
         return f"redis://:{self.REDIS_PASSWORD}@{self.REDIS_HOST}:{self.REDIS_PORT}/{self.REDIS_DB}"
 
     # Pydantic 2.x 验证器语法
-    @field_validator('DB_PORT', 'REDIS_PORT')
+    @field_validator('DB_PORT', 'REDIS_PORT', 'API_PORT')
     @classmethod
     def validate_port(cls, v: int) -> int:
         if not 1 <= v <= 65535:
             raise ValueError('Port must be between 1 and 65535')
         return v
 
-    @field_validator('DB_POOL_SIZE', 'DB_POOL_RECYCLE', 'REDIS_MAX_CONNECTIONS')
+    @field_validator(
+        'DB_POOL_SIZE',
+        'DB_POOL_RECYCLE',
+        'REDIS_MAX_CONNECTIONS',
+        'API_MAX_LIMIT',
+        'API_MAX_CONCURRENT_REQUESTS',
+        'API_QUERY_TIMEOUT',
+        'API_LOG_QUEUE_SIZE',
+    )
     @classmethod
     def validate_positive_int(cls, v: int) -> int:
         if v <= 0:
             raise ValueError('Value must be positive')
         return v
 
+    @field_validator('API_QUEUE_TIMEOUT', 'API_LOG_FLUSH_TIMEOUT')
+    @classmethod
+    def validate_positive_float(cls, v: float) -> float:
+        if v <= 0:
+            raise ValueError('Value must be positive')
+        return v
+
     @field_validator('ROCKETMQ_WAIT_SECONDS')
     @classmethod
     def validate_rocketmq_wait_seconds(cls, v: int) -> int:
@@ -103,4 +129,4 @@ class Settings(BaseSettings):
         return v
 
 
-settings = Settings()
+settings = Settings()

+ 24 - 33
core/base/async_mysql_client.py

@@ -6,35 +6,13 @@ AsyncMySQLClient (基于 asyncmy) + 项目统一日志正式上线版
 - 高并发协程安全
 """
 import traceback
-from typing import List, Dict, Any, Optional, Tuple
+from typing import List, Dict, Any, Optional
 import asyncio
 import asyncmy
 
 from core.utils.log.logger_manager import LoggerManager
 
 class AsyncMySQLClient:
-    # 类变量:同配置单例管理连接池实例
-    _instances: Dict[Tuple, "AsyncMySQLClient"] = {}
-
-    def __new__(cls,
-                host: str,
-                port: int,
-                user: str,
-                password: str,
-                db: str,
-                charset: str,
-                minsize: int = 1,
-                maxsize: int = 5,
-                **kwargs):
-        """
-        单例模式:同配置共享同一连接池实例
-        """
-        key = (host, port, user, db)
-        if key not in cls._instances:
-            instance = super().__new__(cls)
-            cls._instances[key] = instance
-        return cls._instances[key]
-
     def __init__(self,
                  host: str,
                  port: int,
@@ -44,6 +22,7 @@ class AsyncMySQLClient:
                  charset: str,
                  minsize: int = 1,
                  maxsize: int = 5,
+                 pool_recycle: int = 3600,
                  logger: Optional[LoggerManager] = None,
                  aliyun_logr:Optional[LoggerManager] = None,):
         """
@@ -61,6 +40,7 @@ class AsyncMySQLClient:
         }
         self._minsize = minsize
         self._maxsize = maxsize
+        self._pool_recycle = pool_recycle
         self._pool: Optional[asyncmy.Pool] = None
         self._lock = asyncio.Lock()  # 防止并发初始化
 
@@ -90,11 +70,14 @@ class AsyncMySQLClient:
                     **self._db_settings,
                     minsize=self._minsize,
                     maxsize=self._maxsize,
+                    pool_recycle=self._pool_recycle,
                 )
             except Exception as e:
                 msg = f"[AsyncMySQLClient] 连接池初始化失败: {e} \n {traceback.format_exc()}"
-                self.logger.error(msg)
-                self.aliyun_logger.logging(code="9001", message=msg, data={f"error: f{traceback.format_exc()} \n {e}"})
+                if self.logger:
+                    self.logger.error(msg)
+                if self.aliyun_logger:
+                    self.aliyun_logger.logging(code="9001", message=msg)
                 raise
 
     async def close(self):
@@ -121,8 +104,10 @@ class AsyncMySQLClient:
                     return result
         except Exception as e:
             msg = f"[AsyncMySQLClient] fetch_all 执行失败: {e} | SQL: {sql}"
-            self.logger.error(msg)
-            self.aliyun_logger.logging(code="9002", message=msg)
+            if self.logger:
+                self.logger.error(msg)
+            if self.aliyun_logger:
+                self.aliyun_logger.logging(code="9002", message=msg)
             raise
 
     async def fetch_one(self, sql: str, params: Optional[List[Any]] = None) -> Optional[Dict[str, Any]]:
@@ -142,8 +127,10 @@ class AsyncMySQLClient:
                     return result
         except Exception as e:
             msg = f"[AsyncMySQLClient] fetch_one 执行失败: {e} | SQL: {sql}"
-            self.logger.error(msg)
-            self.aliyun_logger.logging(code="9003", message=msg)
+            if self.logger:
+                self.logger.error(msg)
+            if self.aliyun_logger:
+                self.aliyun_logger.logging(code="9003", message=msg)
             raise
 
     async def execute(self, sql: str, params: Optional[List[Any]] = None) -> int:
@@ -158,8 +145,10 @@ class AsyncMySQLClient:
                     return cur.rowcount
         except Exception as e:
             msg = f"[AsyncMySQLClient] execute 执行失败: {e} | SQL: {sql}"
-            self.logger.error(msg)
-            self.aliyun_logger.logging(code="9004", message=msg)
+            if self.logger:
+                self.logger.error(msg)
+            if self.aliyun_logger:
+                self.aliyun_logger.logging(code="9004", message=msg)
             raise
 
     async def executemany(self, sql: str, params_list: List[List[Any]]) -> int:
@@ -174,6 +163,8 @@ class AsyncMySQLClient:
                     return cur.rowcount
         except Exception as e:
             msg = f"[AsyncMySQLClient] executemany 执行失败: {e} | SQL: {sql}"
-            self.logger.error(msg)
-            self.aliyun_logger.logging(code="9005", message=msg)
+            if self.logger:
+                self.logger.error(msg)
+            if self.aliyun_logger:
+                self.aliyun_logger.logging(code="9005", message=msg)
             raise

+ 52 - 32
core/utils/log/aliyun_log.py

@@ -12,6 +12,21 @@ proxies = {"http": None, "https": None}
 from core.utils.trace_utils import get_current_trace_id  # 导入工具函数
 
 
+API_INDEXED_FIELDS = (
+    "request_id",
+    "path",
+    "method",
+    "status_code",
+    "success",
+    "request_duration_ms",
+    "failure_stage",
+    "error_type",
+    "query_result_count",
+    "query_duration_ms",
+    "query_success",
+)
+
+
 class AliyunLogger(object):
     """
     阿里云日志方法
@@ -26,18 +41,21 @@ class AliyunLogger(object):
     def logging(
             self, code, message, data=None, trace_id=None, account=None
     ):
+        """兼容原有单条日志调用。"""
+        self.logging_batch([{
+            "code": code,
+            "message": message,
+            "data": data,
+            "trace_id": trace_id,
+            "account": account,
+        }])
+
+    def logging_batch(self, events):
         """
-        写入阿里云日志
+        单次请求批量写入阿里云日志,降低高并发场景下的网络开销。
         测试库: https://sls.console.aliyun.com/lognext/project/crawler-log-dev/logsearch/crawler-log-dev
         正式库: https://sls.console.aliyun.com/lognext/project/crawler-log-prod/logsearch/crawler-log-prod
         """
-        # 设置阿里云日志服务的访问信息
-        # if message is None:
-        #     message = LOG_CODES.get(code, "未知错误")
-        # if data is None:
-        #     data = {}
-        # 优先使用传入的 trace_id,否则从上下文获取
-        current_trace_id = get_current_trace_id() or ""
         accessKeyId = settings.ALIYUN_ACCESS_KEY_ID
         accessKey = settings.ALIYUN_ACCESS_KEY_SECRET
         if self.env == "dev":
@@ -49,33 +67,35 @@ class AliyunLogger(object):
             logstore = "crawler-fetch"
             endpoint = "cn-hangzhou.log.aliyuncs.com"
 
-        # 创建 LogClient 实例
         client = LogClient(endpoint, accessKeyId, accessKey)
         log_group = []
-        log_item = LogItem()
-
-        """
-        生成日志消息体格式,例如
-        crawler:xigua
-        message:不满足抓取规则 
-        mode:search
-        timestamp:1686656143
-        """
-        message = message.replace("\r", " ").replace("\n", " ")
-        contents = [
-            (f"TraceId", str(current_trace_id)),
-            (f"code", str(code)),
-            (f"platform", str(self.platform)),
-            (f"mode", str(self.mode)),
-            (f"message", str(message)),
-            (f"data", json.dumps(data, ensure_ascii=False) if data else ""),
-            (f"account", str(account)),
-            ("timestamp", str(int(time.time()))),
-        ]
+        for event in events:
+            message = str(event.get("message") or "").replace("\r", " ").replace("\n", " ")
+            trace_id = event.get("trace_id") or get_current_trace_id() or ""
+            data = event.get("data") or {}
+            log_item = LogItem()
+            contents = [
+                ("TraceId", str(trace_id)),
+                ("code", str(event.get("code", ""))),
+                ("platform", str(self.platform)),
+                ("mode", str(self.mode)),
+                ("message", message),
+                ("data", json.dumps(data, ensure_ascii=False, default=str) if data else ""),
+                ("account", str(event.get("account"))),
+                ("timestamp", str(int(time.time()))),
+            ]
+            # API核心指标额外提升为SLS顶层字段,便于直接统计请求量和成功率。
+            if isinstance(data, dict):
+                contents.extend(
+                    (field, str(data[field]))
+                    for field in API_INDEXED_FIELDS
+                    if data.get(field) is not None
+                )
+            log_item.set_contents(contents)
+            log_group.append(log_item)
 
-        log_item.set_contents(contents)
-        log_group.append(log_item)
-        # 写入日志
+        if not log_group:
+            return
         request = PutLogsRequest(
             project=project,
             logstore=logstore,

+ 94 - 21
deploy.sh

@@ -2,13 +2,17 @@
 set -e  # 出错时终止脚本
 
 # 配置信息
-APP_DIR="/root/AutoScraperX"
+APP_DIR="${APP_DIR:-/root/AutoScraperX}"
 LOG_FILE="${APP_DIR}/logs/autoscraperx_deploy.log"  # 部署日志路径(统一到项目logs目录)
-VENV_DIR="${APP_DIR}/venv"  # 虚拟环境路径(与项目中实际venv目录对应)
+VENV_DIR="${VENV_DIR:-${APP_DIR}/venv}"  # 可通过环境变量覆盖
 PYTHON="${VENV_DIR}/bin/python"  # 虚拟环境Python解释器
 PIP="${VENV_DIR}/bin/pip"        # 虚拟环境pip工具
 REQUIREMENTS="${APP_DIR}/requirements.txt"
-STARTUP_LOG="${APP_DIR}/logs/main_startup.log"  # 服务启动日志
+RUNTIME_DIR="${APP_DIR}/run"
+MAIN_STARTUP_LOG="${APP_DIR}/logs/main_startup.log"
+API_STARTUP_LOG="${APP_DIR}/logs/api_startup.log"
+MAIN_PID_FILE="${RUNTIME_DIR}/main.pid"
+API_PID_FILE="${RUNTIME_DIR}/api.pid"
 
 # 日志函数
 log() {
@@ -39,12 +43,65 @@ check_permission() {
     fi
 }
 
+# 优先根据PID文件停止进程,并兼容旧版deploy.sh没有PID文件的进程。
+stop_service() {
+    local service_name="$1"
+    local pid_file="$2"
+    local process_pattern="$3"
+    local pid=""
+
+    if [ -f "$pid_file" ]; then
+        pid="$(tr -cd '0-9' < "$pid_file")"
+        if [ -n "$pid" ] && kill -0 "$pid" 2>/dev/null; then
+            if ps -p "$pid" -o command= | grep -F "$process_pattern" > /dev/null; then
+                log "停止${service_name}: pid=${pid}"
+                kill "$pid" || true
+                for _ in {1..10}; do
+                    if ! kill -0 "$pid" 2>/dev/null; then
+                        break
+                    fi
+                    sleep 1
+                done
+            fi
+        fi
+        rm -f "$pid_file"
+    fi
+
+    # 首次升级时清理旧部署脚本直接启动、但没有记录PID的进程。
+    pkill -f "$process_pattern" || true
+}
+
+show_startup_error() {
+    local service_name="$1"
+    local startup_log="$2"
+    log "${service_name}启动失败,启动日志最后20行:"
+    tail -n 20 "$startup_log" >> "$LOG_FILE" 2>/dev/null || true
+}
+
+wait_api_ready() {
+    local api_pid="$1"
+    for _ in {1..15}; do
+        if ! kill -0 "$api_pid" 2>/dev/null; then
+            return 1
+        fi
+        if "$PYTHON" -c "import urllib.request; urllib.request.urlopen('http://127.0.0.1:${API_PORT}/ready', timeout=3).read()" > /dev/null 2>&1; then
+            return 0
+        fi
+        sleep 1
+    done
+    return 1
+}
+
 # 主函数
 main() {
+    # 先创建基础目录,确保后续所有步骤都能写部署日志。
+    mkdir -p "${APP_DIR}/logs" "$RUNTIME_DIR" || exit 1
     # 前置检查:确保项目目录存在
     ensure_dir "$APP_DIR"
     ensure_dir "${APP_DIR}/logs"  # 确保日志目录存在
+    ensure_dir "$RUNTIME_DIR"
     check_permission "${APP_DIR}/run.sh"  # 修复run.sh权限
+    check_permission "${APP_DIR}/run_api.sh"
 
     log "===== 开始部署 AutoScraperX ====="
     cd "$APP_DIR" || handle_error "应用目录不存在: $APP_DIR"
@@ -69,28 +126,44 @@ main() {
     "$PIP" install --upgrade pip || handle_error "更新pip失败"
     "$PIP" install -r "$REQUIREMENTS" || handle_error "安装依赖失败"
 
-    # 停止现有服务(精准匹配虚拟环境Python)
-    log "停止现有服务..."
-    pkill -f "${PYTHON} main.py" || true  # 忽略"无进程可杀"的错误
-    sleep 2  # 等待进程终止
+    # 在停止旧服务前验证全部配置,避免配置错误扩大停机时间。
+    log "检查API配置和Token..."
+    "$PYTHON" -c "from config import settings; assert len(settings.CHUI_ZHI_API_TOKEN) >= 32, 'CHUI_ZHI_API_TOKEN至少需要32个字符'; print(f'API监听地址: {settings.API_HOST}:{settings.API_PORT}')" >> "$LOG_FILE" \
+        || handle_error "API配置检查失败,请检查${APP_DIR}/.env"
+    API_PORT="$("$PYTHON" -c "from config import settings; print(settings.API_PORT)")" \
+        || handle_error "无法读取API_PORT"
 
-    # 启动服务(输出日志到文件,便于排查)
-    log "启动新服务...(日志: ${STARTUP_LOG})"
-    nohup "${PYTHON}" main.py > "${STARTUP_LOG}" 2>&1 &
-    sleep 3  # 延长等待时间,确保服务有足够时间启动
+    log "停止现有API和主服务..."
+    stop_service "垂直视频API" "$API_PID_FILE" "${PYTHON} -m api.app"
+    stop_service "主爬虫服务" "$MAIN_PID_FILE" "${PYTHON} main.py"
 
-    # 检查服务状态(更可靠的检测方式)
-    if ps aux | grep -v grep | grep -E "${PYTHON}.*main\.py" > /dev/null; then
-        log "服务已成功启动!"
-    else
-        # 启动失败时,输出启动日志片段帮助排查
-        log "服务启动失败!启动日志最后10行:"
-        tail -n 10 "${STARTUP_LOG}" >> "$LOG_FILE"  # 将启动日志尾部写入部署日志
-        handle_error "服务启动失败,详情见启动日志: ${STARTUP_LOG}"
+    # API先启动并通过数据库就绪检查,主爬虫随后启动。
+    log "启动垂直视频API...(日志: ${API_STARTUP_LOG})"
+    nohup "${PYTHON}" -m api.app > "$API_STARTUP_LOG" 2>&1 &
+    API_PID=$!
+    echo "$API_PID" > "$API_PID_FILE"
+    if ! wait_api_ready "$API_PID"; then
+        show_startup_error "垂直视频API" "$API_STARTUP_LOG"
+        stop_service "垂直视频API" "$API_PID_FILE" "${PYTHON} -m api.app"
+        handle_error "垂直视频API启动或数据库就绪检查失败"
+    fi
+    log "垂直视频API已就绪: pid=${API_PID}"
+
+    log "启动主爬虫服务...(日志: ${MAIN_STARTUP_LOG})"
+    nohup "${PYTHON}" main.py > "$MAIN_STARTUP_LOG" 2>&1 &
+    MAIN_PID=$!
+    echo "$MAIN_PID" > "$MAIN_PID_FILE"
+    sleep 3
+    if ! kill -0 "$MAIN_PID" 2>/dev/null; then
+        show_startup_error "主爬虫服务" "$MAIN_STARTUP_LOG"
+        stop_service "主爬虫服务" "$MAIN_PID_FILE" "${PYTHON} main.py"
+        stop_service "垂直视频API" "$API_PID_FILE" "${PYTHON} -m api.app"
+        handle_error "主爬虫服务启动失败"
     fi
+    log "主爬虫服务已启动: pid=${MAIN_PID}"
 
-    log "===== 部署完成 ====="
+    log "===== 部署完成:主爬虫和垂直视频API均已启动 ====="
 }
 
 # 执行主函数
-main
+main

+ 78 - 0
deploy/nginx/chui_zhi_api.conf.example

@@ -0,0 +1,78 @@
+# 此文件应放在 Nginx http {} 上下文中,例如 /etc/nginx/conf.d/chui_zhi_api.conf。
+# 部署前替换 YOUR_API_DOMAIN 和证书路径。
+
+limit_req_zone $binary_remote_addr zone=chui_zhi_api_rate:10m rate=100r/s;
+limit_conn_zone $binary_remote_addr zone=chui_zhi_api_conn:10m;
+
+upstream chui_zhi_api_backend {
+    server 127.0.0.1:8888;
+    keepalive 64;
+}
+
+server {
+    listen 80;
+    listen [::]:80;
+    server_name YOUR_API_DOMAIN;
+
+    return 301 https://$host$request_uri;
+}
+
+server {
+    listen 443 ssl http2;
+    listen [::]:443 ssl http2;
+    server_name YOUR_API_DOMAIN;
+
+    ssl_certificate /etc/letsencrypt/live/YOUR_API_DOMAIN/fullchain.pem;
+    ssl_certificate_key /etc/letsencrypt/live/YOUR_API_DOMAIN/privkey.pem;
+    ssl_protocols TLSv1.2 TLSv1.3;
+    ssl_session_cache shared:SSL:10m;
+    ssl_session_timeout 1d;
+
+    client_max_body_size 1m;
+    gzip on;
+    gzip_min_length 1k;
+    gzip_types application/json;
+    limit_req zone=chui_zhi_api_rate burst=200 nodelay;
+    limit_conn chui_zhi_api_conn 50;
+
+    add_header X-Content-Type-Options nosniff always;
+    add_header Referrer-Policy no-referrer always;
+
+    # 如果新加坡服务器有固定出口 IP,建议取消下面两行注释并替换 IP。
+    # allow 203.0.113.10;
+    # deny all;
+
+    location = /api/v1/crawler/videos/query {
+        limit_except POST {
+            deny all;
+        }
+
+        proxy_pass http://chui_zhi_api_backend;
+        proxy_http_version 1.1;
+        proxy_set_header Connection "";
+        proxy_set_header Host $host;
+        proxy_set_header X-Real-IP $remote_addr;
+        proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
+        proxy_set_header X-Forwarded-Proto $scheme;
+        proxy_set_header X-API-Key $http_x_api_key;
+        proxy_set_header X-Request-ID $http_x_request_id;
+
+        proxy_connect_timeout 5s;
+        proxy_send_timeout 30s;
+        proxy_read_timeout 40s;
+        proxy_buffering on;
+    }
+
+    # 运维探针只允许本机访问,不作为公网业务接口暴露。
+    location ~ ^/(health|ready)$ {
+        allow 127.0.0.1;
+        allow ::1;
+        deny all;
+        proxy_pass http://chui_zhi_api_backend;
+        proxy_set_header Host $host;
+    }
+
+    location / {
+        return 404;
+    }
+}

+ 19 - 0
deploy/systemd/autoscraperx-api.service.example

@@ -0,0 +1,19 @@
+[Unit]
+Description=AutoScraperX Chui Zhi API
+After=network-online.target
+Wants=network-online.target
+
+[Service]
+Type=simple
+User=root
+WorkingDirectory=/root/AutoScraperX
+EnvironmentFile=/root/AutoScraperX/.env
+ExecStart=/root/AutoScraperX/venv/bin/python -m api.app
+Restart=always
+RestartSec=3
+TimeoutStopSec=30
+KillSignal=SIGTERM
+LimitNOFILE=65535
+
+[Install]
+WantedBy=multi-user.target

+ 771 - 0
docs/chui_zhi_video_api.md

@@ -0,0 +1,771 @@
+# 垂直视频查询 API 文档
+
+本文档说明垂直视频查询 API 的调用协议、参数规则、数据库查询逻辑、去重与分页语义、错误处理、日志、部署和性能注意事项。
+
+## 1. 接口概览
+
+| 项目 | 内容 |
+|---|---|
+| 业务接口 | `POST /api/v1/crawler/videos/query` |
+| 请求格式 | `application/json` |
+| 鉴权 | 请求头 `X-API-Key` |
+| 时间标准 | 东八区(UTC+8) |
+| 默认时间范围 | 最近 3 天 |
+| 默认页大小 | 500 条 |
+| 最大页大小 | 由 `API_MAX_LIMIT` 控制,默认 1000 条 |
+| 分页方式 | 基于去重结果最大 `id` 的游标分页 |
+| 去重方式 | 相同 `out_video_id` 保留最大 `id` 对应的记录 |
+| 数据库操作 | 只读,不修改数据库结构和数据 |
+
+此外提供两个运维探针:
+
+| 路径 | 说明 |
+|---|---|
+| `GET /health` | 进程存活检查,不访问数据库 |
+| `GET /ready` | 就绪检查,执行 `SELECT 1 AS ok` 验证数据库连接 |
+
+探针不要求 `X-API-Key`,但 Nginx 示例默认只允许本机访问,不应暴露为公网业务接口。
+
+## 2. 鉴权与请求追踪
+
+### 2.1 API Token
+
+业务请求必须携带:
+
+```http
+X-API-Key: <CHUI_ZHI_API_TOKEN>
+```
+
+Token 只能通过环境变量或服务器密钥管理系统注入,不应写入源码:
+
+```bash
+export CHUI_ZHI_API_TOKEN="$(openssl rand -hex 32)"
+```
+
+Token 为空时,API 拒绝启动。Token 错误时返回 HTTP 401。
+
+### 2.2 Request ID
+
+调用方可以传入:
+
+```http
+X-Request-ID: scheduler-task-123-page-1
+```
+
+允许字符为字母、数字、点、下划线和短横线,最长 128 个字符。未传或格式不正确时,API 自动生成 UUID。
+
+响应头会返回同一个 `X-Request-ID`。本地日志和阿里云 SLS 也记录该值,可用于串联调度日志、API 日志和数据库异常。
+
+## 3. 请求参数
+
+请求体为 JSON 对象。未定义字段会被拒绝,避免拼写错误被静默忽略。
+
+| 参数 | 类型 | 必填 | 默认值 | 限制与说明 |
+|---|---|---:|---|---|
+| `platforms` | `string[]` | 否 | `xiaoniangao`、`xiaoniangaotuijianliu` | 只支持这两个平台;1~20 个;单项最长 50;自动去除首尾空格和重复值 |
+| `start_time` | 时间或时间戳 | 否 | 见时间规则 | 查询 `create_time >= start_time` |
+| `end_time` | 时间或时间戳 | 否 | 见时间规则 | 查询 `create_time < end_time`,结束时间不包含 |
+| `keywords` | `string[]` | 否 | `[]` | 最多 20 个;单项最长 200;按标题模糊匹配 |
+| `filter_match_mode` | `integer` | 否 | `2` | `1` 表示过滤条件 OR;`2` 表示过滤条件 AND |
+| `filters` | `object[]` | 否 | `[]` | 最多 50 个结构化过滤条件 |
+| `limit` | `integer` | 否 | `500` | 范围为 1~`API_MAX_LIMIT`;超出范围直接返回参数错误 |
+| `cursor` | `object` | 否 | `null` | 第一页不传;下一页原样传回响应的 `next_cursor` |
+
+### 3.1 时间参数规则
+
+`start_time` 和 `end_time` 支持:
+
+- 13 位毫秒时间戳,例如 `1786982400000`
+- 10 位秒时间戳,例如 `1786982400`
+- 仅包含数字的时间戳字符串
+- ISO 日期时间字符串,例如 `2026-08-18 00:00:00` 或 `2026-08-18T00:00:00`
+
+时间戳统一按东八区转换为数据库查询时间。
+
+缺省规则:
+
+| 传参情况 | 实际范围 |
+|---|---|
+| 开始和结束均未传 | 当前时间往前 3 天至当前时间 |
+| 只传 `end_time` | `end_time` 往前 3 天至 `end_time` |
+| 只传 `start_time` | `start_time` 至当前时间 |
+| 两者都传 | 使用调用方指定范围 |
+
+`start_time` 必须早于 `end_time`。查询边界为左闭右开:
+
+```sql
+create_time >= start_time AND create_time < end_time
+```
+
+第一页响应会返回实际使用的 `start_time` 和 `end_time`(毫秒时间戳)。调度端会固定该范围用于后续页面,避免默认“当前时间”随分页变化。
+
+### 3.2 关键词规则
+
+关键词只匹配 `video_title`,每个关键词使用:
+
+```sql
+video_title LIKE '%关键词%'
+```
+
+多个关键词之间固定使用 OR。例如:
+
+```json
+{"keywords": ["早", "养生"]}
+```
+
+等价于:
+
+```sql
+(video_title LIKE '%早%' OR video_title LIKE '%养生%')
+```
+
+关键词条件与平台、创建时间、视频地址和 `filters` 之间使用 AND。
+
+`%关键词%` 无法有效利用普通 B-Tree 索引,短关键词或高频词可能扫描较多数据。
+
+### 3.3 结构化过滤条件
+
+每个过滤条件格式为:
+
+```json
+{
+  "field": "like_cnt",
+  "operator": ">",
+  "value": 3
+}
+```
+
+每个 `filters` 元素严格只允许以下三个 key:
+
+| key | 必填 | 说明 |
+|---|---:|---|
+| `field` | 是 | 白名单过滤字段 |
+| `operator` | 是 | 白名单比较操作符 |
+| `value` | 是 | 与操作符匹配的单值或数组 |
+
+缺少任意 key 会返回“不能为空”;增加 `column`、`sql`、`name` 等任何未知 key 会返回 400。例如:
+
+```text
+参数校验失败: 不支持的参数: filters.0.column
+参数校验失败: filters.0.operator不能为空
+```
+
+允许字段:
+
+| 字段 | 含义 | 值类型 |
+|---|---|---|
+| `play_cnt` | 播放数 | 数字 |
+| `like_cnt` | 点赞数 | 数字 |
+| `share_cnt` | 分享数 | 数字 |
+| `collection_cnt` | 收藏数 | 数字 |
+| `comment_cnt` | 评论数 | 数字 |
+| `duration` | 视频时长,单位秒 | 数字 |
+| `publish_time` | 视频发布时间 | 时间或时间戳 |
+
+允许操作符:
+
+| 操作符 | `value` 要求 | 示例 |
+|---|---|---|
+| `>` | 单值 | `{"field":"like_cnt","operator":">","value":3}` |
+| `>=` | 单值 | `{"field":"play_cnt","operator":">=","value":1000}` |
+| `=` | 单值 | `{"field":"duration","operator":"=","value":60}` |
+| `<` | 单值 | `{"field":"duration","operator":"<","value":60}` |
+| `<=` | 单值 | `{"field":"comment_cnt","operator":"<=","value":100}` |
+| `between` | 恰好两个值的数组,包含边界 | `{"field":"duration","operator":"between","value":[60,300]}` |
+| `in` | 非空数组 | `{"field":"play_cnt","operator":"in","value":[100,500]}` |
+| `not_in` | 非空数组 | `{"field":"duration","operator":"not_in","value":[0,1]}` |
+
+数字字段拒绝布尔值、`NaN`、正负无穷、SQL 片段和不能转换为有限数字的字符串。字段名和操作符均为白名单,所有值都通过 SQL 参数绑定,不直接拼接到 SQL。
+
+`between` 的起始值不能大于结束值;`in`、`not_in` 以及其他数组型筛选值最多允许 100 项。过滤条件对象中的未知字段也会被拒绝。
+
+`publish_time` 当前数据库类型为 `varchar(100)`,API 会把时间值格式化成 `YYYY-MM-DD HH:MM:SS` 后比较。数据库中的历史值需要保持相同且可排序的格式,否则筛选结果可能不准确。本接口不修改数据库字段类型。
+
+### 3.4 过滤条件组合
+
+```json
+"filter_match_mode": 2
+```
+
+- `1`:所有 `filters` 使用 OR
+- `2`:所有 `filters` 使用 AND,默认值
+
+例如:
+
+```json
+{
+  "filter_match_mode": 2,
+  "filters": [
+    {"field": "like_cnt", "operator": ">", "value": 3},
+    {"field": "duration", "operator": "between", "value": [60, 300]}
+  ]
+}
+```
+
+对应:
+
+```sql
+like_cnt > 3 AND duration BETWEEN 60 AND 300
+```
+
+`filter_match_mode` 只控制 `filters` 数组内部,不改变关键词之间的 OR 规则。
+
+### 3.5 游标参数
+
+第一页不传 `cursor`:
+
+```json
+{
+  "keywords": ["早"],
+  "limit": 500
+}
+```
+
+如果响应 `has_more=true`,下一页原样传入:
+
+```json
+{
+  "keywords": ["早"],
+  "limit": 500,
+  "cursor": {
+    "id": 6815000
+  }
+}
+```
+
+`cursor.id` 必须是大于 0 的整数。不要自行修改游标,也不要在分页过程中修改其他筛选参数。
+
+## 4. 完整请求示例
+
+```json
+{
+  "platforms": ["xiaoniangao", "xiaoniangaotuijianliu"],
+  "start_time": 1786896000000,
+  "end_time": 1786982400000,
+  "keywords": ["早", "养生"],
+  "filter_match_mode": 2,
+  "filters": [
+    {"field": "play_cnt", "operator": ">=", "value": 1000},
+    {"field": "like_cnt", "operator": ">", "value": 3},
+    {"field": "duration", "operator": "between", "value": [60, 300]},
+    {"field": "publish_time", "operator": ">=", "value": 1785772800000}
+  ],
+  "limit": 500
+}
+```
+
+### 4.1 参数校验失败示例
+
+不支持的平台:
+
+```json
+{
+  "platforms": ["douyin"]
+}
+```
+
+返回:
+
+```json
+{
+  "code": 400,
+  "msg": "参数校验失败: platforms: 不支持的平台: douyin",
+  "data": [
+    {
+      "type": "value_error",
+      "loc": ["platforms"]
+    }
+  ]
+}
+```
+
+不支持的顶层参数:
+
+```json
+{
+  "page_index": 2
+}
+```
+
+返回消息包含:
+
+```text
+参数校验失败: 不支持的参数: page_index
+```
+
+不支持的过滤字段或操作符会指出完整路径,例如:
+
+```text
+参数校验失败: filters.0.field不支持值 unknown_field
+参数校验失败: filters.0.operator不支持值 contains
+```
+
+## 5. 数据库查询逻辑
+
+查询依次应用以下规则:
+
+1. `platform IN (...)`
+2. `create_time >= start_time AND create_time < end_time`
+3. `video_url <> ''`
+4. 有关键词时,执行标题模糊匹配
+5. 有结构化过滤条件时,按 `filter_match_mode` 组合
+6. 在完整筛选范围内按 `out_video_id` 去重
+7. 在去重结果上应用 `id` 游标
+8. 多取 1 条判断是否还有下一页
+9. 通过主键回表读取完整视频字段
+
+核心 SQL 形态:
+
+```sql
+SELECT cv.<视频字段>
+FROM crawler_video AS cv
+INNER JOIN (
+    SELECT MAX(source.id) AS selected_id
+    FROM crawler_video AS source
+    WHERE <平台、时间、视频地址、关键词、过滤条件>
+    GROUP BY
+        CASE
+            WHEN source.out_video_id = '' THEN source.id
+            ELSE 0
+        END,
+        source.out_video_id
+    HAVING MAX(source.id) < :cursor_id
+    ORDER BY selected_id DESC
+    LIMIT :page_size_plus_one
+) AS dedup ON dedup.selected_id = cv.id
+ORDER BY cv.id DESC;
+```
+
+第一页没有 `HAVING` 条件。
+
+### 5.1 去重语义
+
+- 相同非空 `out_video_id` 只保留 `id` 最大的一条,即最后入库记录。
+- `out_video_id=''` 的记录按各自 `id` 分组,因此不会全部合并成一条。
+- 去重发生在分页之前,所以重复记录不占用返回页容量,也不会跨页再次出现。
+- 当前语义是“最大 `id`/最后入库”,不是严格比较 `publish_time`。
+
+### 5.2 分页语义
+
+- 排序方式为 `id DESC`。
+- 下一页通过 `HAVING MAX(id) < cursor.id` 应用于分组后的结果。
+- 不使用 OFFSET,因此翻到后续页面时不会产生越来越大的跳过成本。
+- 固定时间范围后,新插入且 `id` 大于当前游标的数据不会插入到本轮分页中间。
+
+## 6. 响应格式
+
+### 6.1 成功响应
+
+```json
+{
+  "code": 0,
+  "msg": "",
+  "data": {
+    "data": [
+      {
+        "video_id": 10001,
+        "user_id": 20001,
+        "out_user_id": "external-user-id",
+        "platform": "xiaoniangao",
+        "strategy": "recommend",
+        "out_video_id": "external-video-id",
+        "video_title": "早上好",
+        "cover_url": "https://example.com/cover.jpg",
+        "video_url": "https://example.com/video.mp4",
+        "duration": 60,
+        "publish_time": "2026-08-18 08:00:00",
+        "play_cnt": 1000,
+        "like_cnt": 20,
+        "share_cnt": 2,
+        "collection_cnt": 3,
+        "comment_cnt": 5,
+        "width": 1080,
+        "height": 1920,
+        "id": 6815000,
+        "create_time": "2026-08-18 08:10:00"
+      }
+    ],
+    "count": 1,
+    "has_more": true,
+    "next_cursor": {"id": 6815000},
+    "start_time": 1786723200000,
+    "end_time": 1786982400000
+  }
+}
+```
+
+字段说明:
+
+| 字段 | 说明 |
+|---|---|
+| `code` | 业务码;成功固定为 0 |
+| `msg` | 成功时为空字符串 |
+| `data.data` | 当前页视频数组 |
+| `data.count` | 当前页实际返回数量,不包含额外探测行 |
+| `data.has_more` | 是否存在下一页 |
+| `data.next_cursor` | 下一页游标;没有下一页时为 `null` |
+| `data.start_time` | API 实际使用的开始时间,13 位毫秒时间戳 |
+| `data.end_time` | API 实际使用的结束时间,13 位毫秒时间戳 |
+
+响应头包含:
+
+```http
+X-Request-ID: <本次请求ID>
+Cache-Control: no-store
+```
+
+## 7. 错误响应
+
+错误响应统一为:
+
+```json
+{
+  "code": 422,
+  "msg": "like_cnt筛选值必须是数字: invalid",
+  "data": null
+}
+```
+
+| HTTP 状态 | 场景 | 调用方处理建议 |
+|---:|---|---|
+| 400 | JSON 无效、字段类型错误、未知参数、时间范围错误 | 修正请求,不重试 |
+| 401 | Token 缺失或错误 | 修正鉴权,不重试 |
+| 422 | 结构化过滤值不符合业务要求 | 修正参数,不重试 |
+| 500 | 数据库查询异常或内部异常 | 记录 `X-Request-ID`,排查日志;谨慎重试 |
+| 503 | 查询并发已满,或 `/ready` 检查失败 | 指数退避后重试 |
+| 504 | 数据库查询超过 `API_QUERY_TIMEOUT` | 缩小时间范围/关键词或稍后重试 |
+
+调度端当前只自动重试网络异常以及 HTTP 429、502、503、504,最多 3 次;退避约为 1 秒、2 秒、4 秒并附加随机抖动。400、401、422 不会重试。
+
+## 8. 并发、超时与资源保护
+
+默认配置:
+
+| 配置 | 默认值 | 说明 |
+|---|---:|---|
+| `API_HOST` | `127.0.0.1` | 只监听本机,由 Nginx 对外代理 |
+| `API_PORT` | `8888` | API 本地监听端口 |
+| `DB_POOL_SIZE` | `20` | 单进程数据库连接池上限 |
+| `DB_POOL_RECYCLE` | `3600` | 数据库连接回收秒数 |
+| `API_MAX_CONCURRENT_REQUESTS` | `20` | 查询并发配置上限 |
+| `API_QUEUE_TIMEOUT` | `2.0` | 等待查询并发槽位的最长秒数 |
+| `API_QUERY_TIMEOUT` | `30` | 单条查询最长秒数 |
+| `API_MAX_LIMIT` | `1000` | 单页最大返回条数 |
+| `API_LOG_QUEUE_SIZE` | `1000` | SLS 异步日志队列长度 |
+| `API_LOG_FLUSH_TIMEOUT` | `5.0` | 单批 SLS 上报及关闭刷新超时 |
+
+实际查询并发取以下两者最小值:
+
+```text
+min(API_MAX_CONCURRENT_REQUESTS, DB_POOL_SIZE)
+```
+
+这样可以避免大量查询进入数据库连接池内部无界等待。超过并发槽位且 2 秒内未获得执行机会时返回 503。
+
+如果使用多个 API 进程,总数据库连接数约为:
+
+```text
+进程数 × DB_POOL_SIZE
+```
+
+扩容进程前必须确认数据库最大连接数和实际负载。
+
+## 9. 日志与监控
+
+### 9.1 本地日志
+
+本地访问日志记录:
+
+- `request_id`
+- 请求方法和 URL
+- 脱敏后的请求参数
+- HTTP 状态码
+- 请求总耗时
+- 失败时的错误原因和脱敏堆栈
+
+不记录响应视频内容,避免日志量过大。
+
+### 9.2 阿里云 SLS
+
+请求完成后只写入有界内存队列,由后台任务批量上报,业务响应不等待 SLS。
+
+后台策略:
+
+- 每批最多 50 条
+- 上报失败最多重试 3 次
+- 队列满时写本地错误日志并丢弃当前 SLS 事件
+- 服务关闭时尝试刷新队列,超时后取消日志任务并优雅关闭数据库连接池
+
+主要字段:
+
+- `TraceId` / `request_id`
+- `path`、`method`
+- `status_code`、`success`
+- `request_duration_ms`
+- `query_duration_ms`
+- `query_result_count`
+- `query_success`
+- `failure_stage`
+- `error_type`
+- `message`
+- 脱敏后的 `request_params`
+
+其中请求量、状态、耗时和错误相关字段会提升为 SLS 顶层字段,方便直接统计每日请求量、成功率、P95/P99 耗时和错误分布。
+
+## 10. curl 示例
+
+以下示例假设:
+
+```bash
+export CHUI_ZHI_API_URL="https://api.example.com/api/v1/crawler/videos/query"
+export CHUI_ZHI_API_TOKEN="replace-with-real-token"
+```
+
+### 10.1 默认最近 3 天
+
+```bash
+curl --request POST "$CHUI_ZHI_API_URL" \
+  --header "Content-Type: application/json" \
+  --header "X-API-Key: $CHUI_ZHI_API_TOKEN" \
+  --header "X-Request-ID: manual-test-001" \
+  --data '{}'
+```
+
+### 10.2 关键词模糊搜索
+
+```bash
+curl --request POST "$CHUI_ZHI_API_URL" \
+  --header "Content-Type: application/json" \
+  --header "X-API-Key: $CHUI_ZHI_API_TOKEN" \
+  --data '{
+    "keywords": ["早"],
+    "limit": 500
+  }'
+```
+
+### 10.3 指定时间及数值条件
+
+```bash
+curl --request POST "$CHUI_ZHI_API_URL" \
+  --header "Content-Type: application/json" \
+  --header "X-API-Key: $CHUI_ZHI_API_TOKEN" \
+  --data '{
+    "start_time": 1786896000000,
+    "end_time": 1786982400000,
+    "filter_match_mode": 2,
+    "filters": [
+      {"field": "like_cnt", "operator": ">", "value": 3},
+      {"field": "duration", "operator": "between", "value": [60, 300]}
+    ],
+    "limit": 200
+  }'
+```
+
+### 10.4 下一页
+
+```bash
+curl --request POST "$CHUI_ZHI_API_URL" \
+  --header "Content-Type: application/json" \
+  --header "X-API-Key: $CHUI_ZHI_API_TOKEN" \
+  --data '{
+    "keywords": ["早"],
+    "start_time": 1786723200000,
+    "end_time": 1786982400000,
+    "limit": 500,
+    "cursor": {"id": 6815000}
+  }'
+```
+
+后续页必须继续使用第一页响应的时间范围、关键词、过滤条件和 `limit`,只替换 `cursor`。
+
+### 10.5 健康检查
+
+在 API 所在服务器执行:
+
+```bash
+curl --fail http://127.0.0.1:8888/health
+curl --fail http://127.0.0.1:8888/ready
+```
+
+## 11. 启动与公网部署
+
+### 11.1 使用现有 deploy.sh 一起部署(推荐)
+
+现有 `deploy.sh` 已同时管理主爬虫和垂直视频 API。在服务器的 `/root/AutoScraperX/.env` 中增加:
+
+```dotenv
+API_HOST=127.0.0.1
+API_PORT=8888
+CHUI_ZHI_API_TOKEN=至少32个字符的随机Token
+```
+
+生成 Token:
+
+```bash
+openssl rand -hex 32
+```
+
+Token 需要分别配置在 API 服务器和新加坡调度服务中,不能提交到 Git。
+
+执行完整部署:
+
+```bash
+cd /root/AutoScraperX
+bash deploy.sh
+```
+
+部署脚本按以下顺序执行:
+
+1. 拉取 `master` 最新代码。
+2. 创建或更新 `/root/AutoScraperX/venv`。
+3. 安装 `requirements.txt`。
+4. 在停止旧服务前检查全部配置,并要求 Token 长度至少 32。
+5. 根据 PID 文件优雅停止旧 API 和主爬虫;首次升级时兼容清理旧版无 PID 进程。
+6. 启动 `python -m api.app`。
+7. 最多等待 15 秒调用 `/ready`,确认进程和数据库都可用。
+8. API 就绪后启动 `main.py`。
+9. 主进程启动检查失败时,同时停止本次启动的 API,避免留下半部署状态。
+
+部署产物:
+
+| 文件 | 用途 |
+|---|---|
+| `logs/autoscraperx_deploy.log` | 部署过程日志 |
+| `logs/api_startup.log` | API 标准输出和启动异常 |
+| `logs/main_startup.log` | 主爬虫标准输出和启动异常 |
+| `run/api.pid` | API 进程 PID |
+| `run/main.pid` | 主爬虫进程 PID |
+
+部署完成后验证:
+
+```bash
+curl --fail http://127.0.0.1:8888/health
+curl --fail http://127.0.0.1:8888/ready
+cat /root/AutoScraperX/run/api.pid
+cat /root/AutoScraperX/run/main.pid
+tail -n 100 /root/AutoScraperX/logs/autoscraperx_deploy.log
+```
+
+如果服务器目录或虚拟环境路径不同,可以覆盖:
+
+```bash
+APP_DIR=/data/AutoScraperX VENV_DIR=/data/AutoScraperX/venv bash deploy.sh
+```
+
+### 11.2 单独启动 API(仅用于本地调试)
+
+```bash
+cd /root/AutoScraperX
+./run_api.sh prod
+```
+
+`run_api.sh` 优先使用 `venv/bin/python`,不存在时使用 `.venv/bin/python`。
+
+默认监听:
+
+```text
+127.0.0.1:8888
+```
+
+启动窗口会打印所有已注册的路径和处理方法。
+
+### 11.3 Nginx(首次部署配置一次)
+
+示例文件:
+
+```text
+deploy/nginx/chui_zhi_api.conf.example
+```
+
+部署要点:
+
+- API 保持监听 `127.0.0.1:8888`
+- Nginx 对外监听 HTTPS 443
+- 只允许业务路径使用 POST
+- 请求体最大 1 MB
+- 开启 JSON gzip
+- 配置请求速率和单 IP 连接限制
+- 透传 `X-API-Key`、`X-Request-ID`、真实 IP 和协议
+- `/health`、`/ready` 只允许本机访问
+- 如果新加坡调度服务器有固定出口 IP,建议启用 IP 白名单
+
+安装示例:
+
+```bash
+cp deploy/nginx/chui_zhi_api.conf.example /etc/nginx/conf.d/chui_zhi_api.conf
+# 编辑域名、证书路径和可选IP白名单
+nginx -t
+systemctl reload nginx
+```
+
+`deploy.sh` 不会自动覆盖 Nginx 配置、域名和证书。Nginx 首次配置完成后,后续运行 `deploy.sh` 只需重启其后端 API。
+
+### 11.4 systemd(可选替代方案)
+
+示例文件:
+
+```text
+deploy/systemd/autoscraperx-api.service.example
+```
+
+systemd 配置了进程异常自动重启、SIGTERM 优雅退出和文件句柄上限。它是 `deploy.sh + nohup` 的替代方案,同一 API 不要同时由两套方式管理。
+
+## 12. 调度端调用约定
+
+调度端配置:
+
+```bash
+export CHUI_ZHI_API_URL="https://api.example.com/api/v1/crawler/videos/query"
+export CHUI_ZHI_API_TOKEN="same-token-as-api"
+```
+
+调度流程:
+
+1. 把抓取计划转换为 API 的结构化参数。
+2. 请求第一页。
+3. 固定 API 返回的实际 `start_time` 和 `end_time`。
+4. 上传并保存当前页视频。
+5. `has_more=true` 时传递 `next_cursor` 请求下一页。
+6. `has_more=false` 时结束。
+
+调度端复用 `requests.Session`,保持 TCP 连接;每一页处理完成后才请求下一页,避免一次性把全部视频加载到内存。
+
+调度端仍会在入库前检查该计划是否已经存在对应内容,保证任务重试时尽量避免重复上传和保存。
+
+## 13. 当前性能参考
+
+2026-08-18 使用当前实现实测:
+
+- 平台:默认两个平台
+- 时间:最近 3 天
+- 关键词:`早`
+- 页大小:500,SQL 实际多取 1 条
+- 三次纯 SQL 耗时:约 4.03 秒、4.22 秒、4.18 秒
+- 平均纯 SQL:约 4.14 秒
+- 首次数据库连接:约 0.15 秒
+
+该结果只代表当时数据库数据量、网络和负载,不是固定 SLA。短关键词 `%早%` 命中范围较大,通常比长关键词更慢。
+
+建议持续通过 SLS 的 `query_duration_ms` 观察:
+
+- 平均查询耗时
+- P95/P99 查询耗时
+- 超时数量
+- 503 并发拒绝数量
+- 各关键词和时间范围的结果量
+
+如果查询长期接近 30 秒,应优先缩小时间范围或页大小,并结合生产数据执行 `EXPLAIN`。数据库结构优化需单独评估和审批,本接口不会自动执行建索引或字段变更。
+
+## 14. 代码位置
+
+| 功能 | 文件 |
+|---|---|
+| aiohttp 应用、路由、中间件、鉴权、日志和生命周期 | `api/app.py` |
+| 请求模型、筛选、SQL、查询和响应 | `api/chui_zhi/videos.py` |
+| API 和数据库配置 | `config/base.py` |
+| 异步 MySQL 连接池 | `core/base/async_mysql_client.py` |
+| 阿里云 SLS 批量日志 | `core/utils/log/aliyun_log.py` |
+| 启动脚本 | `run_api.sh` |
+| Nginx 示例 | `deploy/nginx/chui_zhi_api.conf.example` |
+| systemd 示例 | `deploy/systemd/autoscraperx-api.service.example` |
+| API 测试 | `test/test_video_query_api.py` |

+ 15 - 0
run_api.sh

@@ -0,0 +1,15 @@
+#!/bin/bash
+set -e
+
+ENV=${1:-prod}
+export ENV
+
+PROJECT_DIR="$(cd "$(dirname "$0")" && pwd)"
+if [ -z "$API_PYTHON_BIN" ]; then
+    if [ -x "${PROJECT_DIR}/venv/bin/python" ]; then
+        API_PYTHON_BIN="${PROJECT_DIR}/venv/bin/python"
+    else
+        API_PYTHON_BIN="${PROJECT_DIR}/.venv/bin/python"
+    fi
+fi
+exec "$API_PYTHON_BIN" -m api.app

+ 478 - 0
test/test_video_query_api.py

@@ -0,0 +1,478 @@
+import asyncio
+from datetime import datetime, timedelta, timezone
+
+import pytest
+from aiohttp.test_utils import TestClient, TestServer
+
+from api.app import ALIYUN_LOGGER_KEY, API_PATH, LOGGER_KEY, create_app
+from api.chui_zhi.videos import (
+    BusinessValidationError,
+    MYSQL_KEY,
+    QUERY_SEMAPHORE_KEY,
+    FilterCondition,
+    VideoQueryRequest,
+    build_query,
+    compile_filter_condition,
+)
+from config import settings
+
+
+def test_build_query_contains_parameterized_filters():
+    request = VideoQueryRequest(
+        platforms=['xiaoniangao', 'xiaoniangaotuijianliu'],
+        start_time=datetime(2026, 8, 4),
+        end_time=datetime(2026, 8, 5),
+        keywords=['养生'],
+        filter_match_mode=2,
+        filters=[
+            FilterCondition(field='like_cnt', operator='>=', value=100),
+            FilterCondition(field='duration', operator='between', value=[60, 300]),
+        ],
+        limit=20,
+    )
+
+    sql, params = build_query(request)
+
+    assert '`platform` IN (%s, %s)' in sql
+    assert '`video_title` LIKE %s' in sql
+    assert '`like_cnt` >= %s' in sql
+    assert '`duration` BETWEEN %s AND %s' in sql
+    assert 'SELECT MAX(`source`.`id`) AS `selected_id`' in sql
+    assert "CASE WHEN `source`.`out_video_id` = '' THEN `source`.`id` ELSE 0 END" in sql
+    assert '`dedup`.`selected_id` = `cv`.`id`' in sql
+    assert 'ORDER BY `cv`.`id` DESC' in sql
+    assert 'OFFSET' not in sql
+    assert params[-1] == 21
+    assert '%养生%' in params
+
+
+def test_deduplication_groups_ids_before_cursor_pagination():
+    request = VideoQueryRequest(
+        platforms=['xiaoniangao'],
+        start_time=datetime(2026, 8, 4),
+        end_time=datetime(2026, 8, 5),
+        filters=[FilterCondition(field='like_cnt', operator='>', value=3)],
+        cursor={'id': 100},
+        limit=50,
+    )
+
+    sql, params = build_query(request)
+
+    # 筛选条件只出现一次;HAVING作用于MAX(id),防止旧重复记录跨页再次出现。
+    assert sql.count('`source`.`like_cnt` > %s') == 1
+    assert 'HAVING MAX(`source`.`id`) < %s' in sql
+    assert params[-2:] == [100, 51]
+
+
+def test_query_limit_above_server_max_is_rejected():
+    with pytest.raises(ValueError):
+        VideoQueryRequest(
+            start_time=datetime(2026, 8, 4),
+            end_time=datetime(2026, 8, 5),
+            limit=settings.API_MAX_LIMIT + 1,
+        )
+
+
+def test_id_cursor_builds_stable_group_pagination_and_fetches_one_extra_row():
+    request = VideoQueryRequest(
+        start_time=datetime(2026, 8, 4),
+        end_time=datetime(2026, 8, 5),
+        limit=100,
+        cursor={'id': 123},
+    )
+
+    sql, params = build_query(request)
+
+    assert 'HAVING MAX(`source`.`id`) < %s' in sql
+    assert params[-2:] == [123, 101]
+    assert 'OFFSET' not in sql
+
+
+def test_millisecond_timestamps_are_converted_to_china_time():
+    china_timezone = timezone(timedelta(hours=8))
+    start_time = datetime(2026, 8, 4, tzinfo=china_timezone)
+    end_time = datetime(2026, 8, 5, tzinfo=china_timezone)
+    request = VideoQueryRequest(
+        start_time=int(start_time.timestamp() * 1000),
+        end_time=int(end_time.timestamp() * 1000),
+    )
+
+    assert request.start_time == datetime(2026, 8, 4)
+    assert request.end_time == datetime(2026, 8, 5)
+
+
+def test_missing_time_defaults_to_latest_three_days():
+    before = datetime.now(timezone(timedelta(hours=8))).replace(tzinfo=None)
+    request = VideoQueryRequest()
+    after = datetime.now(timezone(timedelta(hours=8))).replace(tzinfo=None)
+
+    assert before <= request.end_time <= after
+    assert request.end_time - request.start_time == timedelta(days=3)
+
+    sql, sql_params = build_query(request)
+    assert '`create_time` >= %s' in sql
+    assert '`create_time` < %s' in sql
+    assert request.start_time in sql_params
+    assert request.end_time in sql_params
+
+
+def test_rejects_unsupported_filter_field():
+    with pytest.raises(ValueError):
+        FilterCondition(field='unknown_field', operator='>', value=1)
+
+
+def test_rejects_unsupported_platform_and_unknown_filter_parameter():
+    with pytest.raises(ValueError, match='不支持的平台'):
+        VideoQueryRequest(platforms=['douyin'])
+    with pytest.raises(ValueError):
+        FilterCondition(field='like_cnt', operator='>', value=1, unknown='value')
+
+
+@pytest.mark.parametrize(
+    'field',
+    [
+        'like_cnt',
+        'collection_cnt',
+        'comment_cnt',
+        'share_cnt',
+        'play_cnt',
+        'duration',
+    ],
+)
+def test_supported_numeric_filter_mapping(field):
+    condition = FilterCondition(field=field, operator='>', value=10)
+    sql, params = compile_filter_condition(condition)
+
+    assert f'`{field}` > %s' == sql
+    assert params == [10]
+
+
+def test_empty_filters_are_not_added_to_query():
+    request = VideoQueryRequest(filters=[])
+    sql, _ = build_query(request)
+    where_sql = sql.split('WHERE', 1)[1]
+    assert '`like_cnt` >' not in where_sql
+
+
+def test_publish_time_filter_mapping():
+    condition = FilterCondition(
+        field='publish_time',
+        operator='>=',
+        value='2026-08-01 00:00:00',
+    )
+    sql, params = compile_filter_condition(condition)
+
+    assert sql == '`publish_time` >= %s'
+    assert params == ['2026-08-01 00:00:00']
+
+
+def test_filter_range_and_set_are_parameterized():
+    range_condition = FilterCondition(field='duration', operator='between', value=[60, 300])
+    set_condition = FilterCondition(field='play_cnt', operator='in', value=[10, 20])
+    range_sql, range_params = compile_filter_condition(range_condition)
+    set_sql, set_params = compile_filter_condition(set_condition)
+
+    assert range_sql == '`duration` BETWEEN %s AND %s'
+    assert range_params == [60, 300]
+    assert set_sql == '`play_cnt` IN (%s, %s)'
+    assert set_params == [10, 20]
+
+
+def test_rejects_invalid_filter_range_and_non_finite_number():
+    with pytest.raises(BusinessValidationError, match='起始值不能大于结束值'):
+        compile_filter_condition(
+            FilterCondition(field='duration', operator='between', value=[300, 60])
+        )
+    with pytest.raises(BusinessValidationError, match='有限数字'):
+        compile_filter_condition(
+            FilterCondition(field='like_cnt', operator='>', value='NaN')
+        )
+
+
+def test_rejects_raw_sql_in_filter_value():
+    with pytest.raises(BusinessValidationError, match='必须是数字'):
+        condition = FilterCondition(field='like_cnt', operator='>', value='0 OR 1=1')
+        compile_filter_condition(condition)
+
+
+@pytest.mark.asyncio
+async def test_api_requires_token_and_returns_expected_contract(monkeypatch):
+    local_logs = []
+    cloud_logs = []
+
+    class FakeMySQL:
+        async def fetch_all(self, sql, params):
+            return [
+                {'id': 1, 'platform': params[0], 'create_time': '2026-08-04 12:00:00'},
+                {'id': 2, 'platform': params[0], 'create_time': '2026-08-04 11:00:00'},
+            ]
+
+    class FakeLogger:
+        def info(self, message):
+            local_logs.append(('info', message))
+
+        def error(self, message):
+            local_logs.append(('error', message))
+
+        def exception(self, message):
+            raise AssertionError(message)
+
+    class FakeAliyunLogger:
+        def logging_batch(self, events):
+            cloud_logs.extend(events)
+
+    monkeypatch.setattr(settings, 'CHUI_ZHI_API_TOKEN', 'test-token')
+    app = create_app()
+    app.cleanup_ctx.clear()
+    app[MYSQL_KEY] = FakeMySQL()
+    app[QUERY_SEMAPHORE_KEY] = asyncio.Semaphore(10)
+    app[LOGGER_KEY] = FakeLogger()
+    app[ALIYUN_LOGGER_KEY] = FakeAliyunLogger()
+    client = TestClient(TestServer(app))
+    await client.start_server()
+    payload = {
+        'platforms': ['xiaoniangao'],
+        'start_time': 1785772800000,
+        'end_time': 1785859200000,
+        'filters': [
+            {'field': 'like_cnt', 'operator': '>', 'value': 3},
+        ],
+        'limit': 1,
+    }
+    try:
+        unauthorized = await client.post(API_PATH, json=payload)
+        assert unauthorized.status == 401
+
+        authorized = await client.post(
+            API_PATH,
+            json=payload,
+            headers={'X-API-Key': 'test-token', 'X-Request-ID': 'scheduler-request-1'},
+        )
+        body = await authorized.json()
+        assert authorized.status == 200
+        assert body['code'] == 0
+        assert body['data']['count'] == 1
+        assert body['data']['has_more'] is True
+        assert body['data']['next_cursor'] == {'id': 1}
+        assert body['data']['start_time'] == payload['start_time']
+        assert body['data']['end_time'] == payload['end_time']
+        assert authorized.headers['X-Request-ID'] == 'scheduler-request-1'
+
+        invalid_payload = {
+            **payload,
+            'filters': [{'field': 'like_cnt', 'operator': '>', 'value': 'invalid'}],
+        }
+        failed = await client.post(
+            API_PATH,
+            json=invalid_payload,
+            headers={'X-API-Key': 'test-token'},
+        )
+        assert failed.status == 422
+
+        assert len(cloud_logs) == 3
+        assert all('event_type' not in log['data'] for log in cloud_logs)
+        assert all('request_count' not in log['data'] for log in cloud_logs)
+        assert all(log['data']['path'] == API_PATH for log in cloud_logs)
+        assert all(log['data']['request_duration_ms'] >= 0 for log in cloud_logs)
+        assert cloud_logs[0]['data']['status_code'] == 401
+        assert cloud_logs[0]['data']['request_params']['body'] is None
+        assert cloud_logs[0]['data']['message'] == 'unauthorized'
+        assert cloud_logs[0]['data']['failure_stage'] == 'authentication'
+        assert cloud_logs[0]['data']['error_type'] == 'AuthenticationError'
+        assert cloud_logs[1]['data']['status_code'] == 200
+        assert cloud_logs[1]['trace_id'] == 'scheduler-request-1'
+        assert cloud_logs[1]['data']['request_id'] == 'scheduler-request-1'
+        assert cloud_logs[1]['data']['request_params']['body'] == payload
+        assert 'response' not in cloud_logs[1]['data']
+        assert cloud_logs[1]['data']['success'] is True
+        assert cloud_logs[1]['data']['query_result_count'] == 2
+        assert cloud_logs[1]['data']['query_duration_ms'] >= 0
+        assert cloud_logs[1]['data']['query_success'] is True
+        assert cloud_logs[2]['data']['status_code'] == 422
+        assert cloud_logs[2]['data']['success'] is False
+        assert cloud_logs[2]['data']['query_result_count'] is None
+        assert cloud_logs[2]['data']['query_duration_ms'] is None
+        assert cloud_logs[2]['data']['failure_stage'] == 'business_validation'
+        assert cloud_logs[2]['data']['error_type'] == 'BusinessValidationError'
+        assert 'Traceback (most recent call last)' in cloud_logs[2]['data']['message']
+        assert '必须是数字' in cloud_logs[2]['data']['message']
+        assert cloud_logs[2]['message'] == cloud_logs[2]['data']['message']
+        assert any(level == 'error' and 'status=401' in message for level, message in local_logs)
+        assert any(level == 'info' and 'status=200' in message for level, message in local_logs)
+        assert all('duration_ms=' in message for _, message in local_logs)
+        assert any(level == 'error' and '必须是数字' in message for level, message in local_logs)
+        assert all('response=' not in message for _, message in local_logs)
+    finally:
+        await client.close()
+
+
+@pytest.mark.asyncio
+async def test_cloud_log_failure_does_not_change_api_response(monkeypatch):
+    exception_logs = []
+
+    class FakeMySQL:
+        async def fetch_all(self, sql, params):
+            return []
+
+    class FakeLogger:
+        def info(self, message):
+            pass
+
+        def error(self, message):
+            pass
+
+        def exception(self, message):
+            exception_logs.append(message)
+
+    class BrokenAliyunLogger:
+        def logging_batch(self, events):
+            raise RuntimeError('SLS unavailable')
+
+    monkeypatch.setattr(settings, 'CHUI_ZHI_API_TOKEN', 'test-token')
+    app = create_app()
+    app.cleanup_ctx.clear()
+    app[MYSQL_KEY] = FakeMySQL()
+    app[QUERY_SEMAPHORE_KEY] = asyncio.Semaphore(10)
+    app[LOGGER_KEY] = FakeLogger()
+    app[ALIYUN_LOGGER_KEY] = BrokenAliyunLogger()
+    client = TestClient(TestServer(app))
+    await client.start_server()
+    try:
+        response = await client.post(
+            API_PATH,
+            json={},
+            headers={'X-API-Key': 'test-token'},
+        )
+        assert response.status == 200
+        assert (await response.json())['code'] == 0
+        assert any('阿里云API日志上报失败' in message for message in exception_logs)
+    finally:
+        await client.close()
+
+
+@pytest.mark.asyncio
+async def test_unsupported_parameters_return_specific_error_message(monkeypatch):
+    class FakeMySQL:
+        async def fetch_all(self, sql, params):
+            raise AssertionError('参数校验失败时不应查询数据库')
+
+    class FakeLogger:
+        def info(self, message):
+            pass
+
+        def error(self, message):
+            pass
+
+        def exception(self, message):
+            pass
+
+    class FakeAliyunLogger:
+        def logging_batch(self, events):
+            pass
+
+    monkeypatch.setattr(settings, 'CHUI_ZHI_API_TOKEN', 'test-token')
+    app = create_app()
+    app.cleanup_ctx.clear()
+    app[MYSQL_KEY] = FakeMySQL()
+    app[QUERY_SEMAPHORE_KEY] = asyncio.Semaphore(10)
+    app[LOGGER_KEY] = FakeLogger()
+    app[ALIYUN_LOGGER_KEY] = FakeAliyunLogger()
+    client = TestClient(TestServer(app))
+    await client.start_server()
+    try:
+        unknown_parameter = await client.post(
+            API_PATH,
+            json={'unknown_parameter': 1},
+            headers={'X-API-Key': 'test-token'},
+        )
+        unknown_body = await unknown_parameter.json()
+        assert unknown_parameter.status == 400
+        assert '不支持的参数: unknown_parameter' in unknown_body['msg']
+
+        unsupported_filter = await client.post(
+            API_PATH,
+            json={'filters': [{'field': 'unknown_field', 'operator': '>', 'value': 1}]},
+            headers={'X-API-Key': 'test-token'},
+        )
+        filter_body = await unsupported_filter.json()
+        assert unsupported_filter.status == 400
+        assert 'filters.0.field不支持值 unknown_field' in filter_body['msg']
+
+        unknown_filter_key = await client.post(
+            API_PATH,
+            json={
+                'filters': [
+                    {
+                        'field': 'like_cnt',
+                        'operator': '>',
+                        'value': 1,
+                        'column': 'like_cnt',
+                    }
+                ]
+            },
+            headers={'X-API-Key': 'test-token'},
+        )
+        unknown_key_body = await unknown_filter_key.json()
+        assert unknown_filter_key.status == 400
+        assert '不支持的参数: filters.0.column' in unknown_key_body['msg']
+
+        missing_filter_key = await client.post(
+            API_PATH,
+            json={'filters': [{'field': 'like_cnt', 'value': 1}]},
+            headers={'X-API-Key': 'test-token'},
+        )
+        missing_key_body = await missing_filter_key.json()
+        assert missing_filter_key.status == 400
+        assert 'filters.0.operator不能为空' in missing_key_body['msg']
+    finally:
+        await client.close()
+
+
+def test_api_routes_are_registered():
+    app = create_app()
+    route_names = {route.name for route in app.router.routes()}
+    assert 'chui_zhi_videos' in route_names
+    assert {'health', 'ready'} <= route_names
+
+
+@pytest.mark.asyncio
+async def test_health_and_readiness_do_not_require_business_token():
+    cloud_logs = []
+
+    class FakeMySQL:
+        async def fetch_one(self, sql):
+            assert sql == 'SELECT 1 AS ok'
+            return {'ok': 1}
+
+    class FakeLogger:
+        def info(self, message):
+            pass
+
+        def error(self, message):
+            pass
+
+        def exception(self, message):
+            raise AssertionError(message)
+
+    class FakeAliyunLogger:
+        def logging_batch(self, events):
+            cloud_logs.extend(events)
+
+    app = create_app()
+    app.cleanup_ctx.clear()
+    app[MYSQL_KEY] = FakeMySQL()
+    app[LOGGER_KEY] = FakeLogger()
+    app[ALIYUN_LOGGER_KEY] = FakeAliyunLogger()
+    client = TestClient(TestServer(app))
+    await client.start_server()
+    try:
+        health_response = await client.get('/health')
+        ready_response = await client.get('/ready')
+
+        assert health_response.status == 200
+        assert ready_response.status == 200
+        assert (await health_response.json())['data']['status'] == 'ok'
+        assert (await ready_response.json())['data']['status'] == 'ready'
+        assert len(cloud_logs) == 2
+    finally:
+        await client.close()