mirror of
https://github.com/primedigitaltech/market-assistant.git
synced 2026-07-25 01:34:47 +08:00
fix(pipeline): align models and serializers with pause/resume checkpoint
Made-with: Cursor
This commit is contained in:
parent
b6759039c8
commit
2548ba1df5
65
backend/pipeline/migrations/0012_job_pause_checkpoint.py
Normal file
65
backend/pipeline/migrations/0012_job_pause_checkpoint.py
Normal file
@ -0,0 +1,65 @@
|
|||||||
|
# Generated manually for cookie pause / resume checkpoint
|
||||||
|
|
||||||
|
import django.db.models.deletion
|
||||||
|
from django.db import migrations, models
|
||||||
|
|
||||||
|
|
||||||
|
class Migration(migrations.Migration):
|
||||||
|
|
||||||
|
dependencies = [
|
||||||
|
("pipeline", "0011_job_cancel_and_status"),
|
||||||
|
]
|
||||||
|
|
||||||
|
operations = [
|
||||||
|
migrations.AddField(
|
||||||
|
model_name="pipelinejob",
|
||||||
|
name="resume_from_checkpoint",
|
||||||
|
field=models.BooleanField(db_index=True, default=False),
|
||||||
|
),
|
||||||
|
migrations.AlterField(
|
||||||
|
model_name="pipelinejob",
|
||||||
|
name="status",
|
||||||
|
field=models.CharField(
|
||||||
|
choices=[
|
||||||
|
("pending", "待执行"),
|
||||||
|
("running", "执行中"),
|
||||||
|
("success", "成功"),
|
||||||
|
("failed", "失败"),
|
||||||
|
("cancelled", "已终止"),
|
||||||
|
("paused", "已暂停(待换 Cookie 续跑)"),
|
||||||
|
],
|
||||||
|
db_index=True,
|
||||||
|
default="pending",
|
||||||
|
max_length=16,
|
||||||
|
),
|
||||||
|
),
|
||||||
|
migrations.CreateModel(
|
||||||
|
name="PipelineJobCheckpoint",
|
||||||
|
fields=[
|
||||||
|
(
|
||||||
|
"id",
|
||||||
|
models.BigAutoField(
|
||||||
|
auto_created=True,
|
||||||
|
primary_key=True,
|
||||||
|
serialize=False,
|
||||||
|
verbose_name="ID",
|
||||||
|
),
|
||||||
|
),
|
||||||
|
("phase", models.CharField(db_index=True, max_length=32)),
|
||||||
|
("payload", models.JSONField(blank=True, default=dict)),
|
||||||
|
("hint_zh", models.TextField(blank=True, default="")),
|
||||||
|
("updated_at", models.DateTimeField(auto_now=True)),
|
||||||
|
(
|
||||||
|
"job",
|
||||||
|
models.OneToOneField(
|
||||||
|
on_delete=django.db.models.deletion.CASCADE,
|
||||||
|
related_name="checkpoint_row",
|
||||||
|
to="pipeline.pipelinejob",
|
||||||
|
),
|
||||||
|
),
|
||||||
|
],
|
||||||
|
options={
|
||||||
|
"ordering": ["-updated_at"],
|
||||||
|
},
|
||||||
|
),
|
||||||
|
]
|
||||||
@ -7,6 +7,7 @@ class JobStatus(models.TextChoices):
|
|||||||
SUCCESS = "success", "成功"
|
SUCCESS = "success", "成功"
|
||||||
FAILED = "failed", "失败"
|
FAILED = "failed", "失败"
|
||||||
CANCELLED = "cancelled", "已终止"
|
CANCELLED = "cancelled", "已终止"
|
||||||
|
PAUSED = "paused", "已暂停(待换 Cookie 续跑)"
|
||||||
|
|
||||||
|
|
||||||
class PipelineJob(models.Model):
|
class PipelineJob(models.Model):
|
||||||
@ -43,6 +44,7 @@ class PipelineJob(models.Model):
|
|||||||
db_index=True,
|
db_index=True,
|
||||||
)
|
)
|
||||||
cancellation_requested = models.BooleanField(default=False, db_index=True)
|
cancellation_requested = models.BooleanField(default=False, db_index=True)
|
||||||
|
resume_from_checkpoint = models.BooleanField(default=False, db_index=True)
|
||||||
run_dir = models.TextField(blank=True, default="")
|
run_dir = models.TextField(blank=True, default="")
|
||||||
error_message = models.TextField(blank=True, default="")
|
error_message = models.TextField(blank=True, default="")
|
||||||
created_at = models.DateTimeField(auto_now_add=True)
|
created_at = models.DateTimeField(auto_now_add=True)
|
||||||
@ -55,6 +57,26 @@ class PipelineJob(models.Model):
|
|||||||
return f"[{self.platform}] {self.keyword} ({self.status})"
|
return f"[{self.platform}] {self.keyword} ({self.status})"
|
||||||
|
|
||||||
|
|
||||||
|
class PipelineJobCheckpoint(models.Model):
|
||||||
|
"""Cookie 暂停续跑等场景的断点元数据(与任务一对一)。"""
|
||||||
|
|
||||||
|
job = models.OneToOneField(
|
||||||
|
PipelineJob,
|
||||||
|
on_delete=models.CASCADE,
|
||||||
|
related_name="checkpoint_row",
|
||||||
|
)
|
||||||
|
phase = models.CharField(max_length=32, db_index=True)
|
||||||
|
payload = models.JSONField(default=dict, blank=True)
|
||||||
|
hint_zh = models.TextField(blank=True, default="")
|
||||||
|
updated_at = models.DateTimeField(auto_now=True)
|
||||||
|
|
||||||
|
class Meta:
|
||||||
|
ordering = ["-updated_at"]
|
||||||
|
|
||||||
|
def __str__(self) -> str:
|
||||||
|
return f"checkpoint job={self.job_id} phase={self.phase}"
|
||||||
|
|
||||||
|
|
||||||
class JdProduct(models.Model):
|
class JdProduct(models.Model):
|
||||||
"""京东 SKU 主档:同一 ``platform`` + ``sku_id`` 唯一,多次抓取时覆盖为最新一行合并表数据。"""
|
"""京东 SKU 主档:同一 ``platform`` + ``sku_id`` 唯一,多次抓取时覆盖为最新一行合并表数据。"""
|
||||||
|
|
||||||
|
|||||||
@ -5,7 +5,13 @@ from django.conf import settings
|
|||||||
from rest_framework import serializers
|
from rest_framework import serializers
|
||||||
|
|
||||||
from .cookie_paste import normalize_browser_cookie_paste
|
from .cookie_paste import normalize_browser_cookie_paste
|
||||||
from .models import JdProduct, JdProductSnapshot, JobStatus, PipelineJob
|
from .models import (
|
||||||
|
JdProduct,
|
||||||
|
JdProductSnapshot,
|
||||||
|
JobStatus,
|
||||||
|
PipelineJob,
|
||||||
|
PipelineJobCheckpoint,
|
||||||
|
)
|
||||||
|
|
||||||
# 与 views._safe_file_for_job 中 mapping 一致,供前端展示「数据源是否就绪」
|
# 与 views._safe_file_for_job 中 mapping 一致,供前端展示「数据源是否就绪」
|
||||||
_REPORT_CONFIG_ALLOWED_KEYS = frozenset(
|
_REPORT_CONFIG_ALLOWED_KEYS = frozenset(
|
||||||
@ -54,6 +60,7 @@ class PipelineJobSerializer(serializers.ModelSerializer):
|
|||||||
|
|
||||||
inline_cookie_used = serializers.SerializerMethodField()
|
inline_cookie_used = serializers.SerializerMethodField()
|
||||||
analysis_artifacts = serializers.SerializerMethodField()
|
analysis_artifacts = serializers.SerializerMethodField()
|
||||||
|
checkpoint = serializers.SerializerMethodField()
|
||||||
|
|
||||||
class Meta:
|
class Meta:
|
||||||
model = PipelineJob
|
model = PipelineJob
|
||||||
@ -74,6 +81,8 @@ class PipelineJobSerializer(serializers.ModelSerializer):
|
|||||||
"report_config",
|
"report_config",
|
||||||
"status",
|
"status",
|
||||||
"cancellation_requested",
|
"cancellation_requested",
|
||||||
|
"resume_from_checkpoint",
|
||||||
|
"checkpoint",
|
||||||
"run_dir",
|
"run_dir",
|
||||||
"error_message",
|
"error_message",
|
||||||
"analysis_artifacts",
|
"analysis_artifacts",
|
||||||
@ -84,8 +93,10 @@ class PipelineJobSerializer(serializers.ModelSerializer):
|
|||||||
"id",
|
"id",
|
||||||
"inline_cookie_used",
|
"inline_cookie_used",
|
||||||
"analysis_artifacts",
|
"analysis_artifacts",
|
||||||
|
"checkpoint",
|
||||||
"status",
|
"status",
|
||||||
"cancellation_requested",
|
"cancellation_requested",
|
||||||
|
"resume_from_checkpoint",
|
||||||
"run_dir",
|
"run_dir",
|
||||||
"error_message",
|
"error_message",
|
||||||
"created_at",
|
"created_at",
|
||||||
@ -96,10 +107,24 @@ class PipelineJobSerializer(serializers.ModelSerializer):
|
|||||||
def get_inline_cookie_used(self, obj: PipelineJob) -> bool:
|
def get_inline_cookie_used(self, obj: PipelineJob) -> bool:
|
||||||
return bool((obj.cookie_text or "").strip())
|
return bool((obj.cookie_text or "").strip())
|
||||||
|
|
||||||
|
def get_checkpoint(self, obj: PipelineJob) -> dict | None:
|
||||||
|
try:
|
||||||
|
c = obj.checkpoint_row
|
||||||
|
except PipelineJobCheckpoint.DoesNotExist:
|
||||||
|
return None
|
||||||
|
return {
|
||||||
|
"phase": c.phase,
|
||||||
|
"payload": c.payload,
|
||||||
|
"hint_zh": c.hint_zh,
|
||||||
|
"updated_at": c.updated_at,
|
||||||
|
}
|
||||||
|
|
||||||
def get_analysis_artifacts(self, obj: PipelineJob) -> dict[str, bool] | None:
|
def get_analysis_artifacts(self, obj: PipelineJob) -> dict[str, bool] | None:
|
||||||
if obj.status not in (JobStatus.SUCCESS, JobStatus.CANCELLED) or not (
|
if obj.status not in (
|
||||||
obj.run_dir or ""
|
JobStatus.SUCCESS,
|
||||||
).strip():
|
JobStatus.CANCELLED,
|
||||||
|
JobStatus.PAUSED,
|
||||||
|
) or not (obj.run_dir or "").strip():
|
||||||
return None
|
return None
|
||||||
try:
|
try:
|
||||||
base = Path(obj.run_dir).expanduser().resolve()
|
base = Path(obj.run_dir).expanduser().resolve()
|
||||||
@ -186,6 +211,17 @@ def _jd_data_root() -> Path:
|
|||||||
return (Path(root) / "data" / "JD").resolve()
|
return (Path(root) / "data" / "JD").resolve()
|
||||||
|
|
||||||
|
|
||||||
|
class JobResumeRequestSerializer(serializers.Serializer):
|
||||||
|
"""从断点续跑时可选更新 Cookie。"""
|
||||||
|
|
||||||
|
cookie_text = serializers.CharField(
|
||||||
|
required=False,
|
||||||
|
allow_blank=True,
|
||||||
|
default="",
|
||||||
|
max_length=500_000,
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
class CreatePipelineJobSerializer(serializers.Serializer):
|
class CreatePipelineJobSerializer(serializers.Serializer):
|
||||||
keyword = serializers.CharField(max_length=256, trim_whitespace=True)
|
keyword = serializers.CharField(max_length=256, trim_whitespace=True)
|
||||||
platform = serializers.ChoiceField(choices=["jd"], default="jd")
|
platform = serializers.ChoiceField(choices=["jd"], default="jd")
|
||||||
|
|||||||
@ -15,6 +15,11 @@ urlpatterns = [
|
|||||||
views.JobCancelView.as_view(),
|
views.JobCancelView.as_view(),
|
||||||
name="job-cancel",
|
name="job-cancel",
|
||||||
),
|
),
|
||||||
|
path(
|
||||||
|
"jobs/<int:pk>/resume/",
|
||||||
|
views.JobResumeView.as_view(),
|
||||||
|
name="job-resume",
|
||||||
|
),
|
||||||
path("jobs/<int:pk>/download/", views.JobDownloadView.as_view(), name="job-download"),
|
path("jobs/<int:pk>/download/", views.JobDownloadView.as_view(), name="job-download"),
|
||||||
path("jobs/<int:pk>/preview/", views.JobPreviewView.as_view(), name="job-preview"),
|
path("jobs/<int:pk>/preview/", views.JobPreviewView.as_view(), name="job-preview"),
|
||||||
path(
|
path(
|
||||||
|
|||||||
Loading…
x
Reference in New Issue
Block a user