Sfoglia il codice sorgente

修改时区问题

xueyiming 4 giorni fa
parent
commit
9daf4f259f

+ 1 - 1
alembic/env.py

@@ -41,7 +41,7 @@ def run_migrations_online() -> None:
     )
     with connectable.connect() as connection:
         if connection.dialect.name == "mysql":
-            connection.exec_driver_sql("SET time_zone = '+00:00'")
+            connection.exec_driver_sql("SET time_zone = '+08:00'")
             connection.commit()
         context.configure(
             connection=connection,

+ 4 - 2
api/services/scheduler.py

@@ -6,7 +6,7 @@ from typing import Any
 from api.services.pipeline import get_run, pipeline_health
 from supply_infra.config import get_infra_settings
 from supply_infra.pipeline.dag import PIPELINE_STEPS
-from supply_infra.pipeline.dates import next_schedule_utc
+from supply_infra.pipeline.dates import CHINA_TIMEZONE, next_schedule_china
 from supply_infra.pipeline.run_service import submit_pipeline_run
 from supply_infra.scheduler.constants import SUPPLY_PIPELINE_JOB_ID, SUPPLY_PIPELINE_JOB_NAME
 
@@ -65,7 +65,9 @@ def scheduler_status() -> dict[str, Any]:
             {
                 "id": SUPPLY_PIPELINE_JOB_ID,
                 "name": SUPPLY_PIPELINE_JOB_NAME,
-                "next_run_time": next_schedule_utc(settings=settings).isoformat(),
+                "next_run_time": next_schedule_china(settings=settings)
+                .replace(tzinfo=CHINA_TIMEZONE)
+                .isoformat(),
             }
         ],
         **health,

+ 1 - 1
deploy/README.md

@@ -11,7 +11,7 @@ Scheduler、多个 Pipeline Worker 和 Reconciler。
    默认调度时间为 `Asia/Shanghai` 每日 15:00;
 4. `AIGC_DRY_RUN=true` 且 `PIPELINE_EXTERNAL_EFFECTS_ENABLED=false`;
 5. Docker 停止窗口至少 300 秒。
-6. 容器和 Alembic 连接 MySQL 后会将会话时区固定为 UTC;定时日批有次日
+6. 容器和 Alembic 连接 MySQL 后会将会话时区固定为 `+08:00`(中国标准时间);定时日批有次日
    调度 deadline,手工/API/CLI 补数不设置 deadline。
 
 示例:

+ 4 - 4
supply_infra/db/repositories/demand_grade_plan_repo.py

@@ -2,7 +2,6 @@ from __future__ import annotations
 
 import json
 import uuid
-from datetime import datetime
 from typing import Any
 
 from sqlalchemy import func, select, update
@@ -14,6 +13,7 @@ from supply_infra.db.models.demand_grade_plan import (
 )
 from supply_infra.db.repositories.base import BaseRepository
 from supply_infra.db.repositories.demand_grade_repo import DemandGradeRepository
+from supply_infra.pipeline.dates import china_now
 
 
 class DemandGradePlanRepository(BaseRepository[DemandGradePlan]):
@@ -143,7 +143,7 @@ class DemandGradePlanRepository(BaseRepository[DemandGradePlan]):
     def _mark_claimed(group: DemandGradePlanGroup) -> dict[str, Any]:
         group.status = "running"
         group.attempts += 1
-        group.started_at = datetime.now()
+        group.started_at = china_now()
         group.error_message = None
         group.finished_at = None
         return {
@@ -159,7 +159,7 @@ class DemandGradePlanRepository(BaseRepository[DemandGradePlan]):
             .values(
                 status="finished" if success else "failed",
                 error_message=error_message,
-                finished_at=datetime.now(),
+                finished_at=china_now(),
             )
         )
 
@@ -273,7 +273,7 @@ class DemandGradePlanRepository(BaseRepository[DemandGradePlan]):
             return
         values: dict[str, Any] = {"status": status, "error_message": error_message}
         if status in {"finished", "failed", "skipped"}:
