From 031eefab4d18bed527f7c0f84f057561779c3d5d Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E6=80=BB=E5=B7=A5?= Date: Wed, 2 Sep 2026 12:26:09 +0800 Subject: [PATCH] =?UTF-8?q?fix(HEL-386):=20=E6=B8=85=E7=90=86=E4=BB=BB?= =?UTF-8?q?=E5=8A=A1=E6=8C=89=20ISO=20=E6=88=AA=E6=AD=A2=E6=97=B6=E9=97=B4?= =?UTF-8?q?=E5=88=A0=E9=99=A4=20job=5Fruns/src=5Fcalls?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit YYYYMMDD 与 ISO 字符串比较会把同年保留期内记录全部误删。 Co-authored-by: Cursor Co-authored-by: multica-agent --- xiaobai-datahub/datahub/pipeline.py | 6 ++-- xiaobai-datahub/tests/test_pipeline.py | 47 ++++++++++++++++++++++++-- 2 files changed, 48 insertions(+), 5 deletions(-) diff --git a/xiaobai-datahub/datahub/pipeline.py b/xiaobai-datahub/datahub/pipeline.py index 356ac5f..7f447e9 100644 --- a/xiaobai-datahub/datahub/pipeline.py +++ b/xiaobai-datahub/datahub/pipeline.py @@ -3,6 +3,7 @@ from __future__ import annotations import json import time from collections.abc import Callable +from datetime import timedelta from typing import Any from datahub.adapters.base import AdapterError @@ -342,8 +343,9 @@ class Pipeline: def cleanup(self) -> dict[str, int]: staging_days = int(self.settings.quality.get("staging_retain_days") or 14) job_days = int(self.settings.quality.get("job_run_retain_days") or 90) - cutoff_staging = add_days(yyyymmdd(self.clock()), -staging_days) - cutoff_jobs = add_days(yyyymmdd(self.clock()), -job_days) + now = now_shanghai(self.clock()) + cutoff_staging = add_days(yyyymmdd(now), -staging_days) + cutoff_jobs = isoformat(now - timedelta(days=job_days)) deleted = 0 with self.db.write() as connection: for dataset, (_eod, staging) in DATASET_TABLES.items(): diff --git a/xiaobai-datahub/tests/test_pipeline.py b/xiaobai-datahub/tests/test_pipeline.py index ea3d577..e435f18 100644 --- a/xiaobai-datahub/tests/test_pipeline.py +++ b/xiaobai-datahub/tests/test_pipeline.py @@ -2,6 +2,7 @@ from __future__ import annotations import tempfile import unittest +from datetime import datetime, timedelta from pathlib import Path from datahub.adapters.tushare import TushareAdapter @@ -9,23 +10,34 @@ from datahub.crypto import SecretVault from datahub.db import HubDB from datahub.pipeline import Pipeline, QualityError from datahub.settings import Settings +from datahub.timeutil import SHANGHAI, isoformat from tests.fixtures import TRADE_DATE, fake_transport -def make_pipeline(before_commit=None) -> tuple[Pipeline, HubDB]: +def make_pipeline(before_commit=None, clock=None, quality=None) -> tuple[Pipeline, HubDB]: tmp = tempfile.TemporaryDirectory() db = HubDB(Path(tmp.name) / "hub.db") adapter = TushareAdapter("test-token", transport=fake_transport) + quality_cfg = { + "daily_row_ratio": 0.98, + "null_rate_max": 0.01, + "max_publish_attempts": 3, + "publication_generations": 3, + "job_run_retain_days": 90, + "staging_retain_days": 14, + } + if quality: + quality_cfg.update(quality) settings = Settings( encryption_key=SecretVault.generate_key(), api_token="t" * 32, admin_password="admin-pass", tushare_token="test-token", db_path=db.path, - quality={"daily_row_ratio": 0.98, "null_rate_max": 0.01, "max_publish_attempts": 3, "publication_generations": 3}, + quality=quality_cfg, scheduler_enabled=False, ) - pipe = Pipeline(db, adapter, settings, before_commit=before_commit) + pipe = Pipeline(db, adapter, settings, before_commit=before_commit, clock=clock) pipe._tmp = tmp # keep alive return pipe, db @@ -103,6 +115,35 @@ class PipelineTests(unittest.TestCase): mode = connection.execute("PRAGMA journal_mode").fetchone()[0] self.assertEqual(str(mode).lower(), "wal") + def test_cleanup_iso_timestamps_respect_retention_on_job_and_src(self) -> None: + # job_runs.started_at / src_calls.created_at 存 ISO;旧实现用 YYYYMMDD 比较会误删同年记录。 + frozen = datetime(2026, 9, 2, 0, 30, tzinfo=SHANGHAI) + retain_days = 90 + pipe, db = make_pipeline(clock=lambda: frozen, quality={"job_run_retain_days": retain_days}) + samples = { + "today": isoformat(frozen), + "within": isoformat(frozen - timedelta(days=retain_days - 1)), + "expired": isoformat(frozen - timedelta(days=retain_days + 1)), + } + with db.write() as connection: + for job_id, stamp in samples.items(): + connection.execute( + "INSERT INTO job_runs(job_id, state, started_at, attempt) VALUES (?,?,?,1)", + (job_id, "ok", stamp), + ) + connection.execute( + "INSERT INTO src_calls(provider, endpoint, ok, latency_ms, error, created_at) VALUES (?,?,?,?,?,?)", + ("tushare", job_id, 1, 10, None, stamp), + ) + + pipe.cleanup() + + jobs = {row["job_id"] for row in db.fetchall("SELECT job_id FROM job_runs")} + calls = {row["endpoint"] for row in db.fetchall("SELECT endpoint FROM src_calls")} + kept = {"today", "within"} + self.assertEqual(jobs, kept) + self.assertEqual(calls, kept) + if __name__ == "__main__": unittest.main()