Compare commits
5 Commits
d5d6ba8b7c
...
prod-prepa
| Author | SHA1 | Date | |
|---|---|---|---|
| 6b69b8c002 | |||
|
|
e726461e75 | ||
| f221aebef0 | |||
|
|
6aba4f0dbd | ||
|
|
09fecdae31 |
@@ -14,6 +14,7 @@ from iaai_scraper.core.config import Settings
|
||||
from iaai_scraper.storage.models import Base
|
||||
|
||||
config = context.config
|
||||
VERSION_TABLE = "iaai_parser_alembic_version"
|
||||
|
||||
# Logging
|
||||
if config.config_file_name is not None:
|
||||
@@ -32,6 +33,7 @@ def run_migrations_offline() -> None:
|
||||
context.configure(
|
||||
url=url,
|
||||
target_metadata=target_metadata,
|
||||
version_table=VERSION_TABLE,
|
||||
literal_binds=True,
|
||||
dialect_opts={"paramstyle": "named"},
|
||||
)
|
||||
@@ -47,7 +49,11 @@ def run_migrations_online() -> None:
|
||||
poolclass=pool.NullPool,
|
||||
)
|
||||
with connectable.connect() as connection:
|
||||
context.configure(connection=connection, target_metadata=target_metadata)
|
||||
context.configure(
|
||||
connection=connection,
|
||||
target_metadata=target_metadata,
|
||||
version_table=VERSION_TABLE,
|
||||
)
|
||||
with context.begin_transaction():
|
||||
context.run_migrations()
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# Начальная схема: cars, images, sync_runs
|
||||
# Начальная схема: iaai_cars, iaai_images, iaai_sync_runs
|
||||
from typing import Sequence, Union
|
||||
|
||||
from alembic import op
|
||||
@@ -10,10 +10,29 @@ branch_labels: Union[str, Sequence[str], None] = None
|
||||
depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def _schema_name() -> str | None:
|
||||
bind = op.get_bind()
|
||||
return "public" if bind.dialect.name == "postgresql" else None
|
||||
|
||||
|
||||
def _table_exists(table_name: str) -> bool:
|
||||
inspector = sa.inspect(op.get_bind())
|
||||
return table_name in inspector.get_table_names(schema=_schema_name())
|
||||
|
||||
|
||||
def _index_exists(table_name: str, index_name: str) -> bool:
|
||||
if not _table_exists(table_name):
|
||||
return False
|
||||
inspector = sa.inspect(op.get_bind())
|
||||
indexes = inspector.get_indexes(table_name, schema=_schema_name())
|
||||
return any(index.get("name") == index_name for index in indexes)
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
# Таблица cars — аукционные машины с IAAI
|
||||
# Таблица iaai_cars — аукционные машины с IAAI
|
||||
if not _table_exists("iaai_cars"):
|
||||
op.create_table(
|
||||
"cars",
|
||||
"iaai_cars",
|
||||
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
||||
sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
|
||||
sa.Column("brand", sa.String(50), nullable=False),
|
||||
@@ -45,22 +64,27 @@ def upgrade() -> None:
|
||||
sa.Column("slug", sa.String, nullable=False),
|
||||
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
|
||||
)
|
||||
op.create_index("ix_cars_origin_url", "cars", ["origin_url"])
|
||||
op.create_index("ix_cars_origin_id", "cars", ["origin_id"], unique=True)
|
||||
|
||||
# Таблица images — изображения машин
|
||||
if not _index_exists("iaai_cars", "ix_iaai_cars_origin_url"):
|
||||
op.create_index("ix_iaai_cars_origin_url", "iaai_cars", ["origin_url"])
|
||||
if not _index_exists("iaai_cars", "ix_iaai_cars_origin_id"):
|
||||
op.create_index("ix_iaai_cars_origin_id", "iaai_cars", ["origin_id"], unique=True)
|
||||
|
||||
# Таблица iaai_images — изображения машин
|
||||
if not _table_exists("iaai_images"):
|
||||
op.create_table(
|
||||
"images",
|
||||
"iaai_images",
|
||||
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
||||
sa.Column("fullres_image", sa.String, nullable=False),
|
||||
sa.Column("preview_image", sa.String, nullable=False),
|
||||
sa.Column("order_index", sa.Integer, nullable=False),
|
||||
sa.Column("car_id", sa.Integer, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
|
||||
sa.Column("car_id", sa.Integer, sa.ForeignKey("iaai_cars.id", ondelete="CASCADE"), nullable=False),
|
||||
)
|
||||
|
||||
# Таблица sync_runs — логирование синхронизаций
|
||||
# Таблица iaai_sync_runs — логирование синхронизаций
|
||||
if not _table_exists("iaai_sync_runs"):
|
||||
op.create_table(
|
||||
"sync_runs",
|
||||
"iaai_sync_runs",
|
||||
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
||||
sa.Column("started_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
|
||||
sa.Column("finished_at", sa.DateTime(timezone=True), nullable=True),
|
||||
@@ -73,7 +97,7 @@ def upgrade() -> None:
|
||||
sa.Column("error_summary", sa.Text, nullable=True),
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_table("sync_runs")
|
||||
op.drop_table("images")
|
||||
op.drop_table("cars")
|
||||
# Compatibility-only migration for shared production databases.
|
||||
pass
|
||||
|
||||
@@ -10,20 +10,15 @@ depends_on: Union[str, Sequence[str], None] = None
|
||||
|
||||
|
||||
def upgrade() -> None:
|
||||
op.create_index("ix_cars_brand", "cars", ["brand"])
|
||||
op.create_index("ix_cars_brand_model", "cars", ["brand", "model"])
|
||||
op.create_index("ix_cars_year", "cars", ["year"])
|
||||
op.create_index("ix_cars_is_sold", "cars", ["is_sold"])
|
||||
op.create_index("ix_cars_last_seen_at", "cars", ["last_seen_at"])
|
||||
op.create_index("ix_images_car_id", "images", ["car_id"])
|
||||
op.create_index("ix_sync_runs_status", "sync_runs", ["status"])
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand ON iaai_cars (brand)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand_model ON iaai_cars (brand, model)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_year ON iaai_cars (year)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_is_sold ON iaai_cars (is_sold)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_last_seen_at ON iaai_cars (last_seen_at)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id ON iaai_images (car_id)")
|
||||
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_status ON iaai_sync_runs (status)")
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.drop_index("ix_sync_runs_status", table_name="sync_runs")
|
||||
op.drop_index("ix_images_car_id", table_name="images")
|
||||
op.drop_index("ix_cars_last_seen_at", table_name="cars")
|
||||
op.drop_index("ix_cars_is_sold", table_name="cars")
|
||||
op.drop_index("ix_cars_year", table_name="cars")
|
||||
op.drop_index("ix_cars_brand_model", table_name="cars")
|
||||
op.drop_index("ix_cars_brand", table_name="cars")
|
||||
# Compatibility-only migration for shared production databases.
|
||||
pass
|
||||
|
||||
@@ -13,26 +13,24 @@ def upgrade() -> None:
|
||||
# Partial index для активных машин IAAI (is_sold = FALSE)
|
||||
# Ускоряет поиск при mark_sold операциях
|
||||
op.execute(
|
||||
"CREATE INDEX IF NOT EXISTS ix_cars_active_iaai "
|
||||
"ON cars (origin_url) "
|
||||
"CREATE INDEX IF NOT EXISTS ix_iaai_cars_active_iaai "
|
||||
"ON iaai_cars (origin_url) "
|
||||
"WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'"
|
||||
)
|
||||
|
||||
# Composite index для проверки дубликатов изображений
|
||||
op.create_index(
|
||||
"ix_images_car_id_fullres",
|
||||
"images",
|
||||
["car_id", "fullres_image"],
|
||||
op.execute(
|
||||
"CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id_fullres "
|
||||
"ON iaai_images (car_id, fullres_image)"
|
||||
)
|
||||
|
||||
# Index для быстрого поиска existing машин по origin_url
|
||||
op.execute(
|
||||
"CREATE INDEX IF NOT EXISTS ix_cars_origin_url_id "
|
||||
"ON cars (origin_url, origin_id)"
|
||||
"CREATE INDEX IF NOT EXISTS ix_iaai_cars_origin_url_id "
|
||||
"ON iaai_cars (origin_url, origin_id)"
|
||||
)
|
||||
|
||||
|
||||
def downgrade() -> None:
|
||||
op.execute("DROP INDEX IF EXISTS ix_cars_origin_url_id")
|
||||
op.drop_index("ix_images_car_id_fullres", table_name="images")
|
||||
op.execute("DROP INDEX IF EXISTS ix_cars_active_iaai")
|
||||
# Compatibility-only migration for shared production databases.
|
||||
pass
|
||||
|
||||
@@ -60,6 +60,11 @@ start_self_heal_watchdog_if_worker() {
|
||||
return 0
|
||||
fi
|
||||
|
||||
# После reboot/restart контейнера файловая система контейнера сохраняется,
|
||||
# поэтому stale pidfile может остаться от прошлого запуска и блокировать старт Celery.
|
||||
# Удаляем его перед запуском worker-процесса.
|
||||
rm -f /tmp/celery-worker.pid
|
||||
|
||||
if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then
|
||||
echo "[entrypoint] Self-heal watchdog disabled"
|
||||
return 0
|
||||
|
||||
@@ -11,6 +11,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
|
||||
from concurrent.futures import TimeoutError as FuturesTimeoutError
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from threading import Event, Thread
|
||||
from typing import Any, Callable
|
||||
from urllib.parse import urlsplit, urlunsplit
|
||||
from urllib.request import Request, urlopen
|
||||
@@ -44,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90
|
||||
_PAGE_NEXT_TIMEOUT_S = 60
|
||||
# Ротация контекста выключена по умолчанию.
|
||||
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
|
||||
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT = int(os.environ.get("IAAI_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT", "3") or 3)
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_LINKS = 120
|
||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
|
||||
|
||||
@@ -97,6 +99,33 @@ class _PageOpWatchdog:
|
||||
return False
|
||||
|
||||
|
||||
class _ProgressHeartbeat:
|
||||
"""Периодически пульсует progress во время долгой навигации."""
|
||||
|
||||
def __init__(self, report_progress: Callable[[str, Any], None], stage: str, interval_s: float = 15.0, **meta: Any) -> None:
|
||||
self._report_progress = report_progress
|
||||
self._stage = stage
|
||||
self._interval_s = max(5.0, float(interval_s))
|
||||
self._meta = meta
|
||||
self._stop = Event()
|
||||
self._thread: Thread | None = None
|
||||
|
||||
def _run(self) -> None:
|
||||
while not self._stop.wait(self._interval_s):
|
||||
self._report_progress(self._stage, **self._meta)
|
||||
|
||||
def __enter__(self):
|
||||
self._thread = Thread(target=self._run, name=f"progress-heartbeat-{self._stage}", daemon=True)
|
||||
self._thread.start()
|
||||
return self
|
||||
|
||||
def __exit__(self, exc_type, exc, tb):
|
||||
self._stop.set()
|
||||
if self._thread is not None:
|
||||
self._thread.join(timeout=1.0)
|
||||
return False
|
||||
|
||||
|
||||
class IAAIScraper:
|
||||
|
||||
@staticmethod
|
||||
@@ -453,16 +482,23 @@ class IAAIScraper:
|
||||
)
|
||||
|
||||
for expected_page in range(2, target_page_number + 1):
|
||||
# Heartbeat для watchdog: при deep-resume (100+ кликов) задача может
|
||||
# идти несколько минут без batch_upserted, поэтому пульсуем прогресс.
|
||||
if expected_page == 2 or expected_page == target_page_number or expected_page % 5 == 0:
|
||||
self._report_progress(
|
||||
"listing_resume_progress",
|
||||
current_page=expected_page - 1,
|
||||
target_page=target_page_number,
|
||||
next_page_number=expected_page,
|
||||
)
|
||||
|
||||
if not self.listing_collector.go_to_next_page(page, expected_page_number=expected_page):
|
||||
with _ProgressHeartbeat(
|
||||
self._report_progress,
|
||||
"listing_resume_progress",
|
||||
current_page=expected_page - 1,
|
||||
target_page=target_page_number,
|
||||
next_page_number=expected_page,
|
||||
):
|
||||
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=expected_page)
|
||||
|
||||
if not next_ok:
|
||||
self._report_progress(
|
||||
"listing_resume_failed",
|
||||
current_page=expected_page - 1,
|
||||
@@ -471,6 +507,13 @@ class IAAIScraper:
|
||||
page.close()
|
||||
raise ListingResumeError(f"Failed to resume listing at page {target_page_number}")
|
||||
|
||||
self._report_progress(
|
||||
"listing_resume_progress",
|
||||
current_page=expected_page,
|
||||
target_page=target_page_number,
|
||||
next_page_number=min(target_page_number, expected_page + 1),
|
||||
)
|
||||
|
||||
self._report_progress(
|
||||
"listing_resume_completed",
|
||||
current_page=target_page_number,
|
||||
@@ -552,6 +595,7 @@ class IAAIScraper:
|
||||
failures: list[dict[str, str]] = []
|
||||
early_stopped = False
|
||||
truncated_by_time_budget = False
|
||||
listing_recovery_attempts = 0
|
||||
|
||||
known_origin_ids: set[str] | None = None
|
||||
threshold = self.settings.listing.early_stop_threshold
|
||||
@@ -683,6 +727,21 @@ class IAAIScraper:
|
||||
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
|
||||
return True
|
||||
|
||||
def _consume_listing_recovery_attempt(*, page_number: int, reason: str) -> None:
|
||||
nonlocal listing_recovery_attempts
|
||||
listing_recovery_attempts += 1
|
||||
logger.warning(
|
||||
"Listing recovery attempt %d/%d at page %d (%s)",
|
||||
listing_recovery_attempts,
|
||||
_MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT,
|
||||
page_number,
|
||||
reason,
|
||||
)
|
||||
if listing_recovery_attempts > _MAX_LISTING_RECOVERY_ATTEMPTS_PER_SEGMENT:
|
||||
raise ListingResumeError(
|
||||
f"Listing recovery exhausted at page {page_number}: {reason}"
|
||||
)
|
||||
|
||||
page = None
|
||||
try:
|
||||
page, applied_filters = self._open_listing_for_stream(
|
||||
@@ -721,6 +780,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason="context_rotation",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -752,6 +815,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason=f"collect_failed:{type(op_exc).__name__}",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -797,6 +864,10 @@ class IAAIScraper:
|
||||
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
|
||||
page.close()
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number,
|
||||
reason="empty_page_after_retry",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number,
|
||||
make=make,
|
||||
@@ -973,6 +1044,15 @@ class IAAIScraper:
|
||||
cars_failed=cars_failed,
|
||||
)
|
||||
try:
|
||||
with _ProgressHeartbeat(
|
||||
self._report_progress,
|
||||
"listing_next_page_started",
|
||||
page_number=page_number,
|
||||
next_page_number=page_number + 1,
|
||||
pending_urls=len(pending_urls),
|
||||
cars_upserted=cars_upserted,
|
||||
cars_failed=cars_failed,
|
||||
):
|
||||
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"):
|
||||
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1)
|
||||
except (PageOperationTimeoutError, PlaywrightError) as nav_exc:
|
||||
@@ -1018,6 +1098,10 @@ class IAAIScraper:
|
||||
pass
|
||||
page = None
|
||||
try:
|
||||
_consume_listing_recovery_attempt(
|
||||
page_number=page_number + 1,
|
||||
reason=f"next_page_failed:{type(nav_exc).__name__}",
|
||||
)
|
||||
page, _ = self._reopen_listing_and_resume(
|
||||
target_page_number=page_number + 1,
|
||||
make=make,
|
||||
|
||||
@@ -20,6 +20,7 @@ CAR_DB_FIELDS = {
|
||||
}
|
||||
|
||||
_IN_CHUNK_SIZE = 5000
|
||||
CAR_TABLE_NAME = Car.__tablename__
|
||||
|
||||
|
||||
class PersistenceService:
|
||||
@@ -532,7 +533,7 @@ class PersistenceService:
|
||||
if is_postgres:
|
||||
# Считаем кандидатов на sold.
|
||||
total_active = session.execute(
|
||||
text("SELECT count(*) FROM cars WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%%'")
|
||||
text(f"SELECT count(*) FROM {CAR_TABLE_NAME} WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%%'")
|
||||
).scalar() or 0
|
||||
|
||||
if total_active == 0:
|
||||
@@ -556,12 +557,12 @@ class PersistenceService:
|
||||
# Считаем будущие sold.
|
||||
would_mark = session.execute(text("""
|
||||
SELECT count(*)
|
||||
FROM cars c
|
||||
FROM {car_table} c
|
||||
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
||||
WHERE a.url IS NULL
|
||||
AND c.is_sold = FALSE
|
||||
AND c.origin_id LIKE 'iaai:%%'
|
||||
""")).scalar() or 0
|
||||
""".format(car_table=CAR_TABLE_NAME))).scalar() or 0
|
||||
|
||||
# Защита от аномалии.
|
||||
if total_active > 100 and would_mark > total_active * 0.8:
|
||||
@@ -573,18 +574,18 @@ class PersistenceService:
|
||||
|
||||
# Массовая пометка sold.
|
||||
result = session.execute(text("""
|
||||
UPDATE cars
|
||||
UPDATE {car_table}
|
||||
SET is_sold = TRUE
|
||||
FROM (
|
||||
SELECT c.id
|
||||
FROM cars c
|
||||
FROM {car_table} c
|
||||
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
||||
WHERE a.url IS NULL
|
||||
AND c.is_sold = FALSE
|
||||
AND c.origin_id LIKE 'iaai:%%'
|
||||
) sub
|
||||
WHERE cars.id = sub.id
|
||||
"""))
|
||||
WHERE {car_table}.id = sub.id
|
||||
""".format(car_table=CAR_TABLE_NAME)))
|
||||
count = result.rowcount or 0
|
||||
else:
|
||||
# Упрощённый путь для SQLite.
|
||||
|
||||
@@ -20,10 +20,10 @@ class Base(DeclarativeBase):
|
||||
|
||||
|
||||
class Car(Base):
|
||||
__tablename__ = "cars"
|
||||
__tablename__ = "iaai_cars"
|
||||
__table_args__ = (
|
||||
Index("ix_cars_brand_model", "brand", "model"),
|
||||
Index("ix_cars_origin_id_not_sold", "origin_id", "is_sold"),
|
||||
Index("ix_iaai_cars_brand_model", "brand", "model"),
|
||||
Index("ix_iaai_cars_origin_id_not_sold", "origin_id", "is_sold"),
|
||||
)
|
||||
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
|
||||
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
|
||||
@@ -59,17 +59,17 @@ class Car(Base):
|
||||
|
||||
|
||||
class Image(Base):
|
||||
__tablename__ = "images"
|
||||
__tablename__ = "iaai_images"
|
||||
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
|
||||
fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
|
||||
preview_image: Mapped[str] = mapped_column(String(), nullable=False)
|
||||
order_index: Mapped[int] = mapped_column(Integer, nullable=False)
|
||||
car_id: Mapped[int] = mapped_column(Integer, ForeignKey("cars.id", ondelete="CASCADE"), nullable=False, index=True)
|
||||
car_id: Mapped[int] = mapped_column(Integer, ForeignKey("iaai_cars.id", ondelete="CASCADE"), nullable=False, index=True)
|
||||
car: Mapped[Car] = relationship("Car", back_populates="images")
|
||||
|
||||
|
||||
class SyncRun(Base):
|
||||
__tablename__ = "sync_runs"
|
||||
__tablename__ = "iaai_sync_runs"
|
||||
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
|
||||
started_at: Mapped[datetime] = mapped_column(DateTime(timezone=True), nullable=False, default=func.now())
|
||||
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
# Инициализация Celery-приложения и периодических задач.
|
||||
|
||||
import logging
|
||||
import os
|
||||
|
||||
from celery import Celery
|
||||
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
||||
@@ -11,6 +12,12 @@ from ..core.logs import setup_logging
|
||||
|
||||
logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
raw = os.getenv(name, "true" if default else "false").strip().lower()
|
||||
return raw in {"1", "true", "yes", "on"}
|
||||
|
||||
|
||||
@celery_setup_logging.connect
|
||||
@@ -79,18 +86,19 @@ celery_app.conf.update(
|
||||
worker_hijack_root_logger=False,
|
||||
beat_schedule={
|
||||
"periodic-sync-listing": {
|
||||
"task": "iaai_scraper.worker.tasks.sync_listing_task",
|
||||
"task": "iaai.sync_cars_feed",
|
||||
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
"args": (),
|
||||
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
|
||||
"options": {
|
||||
"queue": "scraping",
|
||||
"queue": IAAI_SYNC_QUEUE,
|
||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
},
|
||||
}
|
||||
},
|
||||
task_routes={
|
||||
"iaai_scraper.worker.tasks.*": {"queue": "scraping"}
|
||||
"iaai.sync_cars_feed": {"queue": IAAI_SYNC_QUEUE},
|
||||
"iaai_scraper.worker.tasks.*": {"queue": IAAI_SYNC_QUEUE},
|
||||
},
|
||||
)
|
||||
|
||||
@@ -100,6 +108,10 @@ celery_app.autodiscover_tasks(["iaai_scraper.worker"])
|
||||
@worker_ready.connect
|
||||
def _on_worker_ready(**kwargs):
|
||||
"""При старте worker отправляем первый sync_listing, если очередь пуста."""
|
||||
if not _env_bool("IAAI_STARTUP_SYNC_ENABLED", True):
|
||||
logger.info("Worker ready: startup sync dispatch disabled by IAAI_STARTUP_SYNC_ENABLED")
|
||||
return
|
||||
|
||||
redis_client = None
|
||||
try:
|
||||
redis_client = Redis.from_url(
|
||||
@@ -121,11 +133,11 @@ def _on_worker_ready(**kwargs):
|
||||
logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True)
|
||||
|
||||
try:
|
||||
queue_len = int(redis_client.llen("scraping") or 0)
|
||||
queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0)
|
||||
except Exception:
|
||||
queue_len = 0
|
||||
if queue_len > 0:
|
||||
logger.info("Worker ready: scraping queue already has %d task(s); skip startup dispatch", queue_len)
|
||||
logger.info("Worker ready: iaai_sync queue already has %d task(s); skip startup dispatch", queue_len)
|
||||
return
|
||||
|
||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
||||
@@ -145,8 +157,8 @@ def _on_worker_ready(**kwargs):
|
||||
|
||||
logger.info("Worker ready — dispatching initial sync_listing task")
|
||||
celery_app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
"iaai.sync_cars_feed",
|
||||
kwargs={"limit": settings.celery.beat_sync_limit, "only_new": False},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
expires=settings.celery.beat_sync_interval_minutes * 60.0,
|
||||
)
|
||||
|
||||
@@ -12,6 +12,8 @@ logger = logging.getLogger("iaai_scraper.worker.self_heal")
|
||||
|
||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||
SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress"
|
||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||
SIGKILL_FALLBACK = getattr(signal, "SIGKILL", signal.SIGTERM)
|
||||
|
||||
|
||||
def _env_bool(name: str, default: bool) -> bool:
|
||||
@@ -76,6 +78,22 @@ def _read_last_progress_ts(redis_client: Redis) -> int | None:
|
||||
return max_ts or None
|
||||
|
||||
|
||||
def _has_inflight_work(redis_client: Redis, queue_name: str) -> tuple[bool, dict[str, int]]:
|
||||
"""Есть ли признаки активной/зависшей работы, даже если очередь пуста."""
|
||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
||||
has_lock = 1 if redis_client.get(SYNC_LISTING_LOCK_KEY) else 0
|
||||
has_task_progress = 0
|
||||
for _ in redis_client.scan_iter(match="iaai:state:task_progress:*"):
|
||||
has_task_progress = 1
|
||||
break
|
||||
flags = {
|
||||
"queue_len": queue_len,
|
||||
"has_lock": has_lock,
|
||||
"has_task_progress": has_task_progress,
|
||||
}
|
||||
return (queue_len > 0 or has_lock == 1 or has_task_progress == 1), flags
|
||||
|
||||
|
||||
def _kill_worker_process() -> None:
|
||||
pid_file = "/tmp/celery-worker.pid"
|
||||
pid: int | None = None
|
||||
@@ -101,7 +119,7 @@ def _kill_worker_process() -> None:
|
||||
# Если процесс ещё жив — принудительно убиваем.
|
||||
os.kill(pid, 0)
|
||||
logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid)
|
||||
os.kill(pid, signal.SIGKILL)
|
||||
os.kill(pid, SIGKILL_FALLBACK)
|
||||
except ProcessLookupError:
|
||||
pass
|
||||
except Exception:
|
||||
@@ -137,8 +155,8 @@ def main() -> None:
|
||||
redis_client = _get_redis()
|
||||
redis_client.ping()
|
||||
|
||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
||||
if queue_len <= 0:
|
||||
has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
|
||||
if not has_inflight:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
|
||||
@@ -151,8 +169,8 @@ def main() -> None:
|
||||
time.sleep(check_interval)
|
||||
continue
|
||||
logger.warning(
|
||||
"Self-heal: queue=%d but no progress timestamp found after startup grace",
|
||||
queue_len,
|
||||
"Self-heal: inflight=%s but no progress timestamp found after startup grace",
|
||||
inflight,
|
||||
)
|
||||
age = stall_seconds + 1
|
||||
|
||||
@@ -174,8 +192,8 @@ def main() -> None:
|
||||
continue
|
||||
|
||||
logger.error(
|
||||
"Self-heal: detected global stall (queue=%d, progress_age=%ss > %ss). Restarting worker process...",
|
||||
queue_len,
|
||||
"Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...",
|
||||
inflight,
|
||||
age,
|
||||
stall_seconds,
|
||||
)
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
# Задачи Celery.
|
||||
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
||||
|
||||
import json
|
||||
import logging
|
||||
@@ -7,7 +7,6 @@ import signal
|
||||
from threading import Event, Thread
|
||||
import time
|
||||
import uuid
|
||||
from typing import Callable
|
||||
|
||||
from billiard.exceptions import SoftTimeLimitExceeded
|
||||
from celery import shared_task
|
||||
@@ -20,29 +19,17 @@ from ..discovery import SitemapDiscoveryError, discover_vehicle_urls_from_sitema
|
||||
|
||||
logger = logging.getLogger("iaai_scraper.worker.tasks")
|
||||
|
||||
# Пауза между батчами.
|
||||
IAAI_SYNC_QUEUE = "iaai_sync"
|
||||
|
||||
# Минимальная пауза между батчами (секунды) — не давит IAAI.
|
||||
_INTER_BATCH_DELAY = max(float(os.getenv("IAAI_INTER_BATCH_DELAY_SECONDS", "0.3")), 0.0)
|
||||
# Порог ошибок в батче.
|
||||
# Если доля failed в батче превышает порог — прерываем (IAAI блокирует).
|
||||
_FAIL_RATE_THRESHOLD = min(max(float(os.getenv("IAAI_FAIL_RATE_THRESHOLD", "0.9")), 0.0), 1.0)
|
||||
|
||||
SYNC_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
||||
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
|
||||
SYNC_LISTING_RESUME_TARGET_KEY = "iaai:state:sync_listing_resume_target_segment"
|
||||
SYNC_LISTING_RESUME_TARGET_STREAK_KEY = "iaai:state:sync_listing_resume_target_streak"
|
||||
SYNC_LISTING_RESUME_TARGET_TTL_SECONDS = 24 * 60 * 60
|
||||
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT = max(
|
||||
1,
|
||||
int(os.getenv("SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT", "2")),
|
||||
)
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_KEY = "iaai:state:sync_listing_repeat_checkpoint"
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY = "iaai:state:sync_listing_repeat_checkpoint_streak"
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS = 24 * 60 * 60
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT = max(
|
||||
1,
|
||||
int(os.getenv("SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT", "2")),
|
||||
)
|
||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
|
||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
|
||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
|
||||
@@ -57,20 +44,25 @@ SYNC_LISTING_BOOTSTRAP_CONTINUATION_STREAK_TTL_SECONDS = max(
|
||||
)
|
||||
HOURLY_FAILURE_STREAK_KEY = "iaai:state:hourly_failure_streak"
|
||||
HOURLY_FAILURE_STREAK_LIMIT = 3
|
||||
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # Сброс через 6 часов.
|
||||
HOURLY_FAILURE_STREAK_TTL_SECONDS = 6 * 60 * 60 # сброс через 6 часов
|
||||
SYNC_LISTING_FOLLOWUP_PENDING_KEY = "iaai:state:sync_listing_followup_pending"
|
||||
SYNC_LISTING_TASK_NAME = "iaai_scraper.worker.tasks.sync_listing_task"
|
||||
SYNC_LISTING_TASK_NAME = "iaai.sync_cars_feed"
|
||||
SYNC_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
|
||||
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
|
||||
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
|
||||
SYNC_SEGMENTS_CYCLE_KEY = "iaai:state:sync_segments_cycle_id"
|
||||
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
||||
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
||||
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
||||
SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment"
|
||||
SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60
|
||||
STALL_WATCHDOG_NAVIGATION_STAGES = {
|
||||
"listing_next_page_started",
|
||||
"listing_resume_progress",
|
||||
}
|
||||
STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS = max(
|
||||
300,
|
||||
int(os.getenv("STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS", "900")),
|
||||
)
|
||||
|
||||
|
||||
def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
@@ -96,7 +88,10 @@ def _retry_with_backoff(func, *, attempts: int = 5, base_delay_s: float = 1.0):
|
||||
|
||||
|
||||
def _run_browser_job(func, *args, **kwargs):
|
||||
# Запускаем в текущем потоке.
|
||||
# С pool=solo Celery worker работает в одном процессе/потоке.
|
||||
# Playwright sync API использует greenlets, которые привязаны к потоку.
|
||||
# Запуск в отдельном потоке вызывает greenlet.error: cannot switch to a different thread.
|
||||
# Поэтому запускаем напрямую в текущем потоке.
|
||||
return func(*args, **kwargs)
|
||||
|
||||
|
||||
@@ -104,7 +99,8 @@ def _sync_listing_lock_ttl_seconds() -> int:
|
||||
settings = Settings()
|
||||
soft = settings.celery.task_soft_time_limit
|
||||
hard = settings.celery.task_time_limit
|
||||
# Ограничиваем hard limit.
|
||||
# Используем clamped hard limit (soft + 120), а не сырой task_time_limit,
|
||||
# чтобы lock не висел 11 дней при CELERY_TASK_TIME_LIMIT=999999.
|
||||
effective_hard = min(hard, soft + 120) if soft else hard
|
||||
return max(effective_hard + 120, 300)
|
||||
|
||||
@@ -156,6 +152,12 @@ def _clear_task_progress(redis_client: Redis, task_id: str) -> None:
|
||||
logger.warning("Failed to clear task progress for %s", task_id, exc_info=True)
|
||||
|
||||
|
||||
def _stall_timeout_for_progress(stage: str | None, default_timeout: int) -> int:
|
||||
if stage in STALL_WATCHDOG_NAVIGATION_STAGES:
|
||||
return max(int(default_timeout), STALL_WATCHDOG_NAVIGATION_GRACE_SECONDS)
|
||||
return int(default_timeout)
|
||||
|
||||
|
||||
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
|
||||
del limit, only_new
|
||||
settings = Settings()
|
||||
@@ -405,7 +407,6 @@ def _start_stall_watchdog(
|
||||
stall_timeout_seconds: int,
|
||||
lock_key: str | None = None,
|
||||
lock_owner: str | None = None,
|
||||
on_stall: Callable[[dict[str, object]], None] | None = None,
|
||||
) -> tuple[Event, Thread]:
|
||||
stop_event = Event()
|
||||
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
||||
@@ -438,25 +439,22 @@ def _start_stall_watchdog(
|
||||
else:
|
||||
no_data_count = 0
|
||||
data = json.loads(raw)
|
||||
stage = data.get("stage")
|
||||
last_ts = int(data.get("ts") or 0)
|
||||
if not last_ts:
|
||||
continue
|
||||
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||
age = int(time.time()) - last_ts
|
||||
if age < stall_timeout_seconds:
|
||||
if age < effective_stall_timeout:
|
||||
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||
continue
|
||||
logger.error(
|
||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||
task_id,
|
||||
age,
|
||||
data.get("stage"),
|
||||
stage,
|
||||
data,
|
||||
)
|
||||
if on_stall is not None:
|
||||
try:
|
||||
on_stall(data)
|
||||
except Exception:
|
||||
logger.warning("Stall watchdog hook failed for task %s", task_id, exc_info=True)
|
||||
except Exception:
|
||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
||||
@@ -486,7 +484,7 @@ def _start_stall_watchdog(
|
||||
celery_app.send_task(
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=15,
|
||||
expires=followup_ttl,
|
||||
)
|
||||
@@ -731,158 +729,6 @@ def _clear_sync_checkpoint(redis_client: Redis) -> None:
|
||||
logger.warning("Failed to clear sync listing checkpoint", exc_info=True)
|
||||
|
||||
|
||||
def _save_stalled_segment(redis_client: Redis, segment_index: int, *, task_id: str, stage: str | None) -> None:
|
||||
try:
|
||||
payload = {
|
||||
"segment_index": int(segment_index),
|
||||
"task_id": task_id,
|
||||
"stage": stage,
|
||||
"ts": int(time.time()),
|
||||
}
|
||||
redis_client.set(
|
||||
SYNC_LISTING_STALLED_SEGMENT_KEY,
|
||||
json.dumps(payload, ensure_ascii=False),
|
||||
ex=SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to persist stalled segment info", exc_info=True)
|
||||
|
||||
|
||||
def _infer_stalled_segment_index(redis_client: Redis, payload: dict[str, object]) -> int | None:
|
||||
"""Пытается определить индекс застрявшего сегмента из payload или checkpoint."""
|
||||
raw_segment = payload.get("segment_index")
|
||||
if raw_segment is not None:
|
||||
try:
|
||||
return max(0, int(raw_segment))
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
last_completed = _load_last_completed_segment(redis_client)
|
||||
if last_completed is None:
|
||||
return 0
|
||||
return max(0, int(last_completed) + 1)
|
||||
|
||||
|
||||
def _load_stalled_segment(redis_client: Redis) -> int | None:
|
||||
try:
|
||||
raw = redis_client.get(SYNC_LISTING_STALLED_SEGMENT_KEY)
|
||||
except Exception:
|
||||
logger.warning("Failed to read stalled segment info", exc_info=True)
|
||||
return None
|
||||
if not raw:
|
||||
return None
|
||||
try:
|
||||
data = json.loads(raw)
|
||||
return int(data.get("segment_index"))
|
||||
except Exception:
|
||||
logger.warning("Invalid stalled segment payload %r; clearing", raw)
|
||||
_clear_stalled_segment(redis_client)
|
||||
return None
|
||||
|
||||
|
||||
def _clear_stalled_segment(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(SYNC_LISTING_STALLED_SEGMENT_KEY)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear stalled segment info", exc_info=True)
|
||||
|
||||
|
||||
def _register_bootstrap_resume_target(redis_client: Redis, segment_index: int) -> tuple[int, bool]:
|
||||
"""Возвращает (streak, should_skip_segment) для защиты от вечного resume на одном сегменте."""
|
||||
ttl = max(60, int(SYNC_LISTING_RESUME_TARGET_TTL_SECONDS))
|
||||
try:
|
||||
raw_target = redis_client.get(SYNC_LISTING_RESUME_TARGET_KEY)
|
||||
raw_streak = redis_client.get(SYNC_LISTING_RESUME_TARGET_STREAK_KEY)
|
||||
prev_target = int(str(raw_target).strip()) if raw_target is not None else None
|
||||
prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0
|
||||
except Exception:
|
||||
logger.warning("Failed to read bootstrap resume target guard", exc_info=True)
|
||||
prev_target = None
|
||||
prev_streak = 0
|
||||
|
||||
streak = (prev_streak + 1) if prev_target == int(segment_index) else 1
|
||||
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.set(SYNC_LISTING_RESUME_TARGET_KEY, str(int(segment_index)), ex=ttl)
|
||||
pipe.set(SYNC_LISTING_RESUME_TARGET_STREAK_KEY, str(int(streak)), ex=ttl)
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
logger.warning("Failed to persist bootstrap resume target guard", exc_info=True)
|
||||
|
||||
# Если bootstrap второй раз подряд возвращается в тот же сегмент,
|
||||
# считаем его застрявшим и отпускаем checkpoint дальше.
|
||||
should_skip_segment = streak >= SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT
|
||||
if should_skip_segment:
|
||||
logger.error(
|
||||
"Bootstrap resume loop guard OPEN: segment=%d streak=%d limit=%d",
|
||||
segment_index,
|
||||
streak,
|
||||
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT,
|
||||
)
|
||||
elif streak > 1:
|
||||
logger.warning(
|
||||
"Bootstrap resume repeated for segment=%d (%d/%d)",
|
||||
segment_index,
|
||||
streak,
|
||||
SYNC_LISTING_RESUME_TARGET_STREAK_LIMIT,
|
||||
)
|
||||
return streak, should_skip_segment
|
||||
|
||||
|
||||
def _register_repeated_bootstrap_checkpoint(redis_client: Redis, checkpoint_segment: int) -> tuple[int, bool]:
|
||||
"""Возвращает (streak, should_skip_next_segment), если bootstrap снова стартует с тем же checkpoint."""
|
||||
ttl = max(60, int(SYNC_LISTING_REPEAT_CHECKPOINT_TTL_SECONDS))
|
||||
try:
|
||||
raw_checkpoint = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_KEY)
|
||||
raw_streak = redis_client.get(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY)
|
||||
prev_checkpoint = int(str(raw_checkpoint).strip()) if raw_checkpoint is not None else None
|
||||
prev_streak = int(str(raw_streak).strip()) if raw_streak is not None else 0
|
||||
except Exception:
|
||||
logger.warning("Failed to read repeated bootstrap checkpoint guard", exc_info=True)
|
||||
prev_checkpoint = None
|
||||
prev_streak = 0
|
||||
|
||||
streak = (prev_streak + 1) if prev_checkpoint == int(checkpoint_segment) else 1
|
||||
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_KEY, str(int(checkpoint_segment)), ex=ttl)
|
||||
pipe.set(SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY, str(int(streak)), ex=ttl)
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
logger.warning("Failed to persist repeated bootstrap checkpoint guard", exc_info=True)
|
||||
|
||||
should_skip_next_segment = streak >= SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT
|
||||
if should_skip_next_segment:
|
||||
logger.error(
|
||||
"Bootstrap checkpoint repeat guard OPEN: checkpoint=%d streak=%d limit=%d",
|
||||
checkpoint_segment,
|
||||
streak,
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT,
|
||||
)
|
||||
elif streak > 1:
|
||||
logger.warning(
|
||||
"Bootstrap repeated with same checkpoint=%d (%d/%d)",
|
||||
checkpoint_segment,
|
||||
streak,
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_LIMIT,
|
||||
)
|
||||
return streak, should_skip_next_segment
|
||||
|
||||
|
||||
def _clear_bootstrap_resume_target(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(
|
||||
SYNC_LISTING_RESUME_TARGET_KEY,
|
||||
SYNC_LISTING_RESUME_TARGET_STREAK_KEY,
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_KEY,
|
||||
SYNC_LISTING_REPEAT_CHECKPOINT_STREAK_KEY,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear bootstrap resume target guard", exc_info=True)
|
||||
|
||||
|
||||
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
|
||||
try:
|
||||
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
|
||||
@@ -1040,59 +886,17 @@ def _reset_segments_progress(redis_client: Redis, total: int) -> None:
|
||||
logger.warning("Failed to reset segments progress", exc_info=True)
|
||||
|
||||
|
||||
def _clear_segments_progress_state(redis_client: Redis) -> None:
|
||||
try:
|
||||
redis_client.delete(
|
||||
SYNC_SEGMENTS_PROGRESS_KEY,
|
||||
SYNC_SEGMENTS_TOTAL_KEY,
|
||||
SYNC_SEGMENTS_CYCLE_KEY,
|
||||
)
|
||||
except Exception:
|
||||
logger.warning("Failed to clear segments progress state", exc_info=True)
|
||||
|
||||
|
||||
def _ensure_segments_cycle(redis_client: Redis, total: int) -> tuple[str, bool]:
|
||||
"""Возвращает (cycle_id, created_new_cycle)."""
|
||||
ttl = max(60, int(SYNC_SEGMENTS_PROGRESS_TTL_SECONDS))
|
||||
try:
|
||||
raw_cycle = redis_client.get(SYNC_SEGMENTS_CYCLE_KEY)
|
||||
cycle_id = str(raw_cycle).strip() if raw_cycle else ""
|
||||
except Exception:
|
||||
logger.warning("Failed to read segments cycle id", exc_info=True)
|
||||
cycle_id = ""
|
||||
|
||||
if cycle_id:
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.expire(SYNC_SEGMENTS_CYCLE_KEY, ttl)
|
||||
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, ttl)
|
||||
pipe.set(SYNC_SEGMENTS_TOTAL_KEY, str(int(total)), ex=ttl)
|
||||
pipe.execute()
|
||||
except Exception:
|
||||
logger.warning("Failed to refresh segments cycle ttl", exc_info=True)
|
||||
return cycle_id, False
|
||||
|
||||
cycle_id = uuid.uuid4().hex
|
||||
try:
|
||||
_reset_segments_progress(redis_client, total)
|
||||
redis_client.set(SYNC_SEGMENTS_CYCLE_KEY, cycle_id, ex=ttl)
|
||||
except Exception:
|
||||
logger.warning("Failed to initialize segments cycle id", exc_info=True)
|
||||
return cycle_id, True
|
||||
|
||||
|
||||
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
|
||||
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
|
||||
try:
|
||||
pipe = redis_client.pipeline()
|
||||
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
|
||||
pipe.expire(SYNC_SEGMENTS_PROGRESS_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
pipe.expire(SYNC_SEGMENTS_CYCLE_KEY, SYNC_SEGMENTS_PROGRESS_TTL_SECONDS)
|
||||
pipe.scard(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
|
||||
results = pipe.execute()
|
||||
completed = int(results[3] or 0)
|
||||
total = int(results[4] or 0) if results[4] else 0
|
||||
completed = int(results[2] or 0)
|
||||
total = int(results[3] or 0) if results[3] else 0
|
||||
return completed, total
|
||||
except Exception:
|
||||
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
|
||||
@@ -1101,6 +905,7 @@ def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[in
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=60,
|
||||
@@ -1274,6 +1079,7 @@ def sync_segment_task(
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_vehicle_task",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=2,
|
||||
default_retry_delay=30,
|
||||
@@ -1304,7 +1110,8 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
|
||||
|
||||
|
||||
@shared_task(
|
||||
name="iaai_scraper.worker.tasks.sync_listing_task",
|
||||
name=SYNC_LISTING_TASK_NAME,
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
bind=True,
|
||||
max_retries=3,
|
||||
default_retry_delay=120,
|
||||
@@ -1376,7 +1183,7 @@ def sync_listing_task(
|
||||
try:
|
||||
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
|
||||
self.app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={
|
||||
"make": make,
|
||||
"model": model,
|
||||
@@ -1384,7 +1191,7 @@ def sync_listing_task(
|
||||
"limit": limit,
|
||||
"only_new": only_new,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=max(0, int(delay_seconds)),
|
||||
expires=followup_expires,
|
||||
)
|
||||
@@ -1436,19 +1243,6 @@ def sync_listing_task(
|
||||
stall_timeout_seconds=stall_timeout,
|
||||
lock_key=SYNC_LISTING_LOCK_KEY,
|
||||
lock_owner=owner_token,
|
||||
on_stall=(
|
||||
lambda payload: (
|
||||
_save_stalled_segment(
|
||||
redis_client,
|
||||
stalled_segment_index,
|
||||
task_id=task_id,
|
||||
stage=str(payload.get("stage") or ""),
|
||||
)
|
||||
if force_bootstrap_full_scan
|
||||
and (stalled_segment_index := _infer_stalled_segment_index(redis_client, payload)) is not None
|
||||
else None
|
||||
)
|
||||
),
|
||||
)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
@@ -1484,11 +1278,9 @@ def sync_listing_task(
|
||||
last_completed_segment: int | None = None
|
||||
if force_bootstrap_full_scan:
|
||||
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||
stalled_segment = _load_stalled_segment(redis_client)
|
||||
else:
|
||||
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
stalled_segment = None
|
||||
|
||||
if force_bootstrap_full_scan:
|
||||
logger.info(
|
||||
@@ -1649,17 +1441,12 @@ def sync_listing_task(
|
||||
if use_segmented and settings.celery.parallel_segments:
|
||||
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
|
||||
already_completed: set[int] = set()
|
||||
cycle_id: str | None = None
|
||||
if force_bootstrap_full_scan:
|
||||
try:
|
||||
cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments))
|
||||
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
|
||||
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
|
||||
if created_new_cycle:
|
||||
logger.warning("Started new bootstrap cycle: %s", cycle_id)
|
||||
except Exception:
|
||||
already_completed = set()
|
||||
cycle_id = None
|
||||
|
||||
pending = [
|
||||
(idx, seg) for idx, seg in enumerate(segments)
|
||||
@@ -1671,7 +1458,6 @@ def sync_listing_task(
|
||||
if force_bootstrap_full_scan:
|
||||
_set_full_scan_done(redis_client, True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_segments_progress_state(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
logger.warning("Parallel segments: nothing to dispatch (all completed)")
|
||||
return {
|
||||
@@ -1683,6 +1469,10 @@ def sync_listing_task(
|
||||
"segments_already_completed": len(already_completed),
|
||||
}
|
||||
|
||||
# При первом запуске bootstrap фиксируем total, чтобы знать когда остановиться.
|
||||
if force_bootstrap_full_scan and not already_completed:
|
||||
_reset_segments_progress(redis_client, len(segments))
|
||||
|
||||
dispatched = 0
|
||||
for idx, seg in pending:
|
||||
try:
|
||||
@@ -1695,15 +1485,15 @@ def sync_listing_task(
|
||||
"only_new": effective_only_new,
|
||||
"is_bootstrap": force_bootstrap_full_scan,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
)
|
||||
dispatched += 1
|
||||
except Exception:
|
||||
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
|
||||
|
||||
logger.warning(
|
||||
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)",
|
||||
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id,
|
||||
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
|
||||
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
|
||||
)
|
||||
return {
|
||||
"status": "success",
|
||||
@@ -1712,11 +1502,9 @@ def sync_listing_task(
|
||||
"segments_total": len(segments),
|
||||
"segments_dispatched": dispatched,
|
||||
"segments_already_completed": len(already_completed),
|
||||
"segments_cycle_id": cycle_id,
|
||||
}
|
||||
# --- конец параллельной ветки ---
|
||||
|
||||
precomputed_result = None
|
||||
resume_from_segment = 0
|
||||
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
||||
resume_from_segment = max(0, last_completed_segment + 1)
|
||||
@@ -1734,117 +1522,6 @@ def sync_listing_task(
|
||||
resume_from_segment, last_completed_segment,
|
||||
)
|
||||
|
||||
checkpoint_streak, should_skip_from_checkpoint = _register_repeated_bootstrap_checkpoint(
|
||||
redis_client,
|
||||
last_completed_segment,
|
||||
)
|
||||
if should_skip_from_checkpoint:
|
||||
skipped_segment = resume_from_segment
|
||||
logger.error(
|
||||
"Checkpoint %d repeated (streak=%d); advancing past segment=%d before resume",
|
||||
last_completed_segment,
|
||||
checkpoint_streak,
|
||||
skipped_segment,
|
||||
)
|
||||
_save_last_completed_segment(redis_client, skipped_segment)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="bootstrap_repeat_checkpoint_segment_skipped",
|
||||
ttl_seconds=progress_ttl,
|
||||
checkpoint_segment=last_completed_segment,
|
||||
segment_index=skipped_segment,
|
||||
checkpoint_streak=checkpoint_streak,
|
||||
)
|
||||
resume_from_segment = skipped_segment + 1
|
||||
if resume_from_segment >= len(segments):
|
||||
precomputed_result = {
|
||||
"status": "partial_success",
|
||||
"run_id": None,
|
||||
"full_scan_completed": True,
|
||||
"cars_upserted": 0,
|
||||
"cars_failed": 0,
|
||||
"images_upserted": 0,
|
||||
"skipped_existing": 0,
|
||||
"elapsed_seconds": 0,
|
||||
"failures": [{
|
||||
"vehicle_url": f"segment_{skipped_segment}",
|
||||
"error": f"Skipped after repeated bootstrap checkpoint loops on segment {skipped_segment}",
|
||||
}],
|
||||
}
|
||||
else:
|
||||
logger.warning(
|
||||
"Bootstrap resume will continue from segment=%d after repeated checkpoint skip of segment=%d",
|
||||
resume_from_segment,
|
||||
skipped_segment,
|
||||
)
|
||||
|
||||
if (
|
||||
force_bootstrap_full_scan
|
||||
and use_segmented
|
||||
and stalled_segment is not None
|
||||
and 0 <= stalled_segment < len(segments)
|
||||
and stalled_segment >= resume_from_segment
|
||||
):
|
||||
logger.error(
|
||||
"Skipping previously stalled segment=%d and advancing checkpoint before resume",
|
||||
stalled_segment,
|
||||
)
|
||||
_save_last_completed_segment(redis_client, stalled_segment)
|
||||
_clear_stalled_segment(redis_client)
|
||||
resume_from_segment = stalled_segment + 1
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="bootstrap_stalled_segment_skipped",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=stalled_segment,
|
||||
)
|
||||
|
||||
if force_bootstrap_full_scan and use_segmented and resume_from_segment < len(segments):
|
||||
resume_streak, should_skip_segment = _register_bootstrap_resume_target(
|
||||
redis_client,
|
||||
resume_from_segment,
|
||||
)
|
||||
if should_skip_segment:
|
||||
skipped_segment = resume_from_segment
|
||||
logger.error(
|
||||
"Segment %d is stuck on bootstrap resume (streak=%d); skipping it and advancing checkpoint",
|
||||
skipped_segment,
|
||||
resume_streak,
|
||||
)
|
||||
_save_last_completed_segment(redis_client, skipped_segment)
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
task_id=task_id,
|
||||
stage="bootstrap_resume_segment_skipped",
|
||||
ttl_seconds=progress_ttl,
|
||||
segment_index=skipped_segment,
|
||||
resume_streak=resume_streak,
|
||||
)
|
||||
resume_from_segment = skipped_segment + 1
|
||||
if resume_from_segment >= len(segments):
|
||||
precomputed_result = {
|
||||
"status": "partial_success",
|
||||
"run_id": None,
|
||||
"full_scan_completed": True,
|
||||
"cars_upserted": 0,
|
||||
"cars_failed": 0,
|
||||
"images_upserted": 0,
|
||||
"skipped_existing": 0,
|
||||
"elapsed_seconds": 0,
|
||||
"failures": [{
|
||||
"vehicle_url": f"segment_{skipped_segment}",
|
||||
"error": f"Skipped after repeated bootstrap resume loops on segment {skipped_segment}",
|
||||
}],
|
||||
}
|
||||
else:
|
||||
logger.warning(
|
||||
"Bootstrap resume will continue from segment=%d after skipping segment=%d",
|
||||
resume_from_segment,
|
||||
skipped_segment,
|
||||
)
|
||||
|
||||
def _progress_cb_main(stage, meta):
|
||||
_update_task_progress(
|
||||
redis_client,
|
||||
@@ -1877,16 +1554,13 @@ def sync_listing_task(
|
||||
only_new=effective_only_new,
|
||||
)
|
||||
|
||||
result = precomputed_result if precomputed_result is not None else _run_browser_job(_job)
|
||||
result = _run_browser_job(_job)
|
||||
|
||||
if force_bootstrap_full_scan:
|
||||
bootstrap_completed = bool(result.get("full_scan_completed"))
|
||||
if bootstrap_completed:
|
||||
_set_full_scan_done(redis_client, False if always_full_scan else True)
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_segments_progress_state(redis_client)
|
||||
_clear_bootstrap_resume_target(redis_client)
|
||||
_clear_stalled_segment(redis_client)
|
||||
_clear_bootstrap_failure_streak(redis_client)
|
||||
_clear_bootstrap_continuation_streak(redis_client)
|
||||
if always_full_scan:
|
||||
@@ -1902,7 +1576,7 @@ def sync_listing_task(
|
||||
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
|
||||
) or int(listing_payload.get("vehicles_collected") or 0) > 0
|
||||
count_as_failure = anti_bot_detected or (str(result.get("status") or "") == "failed" and not had_progress)
|
||||
followup_delay = max(3600, int(settings.celery.beat_sync_interval_minutes) * 60)
|
||||
followup_delay = 5
|
||||
followup_reason = "bootstrap_not_completed"
|
||||
if anti_bot_detected:
|
||||
followup_delay = 180
|
||||
@@ -1922,7 +1596,6 @@ def sync_listing_task(
|
||||
)
|
||||
else:
|
||||
_clear_sync_checkpoint(redis_client)
|
||||
_clear_bootstrap_resume_target(redis_client)
|
||||
|
||||
summary = {
|
||||
"task_id": task_id,
|
||||
@@ -1986,7 +1659,7 @@ def sync_listing_task(
|
||||
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
|
||||
try:
|
||||
self.app.send_task(
|
||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
||||
SYNC_LISTING_TASK_NAME,
|
||||
kwargs={
|
||||
"make": make,
|
||||
"model": model,
|
||||
@@ -1994,7 +1667,7 @@ def sync_listing_task(
|
||||
"limit": limit,
|
||||
"only_new": only_new,
|
||||
},
|
||||
queue="scraping",
|
||||
queue=IAAI_SYNC_QUEUE,
|
||||
countdown=10,
|
||||
expires=followup_ttl,
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user