-            values["finished_at"] = datetime.now()
+            values["finished_at"] = china_now()
         self.session.execute(
             update(DemandGradePlanGroupItem)
             .where(DemandGradePlanGroupItem.id.in_([int(item_id) for item_id in item_ids]))

+ 7 - 4
supply_infra/db/session.py

@@ -14,11 +14,14 @@ _engine: Any = None
 _SessionLocal: sessionmaker[Session] | None = None
 
 
-def _set_mysql_session_utc(dbapi_connection: Any, _connection_record: Any) -> None:
-    """Keep MySQL-generated timestamps aligned with application UTC values."""
+def _set_mysql_session_china_time(
+    dbapi_connection: Any,
+    _connection_record: Any,
+) -> None:
+    """Store MySQL-generated timestamps as China Standard Time."""
     cursor = dbapi_connection.cursor()
     try:
-        cursor.execute("SET time_zone = '+00:00'")
+        cursor.execute("SET time_zone = '+08:00'")
     finally:
         cursor.close()
 
@@ -38,7 +41,7 @@ def get_engine():
             echo=settings.mysql_echo,
         )
         if _engine.dialect.name == "mysql":
-            event.listen(_engine, "connect", _set_mysql_session_utc)
+            event.listen(_engine, "connect", _set_mysql_session_china_time)
         _SessionLocal = sessionmaker(bind=_engine, autoflush=False, autocommit=False)
     return _engine
 

+ 12 - 9
supply_infra/pipeline/dates.py

@@ -1,10 +1,12 @@
 from __future__ import annotations
 
-from datetime import datetime, timedelta, timezone
+from datetime import datetime, timedelta
 from zoneinfo import ZoneInfo
 
 from supply_infra.config import InfraSettings, get_infra_settings
 
+CHINA_TIMEZONE = ZoneInfo("Asia/Shanghai")
+
 
 def validate_biz_dt(value: str) -> str:
     text = str(value).strip()
@@ -41,11 +43,12 @@ def build_date_snapshot(biz_dt: str) -> dict[str, str]:
     }
 
 
