fix(HEL-386): 清理任务按 ISO 截止时间删除 job_runs/src_calls

YYYYMMDD 与 ISO 字符串比较会把同年保留期内记录全部误删。

Co-authored-by: Cursor <cursoragent@cursor.com>
Co-authored-by: multica-agent <github@multica.ai>
This commit is contained in:
总工
2026-09-02 12:26:09 +08:00
co-authored by Cursor multica-agent
parent 3498dd7a4b
commit 031eefab4d
2 changed files with 48 additions and 5 deletions
+4 -2
View File
@@ -3,6 +3,7 @@ from __future__ import annotations
import json import json
import time import time
from collections.abc import Callable from collections.abc import Callable
from datetime import timedelta
from typing import Any from typing import Any
from datahub.adapters.base import AdapterError from datahub.adapters.base import AdapterError
@@ -342,8 +343,9 @@ class Pipeline:
def cleanup(self) -> dict[str, int]: def cleanup(self) -> dict[str, int]:
staging_days = int(self.settings.quality.get("staging_retain_days") or 14) 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) job_days = int(self.settings.quality.get("job_run_retain_days") or 90)
cutoff_staging = add_days(yyyymmdd(self.clock()), -staging_days) now = now_shanghai(self.clock())
cutoff_jobs = add_days(yyyymmdd(self.clock()), -job_days) cutoff_staging = add_days(yyyymmdd(now), -staging_days)
cutoff_jobs = isoformat(now - timedelta(days=job_days))
deleted = 0 deleted = 0
with self.db.write() as connection: with self.db.write() as connection:
for dataset, (_eod, staging) in DATASET_TABLES.items(): for dataset, (_eod, staging) in DATASET_TABLES.items():
+44 -3
View File
@@ -2,6 +2,7 @@ from __future__ import annotations
import tempfile import tempfile
import unittest import unittest
from datetime import datetime, timedelta
from pathlib import Path from pathlib import Path
from datahub.adapters.tushare import TushareAdapter from datahub.adapters.tushare import TushareAdapter
@@ -9,23 +10,34 @@ from datahub.crypto import SecretVault
from datahub.db import HubDB from datahub.db import HubDB
from datahub.pipeline import Pipeline, QualityError from datahub.pipeline import Pipeline, QualityError
from datahub.settings import Settings from datahub.settings import Settings
from datahub.timeutil import SHANGHAI, isoformat
from tests.fixtures import TRADE_DATE, fake_transport 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() tmp = tempfile.TemporaryDirectory()
db = HubDB(Path(tmp.name) / "hub.db") db = HubDB(Path(tmp.name) / "hub.db")
adapter = TushareAdapter("test-token", transport=fake_transport) 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( settings = Settings(
encryption_key=SecretVault.generate_key(), encryption_key=SecretVault.generate_key(),
api_token="t" * 32, api_token="t" * 32,
admin_password="admin-pass", admin_password="admin-pass",
tushare_token="test-token", tushare_token="test-token",
db_path=db.path, 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, 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 pipe._tmp = tmp # keep alive
return pipe, db return pipe, db
@@ -103,6 +115,35 @@ class PipelineTests(unittest.TestCase):
mode = connection.execute("PRAGMA journal_mode").fetchone()[0] mode = connection.execute("PRAGMA journal_mode").fetchone()[0]
self.assertEqual(str(mode).lower(), "wal") 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__": if __name__ == "__main__":
unittest.main() unittest.main()