-def utc_now() -> datetime:
-    return datetime.now(timezone.utc).replace(tzinfo=None)
+def china_now() -> datetime:
+    """Return a naive China Standard Time value for MySQL DATETIME columns."""
+    return datetime.now(CHINA_TIMEZONE).replace(tzinfo=None)
 
 
-def next_schedule_utc(
+def next_schedule_china(
     *,
     settings: InfraSettings | None = None,
     now: datetime | None = None,
@@ -65,15 +68,15 @@ def next_schedule_utc(
     )
     if candidate <= local_now:
         candidate += timedelta(days=1)
-    return candidate.astimezone(timezone.utc).replace(tzinfo=None)
+    return candidate.astimezone(CHINA_TIMEZONE).replace(tzinfo=None)
 
 
-def next_daily_deadline_utc(
+def next_daily_deadline_china(
     *,
     settings: InfraSettings | None = None,
     now: datetime | None = None,
 ) -> datetime:
-    """Return the next local calendar day's scheduler start as UTC-naive."""
+    """Return the next local calendar day's scheduler start as China-time naive."""
     active = settings or get_infra_settings()
     zone = ZoneInfo(active.scheduler_timezone)
     local_now = now or datetime.now(zone)
@@ -87,7 +90,7 @@ def next_daily_deadline_utc(
         second=0,
         microsecond=0,
     )
-    return candidate.astimezone(timezone.utc).replace(tzinfo=None)
+    return candidate.astimezone(CHINA_TIMEZONE).replace(tzinfo=None)
 
 
 def deadline_for_trigger(
@@ -99,4 +102,4 @@ def deadline_for_trigger(
     """Scheduled daily runs have a next-day deadline; manual runs do not."""
     if trigger_type not in {"cron", "reconcile"}:
         return None
-    return next_daily_deadline_utc(settings=settings, now=now)
+    return next_daily_deadline_china(settings=settings, now=now)

+ 3 - 3
supply_infra/pipeline/health.py

@@ -7,7 +7,7 @@ from supply_infra.config import InfraSettings, get_infra_settings
 from supply_infra.db.repositories.pipeline_lock_repo import PipelineLockRepository
 from supply_infra.db.repositories.pipeline_run_repo import PipelineRunRepository
 from supply_infra.db.session import get_session
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import china_now
 
 SCHEDULER_HEARTBEAT_KEY = "pipeline:scheduler"
 _ACTIVE_STATUSES = {
@@ -32,7 +32,7 @@ def touch_scheduler_heartbeat(
                 active.pipeline_lease_seconds,
                 active.pipeline_heartbeat_seconds * 3,
             ),
-            now=utc_now(),
+            now=china_now(),
         )
 
 
@@ -63,7 +63,7 @@ def pipeline_health_snapshot(
     settings: InfraSettings | None = None,
 ) -> dict[str, Any]:
     active = settings or get_infra_settings()
-    now = utc_now()
+    now = china_now()
     with get_session() as session:
         scheduler_lock = PipelineLockRepository(session).get(SCHEDULER_HEARTBEAT_KEY)
         latest = PipelineRunRepository(session).list_recent(limit=1)

+ 2 - 2
supply_infra/pipeline/orchestrator.py

@@ -7,7 +7,7 @@ from sqlalchemy.orm import Session
 from supply_infra.db.repositories.pipeline_run_repo import PipelineRunRepository
 from supply_infra.db.repositories.pipeline_step_run_repo import PipelineStepRunRepository
 from supply_infra.pipeline.contracts import StepResult
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import china_now
 from supply_infra.pipeline.enums import RunStatus, StepStatus
 from supply_infra.pipeline.gates import evaluate_step_gate
 
@@ -25,7 +25,7 @@ def complete_step(
     owner: str,
     result: StepResult,
 ) -> str:
-    now = utc_now()
+    now = china_now()
     step_repo = PipelineStepRunRepository(session)
     run_repo = PipelineRunRepository(session)
     step = step_repo.get(step_run_id, for_update=True)

+ 4 - 4
supply_infra/pipeline/reconciler.py

@@ -3,14 +3,14 @@ from __future__ import annotations
 import logging
 import signal
 import threading
-from datetime import timedelta, timezone
+from datetime import timedelta
 from zoneinfo import ZoneInfo
 
 from supply_infra.config import InfraSettings, get_infra_settings
 from supply_infra.db.repositories.pipeline_run_repo import PipelineRunRepository
 from supply_infra.db.repositories.pipeline_step_run_repo import PipelineStepRunRepository
 from supply_infra.db.session import dispose_engine, get_session
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import CHINA_TIMEZONE, china_now
 from supply_infra.pipeline.enums import RunStatus, StepStatus
 from supply_infra.pipeline.health import run_alert_level
 from supply_infra.pipeline.run_service import PIPELINE_KEY, submit_pipeline_run
@@ -20,7 +20,7 @@ logger = logging.getLogger(__name__)
 
 def reconcile_once() -> dict[str, int]:
     settings = get_infra_settings()
-    now = utc_now()
+    now = china_now()
     stats = {
         "retries_promoted": 0,
         "expired_steps": 0,
@@ -167,7 +167,7 @@ def reconcile_once() -> dict[str, int]:
                 )
                 stats["alerts_emitted"] += 1
 
-    local_now = now.replace(tzinfo=timezone.utc).astimezone(
+    local_now = now.replace(tzinfo=CHINA_TIMEZONE).astimezone(
         ZoneInfo(settings.scheduler_timezone)
     )
     missed_cutoff = local_now.replace(

+ 11 - 10
supply_infra/pipeline/run_service.py

@@ -3,7 +3,7 @@ from __future__ import annotations
 import os
 import uuid
 from dataclasses import dataclass
-from datetime import datetime, timezone
+from datetime import datetime
 from typing import Any
 
 from sqlalchemy.exc import IntegrityError
@@ -19,10 +19,11 @@ from supply_infra.db.session import get_session
 from supply_infra.pipeline.dag import PIPELINE_STEPS
 from supply_infra.pipeline.dates import (
     build_date_snapshot,
+    CHINA_TIMEZONE,
+    china_now,
     deadline_for_trigger,
-    next_schedule_utc,
+    next_schedule_china,
     resolve_biz_dt,
-    utc_now,
 )
 from supply_infra.pipeline.enums import RunStatus, StepStatus
 
@@ -84,7 +85,7 @@ def submit_pipeline_run(
     try:
         with get_session() as session:
             run_repo = PipelineRunRepository(session)
-            PipelineLockRepository(session).lock_control_plane(now=utc_now())
+            PipelineLockRepository(session).lock_control_plane(now=china_now())
             existing = run_repo.get_by_dedupe_key(dedupe_key)
             if existing is not None:
                 return RunSubmission(
@@ -225,7 +226,7 @@ def list_pipeline_runs(
 
 
 def cancel_pipeline_run(run_id: str) -> bool:
-    now = utc_now()
+    now = china_now()
     with get_session() as session:
         run_repo = PipelineRunRepository(session)
         step_repo = PipelineStepRunRepository(session)
@@ -313,8 +314,8 @@ def resume_pipeline_run(run_id: str, *, step_key: str | None = None) -> bool:
                 item.error_message = None
         run.status = RunStatus.QUEUED.value
         run.current_step = retry.step_key
-        if run.deadline_at is not None and run.deadline_at <= utc_now():
-            run.deadline_at = next_schedule_utc()
+        if run.deadline_at is not None and run.deadline_at <= china_now():
+            run.deadline_at = next_schedule_china()
         run.finished_at = None
         run.summary_json = None
         run.error_code = None
@@ -374,8 +375,8 @@ def _iso(value: datetime | None) -> str | None:
     if value is None:
         return None
     aware = (
-        value.replace(tzinfo=timezone.utc)
+        value.replace(tzinfo=CHINA_TIMEZONE)
         if value.tzinfo is None
-        else value.astimezone(timezone.utc)
+        else value.astimezone(CHINA_TIMEZONE)
     )
-    return aware.isoformat().replace("+00:00", "Z")
+    return aware.isoformat()

+ 3 - 3
supply_infra/pipeline/worker.py

@@ -13,7 +13,7 @@ from supply_infra.config import InfraSettings, get_infra_settings
 from supply_infra.db.repositories.pipeline_run_repo import PipelineRunRepository
 from supply_infra.db.repositories.pipeline_step_run_repo import PipelineStepRunRepository
 from supply_infra.db.session import dispose_engine, get_session
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import china_now
 from supply_infra.pipeline.orchestrator import complete_step
 from supply_infra.pipeline.contracts import StepResult
 from supply_infra.pipeline.step_runner import StepExecutionRequest, run_step_subprocess
@@ -47,7 +47,7 @@ class PipelineWorker:
         self._stop.set()
 
     def claim(self) -> ClaimedStep | None:
-        now = utc_now()
+        now = china_now()
         with get_session() as session:
             step = PipelineStepRunRepository(session).claim_next_ready(
                 owner=self.worker_id,
@@ -76,7 +76,7 @@ class PipelineWorker:
 
     def _heartbeat_loop(self, claimed: ClaimedStep, stopped: threading.Event) -> None:
         while not stopped.wait(self.settings.pipeline_heartbeat_seconds):
-            now = utc_now()
+            now = china_now()
             lease_until = now + timedelta(seconds=self.settings.pipeline_lease_seconds)
             try:
                 with get_session() as session:

+ 4 - 4
supply_infra/scheduler/job_execution.py

@@ -4,13 +4,13 @@ from __future__ import annotations
 import json
 import logging
 import uuid
-from datetime import datetime
 from typing import Any
 
 from supply_infra.db.repositories.scheduler_job_execution_repo import (
     SchedulerJobExecutionRepository,
 )
 from supply_infra.db.session import get_session
+from supply_infra.pipeline.dates import china_now
 
 logger = logging.getLogger(__name__)
 
@@ -40,7 +40,7 @@ class JobExecutionRecorder:
         self.job_name = job_name
         self.job_id = job_id
         self.biz_dt = biz_dt
-        self.started_at = datetime.now()
+        self.started_at = china_now()
 
     def record_started(self) -> None:
         try:
@@ -64,7 +64,7 @@ class JobExecutionRecorder:
         result: dict[str, Any] | None = None,
         error_message: str | None = None,
     ) -> None:
-        finished_at = datetime.now()
+        finished_at = china_now()
         duration_seconds = round((finished_at - self.started_at).total_seconds(), 2)
         status = STATUS_FINISHED if success else STATUS_FAILED
 
@@ -100,7 +100,7 @@ def record_skipped(
     reason: str,
 ) -> None:
     """记录因并发等原因跳过的执行。"""
-    now = datetime.now()
+    now = china_now()
     run_id = str(uuid.uuid4())
     try:
         with get_session() as session:

+ 10 - 10
tests/supply_infra/pipeline/test_dates.py

@@ -7,8 +7,8 @@ from supply_infra.config import InfraSettings
 from supply_infra.pipeline.dates import (
     build_date_snapshot,
     deadline_for_trigger,
-    next_daily_deadline_utc,
-    next_schedule_utc,
+    next_daily_deadline_china,
+    next_schedule_china,
     resolve_biz_dt,
 )
 
@@ -39,22 +39,22 @@ def test_resolve_biz_dt_uses_scheduler_timezone() -> None:
 
 def test_deadline_is_next_scheduler_start() -> None:
     now = datetime(2026, 7, 27, 15, 1, tzinfo=ZoneInfo("Asia/Shanghai"))
-    deadline = next_schedule_utc(settings=_settings(), now=now)
-    assert deadline.isoformat() == "2026-07-28T07:00:00"
+    deadline = next_schedule_china(settings=_settings(), now=now)
+    assert deadline.isoformat() == "2026-07-28T15:00:00"
 
 
 def test_scheduled_daily_deadline_is_always_next_local_day() -> None:
     before_cron = datetime(2026, 7, 27, 14, 0, tzinfo=ZoneInfo("Asia/Shanghai"))
     after_cron = datetime(2026, 7, 27, 16, 0, tzinfo=ZoneInfo("Asia/Shanghai"))
 
-    assert next_daily_deadline_utc(
+    assert next_daily_deadline_china(
         settings=_settings(),
         now=before_cron,
-    ).isoformat() == "2026-07-28T07:00:00"
-    assert next_daily_deadline_utc(
+    ).isoformat() == "2026-07-28T15:00:00"
+    assert next_daily_deadline_china(
         settings=_settings(),
         now=after_cron,
-    ).isoformat() == "2026-07-28T07:00:00"
+    ).isoformat() == "2026-07-28T15:00:00"
 
 
 def test_only_scheduled_triggers_receive_a_deadline() -> None:
@@ -64,11 +64,11 @@ def test_only_scheduled_triggers_receive_a_deadline() -> None:
         "cron",
         settings=_settings(),
         now=now,
-    ).isoformat() == "2026-07-28T07:00:00"
+    ).isoformat() == "2026-07-28T15:00:00"
     assert deadline_for_trigger(
         "reconcile",
         settings=_settings(),
         now=now,
-    ).isoformat() == "2026-07-28T07:00:00"
+    ).isoformat() == "2026-07-28T15:00:00"
     assert deadline_for_trigger("api", settings=_settings(), now=now) is None
     assert deadline_for_trigger("cli", settings=_settings(), now=now) is None

+ 4 - 4
tests/supply_infra/pipeline/test_health.py

@@ -3,7 +3,7 @@ from __future__ import annotations
 from datetime import timedelta
 
 from supply_infra.config import InfraSettings
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import china_now
 from supply_infra.pipeline.health import run_alert_level
 
 
@@ -17,7 +17,7 @@ def _settings() -> InfraSettings:
 
 def test_alert_level_escalates_against_next_schedule_deadline() -> None:
     settings = _settings()
-    now = utc_now()
+    now = china_now()
 
     assert (
         run_alert_level(
@@ -55,7 +55,7 @@ def test_alert_level_escalates_against_next_schedule_deadline() -> None:
 
 
 def test_terminal_run_has_no_duration_alert() -> None:
-    now = utc_now()
+    now = china_now()
     assert (
         run_alert_level(
             status="failed",
@@ -70,7 +70,7 @@ def test_terminal_run_has_no_duration_alert() -> None:
 
 
 def test_manual_run_without_deadline_only_uses_duration_warning() -> None:
-    now = utc_now()
+    now = china_now()
     settings = _settings()
 
     assert (

+ 27 - 4
tests/supply_infra/pipeline/test_run_serialization.py

@@ -1,12 +1,12 @@
 from __future__ import annotations
 
-from datetime import datetime
+from datetime import datetime, timezone
 
 from supply_infra.db.models.pipeline_run import PipelineRun
 from supply_infra.pipeline.run_service import serialize_run
 
 
-def test_manual_run_serializes_null_deadline_and_utc_timestamps() -> None:
+def test_manual_run_serializes_null_deadline_and_china_timestamps() -> None:
     run = PipelineRun(
         run_id="manual-run",
         dedupe_key="supply_pipeline:20260727:full",
@@ -26,5 +26,28 @@ def test_manual_run_serializes_null_deadline_and_utc_timestamps() -> None:
     payload = serialize_run(run)
 
     assert payload["deadline_at"] is None
-    assert payload["created_at"] == "2026-07-27T06:30:00Z"
-    assert payload["updated_at"] == "2026-07-27T06:31:00Z"
+    assert payload["created_at"] == "2026-07-27T06:30:00+08:00"
+    assert payload["updated_at"] == "2026-07-27T06:31:00+08:00"
+
+
+def test_aware_timestamp_is_converted_to_china_time() -> None:
+    run = PipelineRun(
+        run_id="aware-run",
+        dedupe_key="supply_pipeline:20260727:aware",
+        pipeline_key="supply_pipeline",
+        biz_dt="20260727",
+        trigger_type="api",
+        run_mode="full",
+        dry_run=True,
+        status="queued",
+        deadline_at=None,
+        config_snapshot_json={},
+        date_snapshot_json={},
+        created_at=datetime(2026, 7, 27, 6, 30, tzinfo=timezone.utc),
+        updated_at=datetime(2026, 7, 27, 6, 31, tzinfo=timezone.utc),
+    )
+
+    payload = serialize_run(run)
+
+    assert payload["created_at"] == "2026-07-27T14:30:00+08:00"
+    assert payload["updated_at"] == "2026-07-27T14:31:00+08:00"

+ 4 - 4
tests/supply_infra/pipeline/test_step_claim.py

@@ -11,7 +11,7 @@ from supply_infra.db.models.pipeline_step_run import PipelineStepRun
 from supply_infra.db.repositories.pipeline_step_run_repo import (
     PipelineStepRunRepository,
 )
-from supply_infra.pipeline.dates import utc_now
+from supply_infra.pipeline.dates import china_now
 
 
 def _session() -> Session:
@@ -23,7 +23,7 @@ def _session() -> Session:
 
 
 def _run(run_id: str) -> PipelineRun:
-    now = utc_now()
+    now = china_now()
     return PipelineRun(
         run_id=run_id,
         dedupe_key=f"supply_pipeline:20260727:{run_id}",
@@ -88,7 +88,7 @@ def test_claim_respects_global_active_step_cap() -> None:
             owner="worker-2",
             lease_seconds=120,
             max_active_steps=1,
-            now=utc_now(),
+            now=china_now(),
         )
 
         assert claimed is None
@@ -123,7 +123,7 @@ def test_claim_does_not_overlap_a_different_pipeline_run() -> None:
             owner="worker-2",
             lease_seconds=120,
             max_active_steps=2,
-            now=utc_now(),
+            now=china_now(),
         )
 
         assert claimed is None