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
|
from iaai_scraper.storage.models import Base
|
||||||
|
|
||||||
config = context.config
|
config = context.config
|
||||||
|
VERSION_TABLE = "iaai_parser_alembic_version"
|
||||||
|
|
||||||
# Logging
|
# Logging
|
||||||
if config.config_file_name is not None:
|
if config.config_file_name is not None:
|
||||||
@@ -32,6 +33,7 @@ def run_migrations_offline() -> None:
|
|||||||
context.configure(
|
context.configure(
|
||||||
url=url,
|
url=url,
|
||||||
target_metadata=target_metadata,
|
target_metadata=target_metadata,
|
||||||
|
version_table=VERSION_TABLE,
|
||||||
literal_binds=True,
|
literal_binds=True,
|
||||||
dialect_opts={"paramstyle": "named"},
|
dialect_opts={"paramstyle": "named"},
|
||||||
)
|
)
|
||||||
@@ -47,7 +49,11 @@ def run_migrations_online() -> None:
|
|||||||
poolclass=pool.NullPool,
|
poolclass=pool.NullPool,
|
||||||
)
|
)
|
||||||
with connectable.connect() as connection:
|
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():
|
with context.begin_transaction():
|
||||||
context.run_migrations()
|
context.run_migrations()
|
||||||
|
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# Начальная схема: cars, images, sync_runs
|
# Начальная схема: iaai_cars, iaai_images, iaai_sync_runs
|
||||||
from typing import Sequence, Union
|
from typing import Sequence, Union
|
||||||
|
|
||||||
from alembic import op
|
from alembic import op
|
||||||
@@ -10,70 +10,94 @@ branch_labels: Union[str, Sequence[str], None] = None
|
|||||||
depends_on: 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:
|
def upgrade() -> None:
|
||||||
# Таблица cars — аукционные машины с IAAI
|
# Таблица iaai_cars — аукционные машины с IAAI
|
||||||
op.create_table(
|
if not _table_exists("iaai_cars"):
|
||||||
"cars",
|
op.create_table(
|
||||||
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
"iaai_cars",
|
||||||
sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
|
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
||||||
sa.Column("brand", sa.String(50), nullable=False),
|
sa.Column("parser_id", sa.String(50), nullable=False, unique=True),
|
||||||
sa.Column("model", sa.String(50), nullable=False),
|
sa.Column("brand", sa.String(50), nullable=False),
|
||||||
sa.Column("year", sa.Integer, nullable=True),
|
sa.Column("model", sa.String(50), nullable=False),
|
||||||
sa.Column("price", sa.BigInteger, nullable=True),
|
sa.Column("year", sa.Integer, nullable=True),
|
||||||
sa.Column("currency", sa.String(10), nullable=False, server_default="USD"),
|
sa.Column("price", sa.BigInteger, nullable=True),
|
||||||
sa.Column("mileage", sa.Integer, nullable=False, server_default="0"),
|
sa.Column("currency", sa.String(10), nullable=False, server_default="USD"),
|
||||||
sa.Column("country", sa.String(10), nullable=False, server_default="NA"),
|
sa.Column("mileage", sa.Integer, nullable=False, server_default="0"),
|
||||||
sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("country", sa.String(10), nullable=False, server_default="NA"),
|
||||||
sa.Column("color", sa.String, nullable=False, server_default="other"),
|
sa.Column("is_sold", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("drive", sa.String(10), nullable=True),
|
sa.Column("color", sa.String, nullable=False, server_default="other"),
|
||||||
sa.Column("gearbox", sa.String(10), nullable=True),
|
sa.Column("drive", sa.String(10), nullable=True),
|
||||||
sa.Column("steering_wheel", sa.String(10), nullable=True),
|
sa.Column("gearbox", sa.String(10), nullable=True),
|
||||||
sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"),
|
sa.Column("steering_wheel", sa.String(10), nullable=True),
|
||||||
sa.Column("engine_volume", sa.Integer, nullable=True),
|
sa.Column("body_type", sa.String(20), nullable=False, server_default="OTHER"),
|
||||||
sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"),
|
sa.Column("engine_volume", sa.Integer, nullable=True),
|
||||||
sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("selling_type", sa.String(20), nullable=False, server_default="NA"),
|
||||||
sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("one_owner", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("new_car", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("origin", sa.String(20), nullable=False, server_default="NA"),
|
sa.Column("is_hidden", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("origin_url", sa.String, nullable=False),
|
sa.Column("origin", sa.String(20), nullable=False, server_default="NA"),
|
||||||
sa.Column("origin_id", sa.String, nullable=False, unique=True),
|
sa.Column("origin_url", sa.String, nullable=False),
|
||||||
sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("origin_id", sa.String, nullable=False, unique=True),
|
||||||
sa.Column("evaluation", sa.String, nullable=True),
|
sa.Column("is_damaged", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")),
|
sa.Column("evaluation", sa.String, nullable=True),
|
||||||
sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("non_smoking", sa.Boolean, nullable=False, server_default=sa.text("true")),
|
||||||
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
sa.Column("rental", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("slug", sa.String, nullable=False),
|
sa.Column("repair_history", sa.Boolean, nullable=False, server_default=sa.text("false")),
|
||||||
sa.Column("last_seen_at", sa.DateTime(timezone=True), nullable=False, server_default=sa.func.now()),
|
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_table(
|
op.create_index("ix_iaai_cars_origin_url", "iaai_cars", ["origin_url"])
|
||||||
"images",
|
if not _index_exists("iaai_cars", "ix_iaai_cars_origin_id"):
|
||||||
sa.Column("id", sa.BigInteger, primary_key=True, autoincrement=True),
|
op.create_index("ix_iaai_cars_origin_id", "iaai_cars", ["origin_id"], unique=True)
|
||||||
sa.Column("fullres_image", sa.String, nullable=False),
|
|
||||||
sa.Column("preview_image", sa.String, nullable=False),
|
# Таблица iaai_images — изображения машин
|
||||||
sa.Column("order_index", sa.Integer, nullable=False),
|
if not _table_exists("iaai_images"):
|
||||||
sa.Column("car_id", sa.Integer, sa.ForeignKey("cars.id", ondelete="CASCADE"), nullable=False),
|
op.create_table(
|
||||||
)
|
"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("iaai_cars.id", ondelete="CASCADE"), nullable=False),
|
||||||
|
)
|
||||||
|
|
||||||
|
# Таблица iaai_sync_runs — логирование синхронизаций
|
||||||
|
if not _table_exists("iaai_sync_runs"):
|
||||||
|
op.create_table(
|
||||||
|
"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),
|
||||||
|
sa.Column("status", sa.Text, nullable=False),
|
||||||
|
sa.Column("lane", sa.Text, nullable=False),
|
||||||
|
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
|
||||||
|
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
|
||||||
|
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
|
||||||
|
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
|
||||||
|
sa.Column("error_summary", sa.Text, nullable=True),
|
||||||
|
)
|
||||||
|
|
||||||
# Таблица sync_runs — логирование синхронизаций
|
|
||||||
op.create_table(
|
|
||||||
"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),
|
|
||||||
sa.Column("status", sa.Text, nullable=False),
|
|
||||||
sa.Column("lane", sa.Text, nullable=False),
|
|
||||||
sa.Column("ids_fetched", sa.Integer, nullable=False, server_default="0"),
|
|
||||||
sa.Column("cars_upserted", sa.Integer, nullable=False, server_default="0"),
|
|
||||||
sa.Column("cars_failed", sa.Integer, nullable=False, server_default="0"),
|
|
||||||
sa.Column("images_upserted", sa.Integer, nullable=False, server_default="0"),
|
|
||||||
sa.Column("error_summary", sa.Text, nullable=True),
|
|
||||||
)
|
|
||||||
|
|
||||||
def downgrade() -> None:
|
def downgrade() -> None:
|
||||||
op.drop_table("sync_runs")
|
# Compatibility-only migration for shared production databases.
|
||||||
op.drop_table("images")
|
pass
|
||||||
op.drop_table("cars")
|
|
||||||
|
|||||||
@@ -10,20 +10,15 @@ depends_on: Union[str, Sequence[str], None] = None
|
|||||||
|
|
||||||
|
|
||||||
def upgrade() -> None:
|
def upgrade() -> None:
|
||||||
op.create_index("ix_cars_brand", "cars", ["brand"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand ON iaai_cars (brand)")
|
||||||
op.create_index("ix_cars_brand_model", "cars", ["brand", "model"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_brand_model ON iaai_cars (brand, model)")
|
||||||
op.create_index("ix_cars_year", "cars", ["year"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_year ON iaai_cars (year)")
|
||||||
op.create_index("ix_cars_is_sold", "cars", ["is_sold"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_is_sold ON iaai_cars (is_sold)")
|
||||||
op.create_index("ix_cars_last_seen_at", "cars", ["last_seen_at"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_cars_last_seen_at ON iaai_cars (last_seen_at)")
|
||||||
op.create_index("ix_images_car_id", "images", ["car_id"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id ON iaai_images (car_id)")
|
||||||
op.create_index("ix_sync_runs_status", "sync_runs", ["status"])
|
op.execute("CREATE INDEX IF NOT EXISTS ix_iaai_sync_runs_status ON iaai_sync_runs (status)")
|
||||||
|
|
||||||
|
|
||||||
def downgrade() -> None:
|
def downgrade() -> None:
|
||||||
op.drop_index("ix_sync_runs_status", table_name="sync_runs")
|
# Compatibility-only migration for shared production databases.
|
||||||
op.drop_index("ix_images_car_id", table_name="images")
|
pass
|
||||||
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")
|
|
||||||
|
|||||||
@@ -13,26 +13,24 @@ def upgrade() -> None:
|
|||||||
# Partial index для активных машин IAAI (is_sold = FALSE)
|
# Partial index для активных машин IAAI (is_sold = FALSE)
|
||||||
# Ускоряет поиск при mark_sold операциях
|
# Ускоряет поиск при mark_sold операциях
|
||||||
op.execute(
|
op.execute(
|
||||||
"CREATE INDEX IF NOT EXISTS ix_cars_active_iaai "
|
"CREATE INDEX IF NOT EXISTS ix_iaai_cars_active_iaai "
|
||||||
"ON cars (origin_url) "
|
"ON iaai_cars (origin_url) "
|
||||||
"WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'"
|
"WHERE is_sold = FALSE AND origin_id LIKE 'iaai:%'"
|
||||||
)
|
)
|
||||||
|
|
||||||
# Composite index для проверки дубликатов изображений
|
# Composite index для проверки дубликатов изображений
|
||||||
op.create_index(
|
op.execute(
|
||||||
"ix_images_car_id_fullres",
|
"CREATE INDEX IF NOT EXISTS ix_iaai_images_car_id_fullres "
|
||||||
"images",
|
"ON iaai_images (car_id, fullres_image)"
|
||||||
["car_id", "fullres_image"],
|
|
||||||
)
|
)
|
||||||
|
|
||||||
# Index для быстрого поиска existing машин по origin_url
|
# Index для быстрого поиска existing машин по origin_url
|
||||||
op.execute(
|
op.execute(
|
||||||
"CREATE INDEX IF NOT EXISTS ix_cars_origin_url_id "
|
"CREATE INDEX IF NOT EXISTS ix_iaai_cars_origin_url_id "
|
||||||
"ON cars (origin_url, origin_id)"
|
"ON iaai_cars (origin_url, origin_id)"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def downgrade() -> None:
|
def downgrade() -> None:
|
||||||
op.execute("DROP INDEX IF EXISTS ix_cars_origin_url_id")
|
# Compatibility-only migration for shared production databases.
|
||||||
op.drop_index("ix_images_car_id_fullres", table_name="images")
|
pass
|
||||||
op.execute("DROP INDEX IF EXISTS ix_cars_active_iaai")
|
|
||||||
|
|||||||
@@ -60,6 +60,11 @@ start_self_heal_watchdog_if_worker() {
|
|||||||
return 0
|
return 0
|
||||||
fi
|
fi
|
||||||
|
|
||||||
|
# После reboot/restart контейнера файловая система контейнера сохраняется,
|
||||||
|
# поэтому stale pidfile может остаться от прошлого запуска и блокировать старт Celery.
|
||||||
|
# Удаляем его перед запуском worker-процесса.
|
||||||
|
rm -f /tmp/celery-worker.pid
|
||||||
|
|
||||||
if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then
|
if [ "${IAAI_SELF_HEAL_ENABLED:-true}" = "false" ]; then
|
||||||
echo "[entrypoint] Self-heal watchdog disabled"
|
echo "[entrypoint] Self-heal watchdog disabled"
|
||||||
return 0
|
return 0
|
||||||
|
|||||||
@@ -11,6 +11,7 @@ from concurrent.futures import ThreadPoolExecutor, as_completed
|
|||||||
from concurrent.futures import TimeoutError as FuturesTimeoutError
|
from concurrent.futures import TimeoutError as FuturesTimeoutError
|
||||||
from datetime import datetime, timezone
|
from datetime import datetime, timezone
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
|
from threading import Event, Thread
|
||||||
from typing import Any, Callable
|
from typing import Any, Callable
|
||||||
from urllib.parse import urlsplit, urlunsplit
|
from urllib.parse import urlsplit, urlunsplit
|
||||||
from urllib.request import Request, urlopen
|
from urllib.request import Request, urlopen
|
||||||
@@ -44,6 +45,7 @@ _PAGE_COLLECT_TIMEOUT_S = 90
|
|||||||
_PAGE_NEXT_TIMEOUT_S = 60
|
_PAGE_NEXT_TIMEOUT_S = 60
|
||||||
# Ротация контекста выключена по умолчанию.
|
# Ротация контекста выключена по умолчанию.
|
||||||
_CONTEXT_ROTATE_EVERY_PAGES = int(os.environ.get("IAAI_CONTEXT_ROTATE_EVERY_PAGES", "0") or 0)
|
_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_LINKS = 120
|
||||||
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
|
_SMALL_SEGMENT_SUSPICIOUS_PAGINATION_MAX_PAGE = 2
|
||||||
|
|
||||||
@@ -97,6 +99,33 @@ class _PageOpWatchdog:
|
|||||||
return False
|
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:
|
class IAAIScraper:
|
||||||
|
|
||||||
@staticmethod
|
@staticmethod
|
||||||
@@ -453,16 +482,23 @@ class IAAIScraper:
|
|||||||
)
|
)
|
||||||
|
|
||||||
for expected_page in range(2, target_page_number + 1):
|
for expected_page in range(2, target_page_number + 1):
|
||||||
# Heartbeat для watchdog: при deep-resume (100+ кликов) задача может
|
self._report_progress(
|
||||||
# идти несколько минут без batch_upserted, поэтому пульсуем прогресс.
|
"listing_resume_progress",
|
||||||
if expected_page == 2 or expected_page == target_page_number or expected_page % 5 == 0:
|
current_page=expected_page - 1,
|
||||||
self._report_progress(
|
target_page=target_page_number,
|
||||||
"listing_resume_progress",
|
next_page_number=expected_page,
|
||||||
current_page=expected_page - 1,
|
)
|
||||||
target_page=target_page_number,
|
|
||||||
)
|
|
||||||
|
|
||||||
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(
|
self._report_progress(
|
||||||
"listing_resume_failed",
|
"listing_resume_failed",
|
||||||
current_page=expected_page - 1,
|
current_page=expected_page - 1,
|
||||||
@@ -471,6 +507,13 @@ class IAAIScraper:
|
|||||||
page.close()
|
page.close()
|
||||||
raise ListingResumeError(f"Failed to resume listing at page {target_page_number}")
|
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(
|
self._report_progress(
|
||||||
"listing_resume_completed",
|
"listing_resume_completed",
|
||||||
current_page=target_page_number,
|
current_page=target_page_number,
|
||||||
@@ -552,6 +595,7 @@ class IAAIScraper:
|
|||||||
failures: list[dict[str, str]] = []
|
failures: list[dict[str, str]] = []
|
||||||
early_stopped = False
|
early_stopped = False
|
||||||
truncated_by_time_budget = False
|
truncated_by_time_budget = False
|
||||||
|
listing_recovery_attempts = 0
|
||||||
|
|
||||||
known_origin_ids: set[str] | None = None
|
known_origin_ids: set[str] | None = None
|
||||||
threshold = self.settings.listing.early_stop_threshold
|
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)})
|
failures.append({"vehicle_url": f"batch_{batch_start}", "error": str(batch_exc)})
|
||||||
return True
|
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
|
page = None
|
||||||
try:
|
try:
|
||||||
page, applied_filters = self._open_listing_for_stream(
|
page, applied_filters = self._open_listing_for_stream(
|
||||||
@@ -721,6 +780,10 @@ class IAAIScraper:
|
|||||||
pass
|
pass
|
||||||
page = None
|
page = None
|
||||||
try:
|
try:
|
||||||
|
_consume_listing_recovery_attempt(
|
||||||
|
page_number=page_number,
|
||||||
|
reason="context_rotation",
|
||||||
|
)
|
||||||
page, _ = self._reopen_listing_and_resume(
|
page, _ = self._reopen_listing_and_resume(
|
||||||
target_page_number=page_number,
|
target_page_number=page_number,
|
||||||
make=make,
|
make=make,
|
||||||
@@ -752,6 +815,10 @@ class IAAIScraper:
|
|||||||
pass
|
pass
|
||||||
page = None
|
page = None
|
||||||
try:
|
try:
|
||||||
|
_consume_listing_recovery_attempt(
|
||||||
|
page_number=page_number,
|
||||||
|
reason=f"collect_failed:{type(op_exc).__name__}",
|
||||||
|
)
|
||||||
page, _ = self._reopen_listing_and_resume(
|
page, _ = self._reopen_listing_and_resume(
|
||||||
target_page_number=page_number,
|
target_page_number=page_number,
|
||||||
make=make,
|
make=make,
|
||||||
@@ -797,6 +864,10 @@ class IAAIScraper:
|
|||||||
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
|
logger.warning("Page %d still empty after retry — reopening listing and resuming", page_number)
|
||||||
page.close()
|
page.close()
|
||||||
try:
|
try:
|
||||||
|
_consume_listing_recovery_attempt(
|
||||||
|
page_number=page_number,
|
||||||
|
reason="empty_page_after_retry",
|
||||||
|
)
|
||||||
page, _ = self._reopen_listing_and_resume(
|
page, _ = self._reopen_listing_and_resume(
|
||||||
target_page_number=page_number,
|
target_page_number=page_number,
|
||||||
make=make,
|
make=make,
|
||||||
@@ -973,8 +1044,17 @@ class IAAIScraper:
|
|||||||
cars_failed=cars_failed,
|
cars_failed=cars_failed,
|
||||||
)
|
)
|
||||||
try:
|
try:
|
||||||
with _PageOpWatchdog(_PAGE_NEXT_TIMEOUT_S, f"next page {page_number + 1}"):
|
with _ProgressHeartbeat(
|
||||||
next_ok = self.listing_collector.go_to_next_page(page, expected_page_number=page_number + 1)
|
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:
|
except (PageOperationTimeoutError, PlaywrightError) as nav_exc:
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"go_to_next_page stalled/failed at page %d (%s) — will reopen",
|
"go_to_next_page stalled/failed at page %d (%s) — will reopen",
|
||||||
@@ -1018,6 +1098,10 @@ class IAAIScraper:
|
|||||||
pass
|
pass
|
||||||
page = None
|
page = None
|
||||||
try:
|
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(
|
page, _ = self._reopen_listing_and_resume(
|
||||||
target_page_number=page_number + 1,
|
target_page_number=page_number + 1,
|
||||||
make=make,
|
make=make,
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ CAR_DB_FIELDS = {
|
|||||||
}
|
}
|
||||||
|
|
||||||
_IN_CHUNK_SIZE = 5000
|
_IN_CHUNK_SIZE = 5000
|
||||||
|
CAR_TABLE_NAME = Car.__tablename__
|
||||||
|
|
||||||
|
|
||||||
class PersistenceService:
|
class PersistenceService:
|
||||||
@@ -532,7 +533,7 @@ class PersistenceService:
|
|||||||
if is_postgres:
|
if is_postgres:
|
||||||
# Считаем кандидатов на sold.
|
# Считаем кандидатов на sold.
|
||||||
total_active = session.execute(
|
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
|
).scalar() or 0
|
||||||
|
|
||||||
if total_active == 0:
|
if total_active == 0:
|
||||||
@@ -556,12 +557,12 @@ class PersistenceService:
|
|||||||
# Считаем будущие sold.
|
# Считаем будущие sold.
|
||||||
would_mark = session.execute(text("""
|
would_mark = session.execute(text("""
|
||||||
SELECT count(*)
|
SELECT count(*)
|
||||||
FROM cars c
|
FROM {car_table} c
|
||||||
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
||||||
WHERE a.url IS NULL
|
WHERE a.url IS NULL
|
||||||
AND c.is_sold = FALSE
|
AND c.is_sold = FALSE
|
||||||
AND c.origin_id LIKE 'iaai:%%'
|
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:
|
if total_active > 100 and would_mark > total_active * 0.8:
|
||||||
@@ -573,18 +574,18 @@ class PersistenceService:
|
|||||||
|
|
||||||
# Массовая пометка sold.
|
# Массовая пометка sold.
|
||||||
result = session.execute(text("""
|
result = session.execute(text("""
|
||||||
UPDATE cars
|
UPDATE {car_table}
|
||||||
SET is_sold = TRUE
|
SET is_sold = TRUE
|
||||||
FROM (
|
FROM (
|
||||||
SELECT c.id
|
SELECT c.id
|
||||||
FROM cars c
|
FROM {car_table} c
|
||||||
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
LEFT JOIN _active_urls a ON replace(c.origin_url, '~', '-') = a.url
|
||||||
WHERE a.url IS NULL
|
WHERE a.url IS NULL
|
||||||
AND c.is_sold = FALSE
|
AND c.is_sold = FALSE
|
||||||
AND c.origin_id LIKE 'iaai:%%'
|
AND c.origin_id LIKE 'iaai:%%'
|
||||||
) sub
|
) sub
|
||||||
WHERE cars.id = sub.id
|
WHERE {car_table}.id = sub.id
|
||||||
"""))
|
""".format(car_table=CAR_TABLE_NAME)))
|
||||||
count = result.rowcount or 0
|
count = result.rowcount or 0
|
||||||
else:
|
else:
|
||||||
# Упрощённый путь для SQLite.
|
# Упрощённый путь для SQLite.
|
||||||
|
|||||||
@@ -20,10 +20,10 @@ class Base(DeclarativeBase):
|
|||||||
|
|
||||||
|
|
||||||
class Car(Base):
|
class Car(Base):
|
||||||
__tablename__ = "cars"
|
__tablename__ = "iaai_cars"
|
||||||
__table_args__ = (
|
__table_args__ = (
|
||||||
Index("ix_cars_brand_model", "brand", "model"),
|
Index("ix_iaai_cars_brand_model", "brand", "model"),
|
||||||
Index("ix_cars_origin_id_not_sold", "origin_id", "is_sold"),
|
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)
|
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)
|
parser_id: Mapped[str] = mapped_column(String(50), nullable=False, unique=True)
|
||||||
@@ -59,17 +59,17 @@ class Car(Base):
|
|||||||
|
|
||||||
|
|
||||||
class Image(Base):
|
class Image(Base):
|
||||||
__tablename__ = "images"
|
__tablename__ = "iaai_images"
|
||||||
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
|
id: Mapped[int] = mapped_column(BigInteger().with_variant(Integer, "sqlite"), primary_key=True, autoincrement=True)
|
||||||
fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
|
fullres_image: Mapped[str] = mapped_column(String(), nullable=False)
|
||||||
preview_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)
|
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")
|
car: Mapped[Car] = relationship("Car", back_populates="images")
|
||||||
|
|
||||||
|
|
||||||
class SyncRun(Base):
|
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)
|
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())
|
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)
|
finished_at: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True)
|
||||||
|
|||||||
@@ -1,6 +1,7 @@
|
|||||||
# Инициализация Celery-приложения и периодических задач.
|
# Инициализация Celery-приложения и периодических задач.
|
||||||
|
|
||||||
import logging
|
import logging
|
||||||
|
import os
|
||||||
|
|
||||||
from celery import Celery
|
from celery import Celery
|
||||||
from celery.signals import worker_process_init, worker_ready, setup_logging as celery_setup_logging
|
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")
|
logger = logging.getLogger("iaai_scraper.worker.celery_app")
|
||||||
STARTUP_SYNC_DISPATCH_KEY = "iaai:state:startup_sync_dispatched"
|
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
|
@celery_setup_logging.connect
|
||||||
@@ -79,18 +86,19 @@ celery_app.conf.update(
|
|||||||
worker_hijack_root_logger=False,
|
worker_hijack_root_logger=False,
|
||||||
beat_schedule={
|
beat_schedule={
|
||||||
"periodic-sync-listing": {
|
"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,
|
"schedule": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||||
"args": (),
|
"args": (),
|
||||||
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
|
"kwargs": {"limit": settings.celery.beat_sync_limit, "only_new": True},
|
||||||
"options": {
|
"options": {
|
||||||
"queue": "scraping",
|
"queue": IAAI_SYNC_QUEUE,
|
||||||
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
"expires": settings.celery.beat_sync_interval_minutes * 60.0,
|
||||||
},
|
},
|
||||||
}
|
}
|
||||||
},
|
},
|
||||||
task_routes={
|
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
|
@worker_ready.connect
|
||||||
def _on_worker_ready(**kwargs):
|
def _on_worker_ready(**kwargs):
|
||||||
"""При старте worker отправляем первый sync_listing, если очередь пуста."""
|
"""При старте 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
|
redis_client = None
|
||||||
try:
|
try:
|
||||||
redis_client = Redis.from_url(
|
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)
|
logger.warning("Failed to clear stale lock %s on startup", stale_key, exc_info=True)
|
||||||
|
|
||||||
try:
|
try:
|
||||||
queue_len = int(redis_client.llen("scraping") or 0)
|
queue_len = int(redis_client.llen(IAAI_SYNC_QUEUE) or 0)
|
||||||
except Exception:
|
except Exception:
|
||||||
queue_len = 0
|
queue_len = 0
|
||||||
if 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
|
return
|
||||||
|
|
||||||
should_dispatch = bool(redis_client.set(STARTUP_SYNC_DISPATCH_KEY, "1", nx=True, ex=600))
|
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")
|
logger.info("Worker ready — dispatching initial sync_listing task")
|
||||||
celery_app.send_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},
|
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,
|
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"
|
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||||
SELF_HEAL_RESTART_LOCK_KEY = "iaai:state:self_heal_restart_in_progress"
|
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:
|
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
|
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:
|
def _kill_worker_process() -> None:
|
||||||
pid_file = "/tmp/celery-worker.pid"
|
pid_file = "/tmp/celery-worker.pid"
|
||||||
pid: int | None = None
|
pid: int | None = None
|
||||||
@@ -101,7 +119,7 @@ def _kill_worker_process() -> None:
|
|||||||
# Если процесс ещё жив — принудительно убиваем.
|
# Если процесс ещё жив — принудительно убиваем.
|
||||||
os.kill(pid, 0)
|
os.kill(pid, 0)
|
||||||
logger.error("Self-heal: worker pid=%s did not stop after SIGTERM; sending SIGKILL", pid)
|
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:
|
except ProcessLookupError:
|
||||||
pass
|
pass
|
||||||
except Exception:
|
except Exception:
|
||||||
@@ -137,8 +155,8 @@ def main() -> None:
|
|||||||
redis_client = _get_redis()
|
redis_client = _get_redis()
|
||||||
redis_client.ping()
|
redis_client.ping()
|
||||||
|
|
||||||
queue_len = _safe_int(redis_client.llen(queue_name), 0)
|
has_inflight, inflight = _has_inflight_work(redis_client, queue_name)
|
||||||
if queue_len <= 0:
|
if not has_inflight:
|
||||||
time.sleep(check_interval)
|
time.sleep(check_interval)
|
||||||
continue
|
continue
|
||||||
|
|
||||||
@@ -151,8 +169,8 @@ def main() -> None:
|
|||||||
time.sleep(check_interval)
|
time.sleep(check_interval)
|
||||||
continue
|
continue
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Self-heal: queue=%d but no progress timestamp found after startup grace",
|
"Self-heal: inflight=%s but no progress timestamp found after startup grace",
|
||||||
queue_len,
|
inflight,
|
||||||
)
|
)
|
||||||
age = stall_seconds + 1
|
age = stall_seconds + 1
|
||||||
|
|
||||||
@@ -174,8 +192,8 @@ def main() -> None:
|
|||||||
continue
|
continue
|
||||||
|
|
||||||
logger.error(
|
logger.error(
|
||||||
"Self-heal: detected global stall (queue=%d, progress_age=%ss > %ss). Restarting worker process...",
|
"Self-heal: detected global stall (inflight=%s, progress_age=%ss > %ss). Restarting worker process...",
|
||||||
queue_len,
|
inflight,
|
||||||
age,
|
age,
|
||||||
stall_seconds,
|
stall_seconds,
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
# Задачи Celery.
|
# Задачи Celery для синхронизации автомобилей и листинга IAAI.
|
||||||
|
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
@@ -7,7 +7,6 @@ import signal
|
|||||||
from threading import Event, Thread
|
from threading import Event, Thread
|
||||||
import time
|
import time
|
||||||
import uuid
|
import uuid
|
||||||
from typing import Callable
|
|
||||||
|
|
||||||
from billiard.exceptions import SoftTimeLimitExceeded
|
from billiard.exceptions import SoftTimeLimitExceeded
|
||||||
from celery import shared_task
|
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")
|
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)
|
_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)
|
_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_LISTING_LOCK_KEY = "iaai:locks:sync_listing"
|
||||||
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
SYNC_FULL_SCAN_DONE_KEY = "iaai:state:sync_full_scan_done"
|
||||||
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
SYNC_LISTING_CHECKPOINT_KEY = "iaai:state:sync_listing_checkpoint"
|
||||||
SYNC_LISTING_CHECKPOINT_TTL_SECONDS = 7 * 24 * 60 * 60
|
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_KEY = "iaai:state:sync_listing_bootstrap_failure_streak"
|
||||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
|
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_LIMIT = 3
|
||||||
SYNC_LISTING_BOOTSTRAP_FAILURE_STREAK_TTL_SECONDS = 24 * 60 * 60
|
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_KEY = "iaai:state:hourly_failure_streak"
|
||||||
HOURLY_FAILURE_STREAK_LIMIT = 3
|
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_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_SEGMENT_LOCK_KEY_FMT = "iaai:locks:sync_segment:{idx}"
|
||||||
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
|
SYNC_SEGMENTS_PROGRESS_KEY = "iaai:state:sync_segments_progress"
|
||||||
SYNC_SEGMENTS_TOTAL_KEY = "iaai:state:sync_segments_total"
|
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
|
SYNC_SEGMENTS_PROGRESS_TTL_SECONDS = 24 * 60 * 60
|
||||||
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
TASK_PROGRESS_KEY_FMT = "iaai:state:task_progress:{task_id}"
|
||||||
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
GLOBAL_PROGRESS_TS_KEY = "iaai:state:last_progress_ts"
|
||||||
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
SITEMAP_HOURLY_LAST_COUNT_KEY = "iaai:state:sitemap_hourly_last_count"
|
||||||
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
SITEMAP_HOURLY_REFRESH_OFFSET_KEY = "iaai:state:sitemap_hourly_refresh_offset"
|
||||||
SYNC_LISTING_STALLED_SEGMENT_KEY = "iaai:state:sync_listing_stalled_segment"
|
STALL_WATCHDOG_NAVIGATION_STAGES = {
|
||||||
SYNC_LISTING_STALLED_SEGMENT_TTL_SECONDS = 24 * 60 * 60
|
"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):
|
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):
|
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)
|
return func(*args, **kwargs)
|
||||||
|
|
||||||
|
|
||||||
@@ -104,7 +99,8 @@ def _sync_listing_lock_ttl_seconds() -> int:
|
|||||||
settings = Settings()
|
settings = Settings()
|
||||||
soft = settings.celery.task_soft_time_limit
|
soft = settings.celery.task_soft_time_limit
|
||||||
hard = settings.celery.task_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
|
effective_hard = min(hard, soft + 120) if soft else hard
|
||||||
return max(effective_hard + 120, 300)
|
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)
|
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]:
|
def _hourly_sitemap_diff_sync(*, lane: str, limit: int | None, only_new: bool | None, progress_callback=None) -> dict[str, object]:
|
||||||
del limit, only_new
|
del limit, only_new
|
||||||
settings = Settings()
|
settings = Settings()
|
||||||
@@ -405,7 +407,6 @@ def _start_stall_watchdog(
|
|||||||
stall_timeout_seconds: int,
|
stall_timeout_seconds: int,
|
||||||
lock_key: str | None = None,
|
lock_key: str | None = None,
|
||||||
lock_owner: str | None = None,
|
lock_owner: str | None = None,
|
||||||
on_stall: Callable[[dict[str, object]], None] | None = None,
|
|
||||||
) -> tuple[Event, Thread]:
|
) -> tuple[Event, Thread]:
|
||||||
stop_event = Event()
|
stop_event = Event()
|
||||||
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
interval_seconds = max(5.0, min(30.0, stall_timeout_seconds / 3))
|
||||||
@@ -438,25 +439,22 @@ def _start_stall_watchdog(
|
|||||||
else:
|
else:
|
||||||
no_data_count = 0
|
no_data_count = 0
|
||||||
data = json.loads(raw)
|
data = json.loads(raw)
|
||||||
|
stage = data.get("stage")
|
||||||
last_ts = int(data.get("ts") or 0)
|
last_ts = int(data.get("ts") or 0)
|
||||||
if not last_ts:
|
if not last_ts:
|
||||||
continue
|
continue
|
||||||
|
effective_stall_timeout = _stall_timeout_for_progress(stage, stall_timeout_seconds)
|
||||||
age = int(time.time()) - last_ts
|
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
|
watchdog_born = time.monotonic() # reset absolute deadline on real progress
|
||||||
continue
|
continue
|
||||||
logger.error(
|
logger.error(
|
||||||
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
"Task %s stalled for %ss at stage=%s payload=%s; cleaning up and restarting",
|
||||||
task_id,
|
task_id,
|
||||||
age,
|
age,
|
||||||
data.get("stage"),
|
stage,
|
||||||
data,
|
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:
|
except Exception:
|
||||||
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
logger.warning("Failed to inspect task progress for stall watchdog", exc_info=True)
|
||||||
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
# Если Redis тоже не отвечает дольше дедлайна — убиваем.
|
||||||
@@ -486,7 +484,7 @@ def _start_stall_watchdog(
|
|||||||
celery_app.send_task(
|
celery_app.send_task(
|
||||||
SYNC_LISTING_TASK_NAME,
|
SYNC_LISTING_TASK_NAME,
|
||||||
kwargs={},
|
kwargs={},
|
||||||
queue="scraping",
|
queue=IAAI_SYNC_QUEUE,
|
||||||
countdown=15,
|
countdown=15,
|
||||||
expires=followup_ttl,
|
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)
|
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:
|
def _try_set_followup_pending(redis_client: Redis, *, ttl_seconds: int) -> bool:
|
||||||
try:
|
try:
|
||||||
return bool(redis_client.set(SYNC_LISTING_FOLLOWUP_PENDING_KEY, "1", nx=True, ex=max(60, int(ttl_seconds))))
|
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)
|
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]:
|
def _mark_segment_completed(redis_client: Redis, segment_index: int) -> tuple[int, int]:
|
||||||
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
|
"""Помечает сегмент завершённым. Возвращает (completed_count, total)."""
|
||||||
try:
|
try:
|
||||||
pipe = redis_client.pipeline()
|
pipe = redis_client.pipeline()
|
||||||
pipe.sadd(SYNC_SEGMENTS_PROGRESS_KEY, str(int(segment_index)))
|
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_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.scard(SYNC_SEGMENTS_PROGRESS_KEY)
|
||||||
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
|
pipe.get(SYNC_SEGMENTS_TOTAL_KEY)
|
||||||
results = pipe.execute()
|
results = pipe.execute()
|
||||||
completed = int(results[3] or 0)
|
completed = int(results[2] or 0)
|
||||||
total = int(results[4] or 0) if results[4] else 0
|
total = int(results[3] or 0) if results[3] else 0
|
||||||
return completed, total
|
return completed, total
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to mark segment %d completed", segment_index, exc_info=True)
|
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(
|
@shared_task(
|
||||||
name="iaai_scraper.worker.tasks.sync_segment_task",
|
name="iaai_scraper.worker.tasks.sync_segment_task",
|
||||||
|
queue=IAAI_SYNC_QUEUE,
|
||||||
bind=True,
|
bind=True,
|
||||||
max_retries=2,
|
max_retries=2,
|
||||||
default_retry_delay=60,
|
default_retry_delay=60,
|
||||||
@@ -1274,6 +1079,7 @@ def sync_segment_task(
|
|||||||
|
|
||||||
@shared_task(
|
@shared_task(
|
||||||
name="iaai_scraper.worker.tasks.sync_vehicle_task",
|
name="iaai_scraper.worker.tasks.sync_vehicle_task",
|
||||||
|
queue=IAAI_SYNC_QUEUE,
|
||||||
bind=True,
|
bind=True,
|
||||||
max_retries=2,
|
max_retries=2,
|
||||||
default_retry_delay=30,
|
default_retry_delay=30,
|
||||||
@@ -1304,7 +1110,8 @@ def sync_vehicle_task(self, vehicle_url: str, lane: str = "iaai"):
|
|||||||
|
|
||||||
|
|
||||||
@shared_task(
|
@shared_task(
|
||||||
name="iaai_scraper.worker.tasks.sync_listing_task",
|
name=SYNC_LISTING_TASK_NAME,
|
||||||
|
queue=IAAI_SYNC_QUEUE,
|
||||||
bind=True,
|
bind=True,
|
||||||
max_retries=3,
|
max_retries=3,
|
||||||
default_retry_delay=120,
|
default_retry_delay=120,
|
||||||
@@ -1376,7 +1183,7 @@ def sync_listing_task(
|
|||||||
try:
|
try:
|
||||||
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
|
followup_expires = max(int(lock_ttl), int(delay_seconds) + 300)
|
||||||
self.app.send_task(
|
self.app.send_task(
|
||||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
SYNC_LISTING_TASK_NAME,
|
||||||
kwargs={
|
kwargs={
|
||||||
"make": make,
|
"make": make,
|
||||||
"model": model,
|
"model": model,
|
||||||
@@ -1384,7 +1191,7 @@ def sync_listing_task(
|
|||||||
"limit": limit,
|
"limit": limit,
|
||||||
"only_new": only_new,
|
"only_new": only_new,
|
||||||
},
|
},
|
||||||
queue="scraping",
|
queue=IAAI_SYNC_QUEUE,
|
||||||
countdown=max(0, int(delay_seconds)),
|
countdown=max(0, int(delay_seconds)),
|
||||||
expires=followup_expires,
|
expires=followup_expires,
|
||||||
)
|
)
|
||||||
@@ -1436,19 +1243,6 @@ def sync_listing_task(
|
|||||||
stall_timeout_seconds=stall_timeout,
|
stall_timeout_seconds=stall_timeout,
|
||||||
lock_key=SYNC_LISTING_LOCK_KEY,
|
lock_key=SYNC_LISTING_LOCK_KEY,
|
||||||
lock_owner=owner_token,
|
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(
|
_update_task_progress(
|
||||||
redis_client,
|
redis_client,
|
||||||
@@ -1484,11 +1278,9 @@ def sync_listing_task(
|
|||||||
last_completed_segment: int | None = None
|
last_completed_segment: int | None = None
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
last_completed_segment = _load_last_completed_segment(redis_client)
|
last_completed_segment = _load_last_completed_segment(redis_client)
|
||||||
stalled_segment = _load_stalled_segment(redis_client)
|
|
||||||
else:
|
else:
|
||||||
# После завершения bootstrap чекпоинт не нужен никогда.
|
# После завершения bootstrap чекпоинт не нужен никогда.
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
stalled_segment = None
|
|
||||||
|
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
logger.info(
|
logger.info(
|
||||||
@@ -1649,17 +1441,12 @@ def sync_listing_task(
|
|||||||
if use_segmented and settings.celery.parallel_segments:
|
if use_segmented and settings.celery.parallel_segments:
|
||||||
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
|
# Сегменты уже завершённые (для bootstrap resume) пропускаем по Redis SET.
|
||||||
already_completed: set[int] = set()
|
already_completed: set[int] = set()
|
||||||
cycle_id: str | None = None
|
|
||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
try:
|
try:
|
||||||
cycle_id, created_new_cycle = _ensure_segments_cycle(redis_client, len(segments))
|
|
||||||
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
|
raw = redis_client.smembers(SYNC_SEGMENTS_PROGRESS_KEY) or set()
|
||||||
already_completed = {int(x) for x in raw if str(x).strip().lstrip("-").isdigit()}
|
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:
|
except Exception:
|
||||||
already_completed = set()
|
already_completed = set()
|
||||||
cycle_id = None
|
|
||||||
|
|
||||||
pending = [
|
pending = [
|
||||||
(idx, seg) for idx, seg in enumerate(segments)
|
(idx, seg) for idx, seg in enumerate(segments)
|
||||||
@@ -1671,7 +1458,6 @@ def sync_listing_task(
|
|||||||
if force_bootstrap_full_scan:
|
if force_bootstrap_full_scan:
|
||||||
_set_full_scan_done(redis_client, True)
|
_set_full_scan_done(redis_client, True)
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
_clear_segments_progress_state(redis_client)
|
|
||||||
_clear_bootstrap_failure_streak(redis_client)
|
_clear_bootstrap_failure_streak(redis_client)
|
||||||
logger.warning("Parallel segments: nothing to dispatch (all completed)")
|
logger.warning("Parallel segments: nothing to dispatch (all completed)")
|
||||||
return {
|
return {
|
||||||
@@ -1683,6 +1469,10 @@ def sync_listing_task(
|
|||||||
"segments_already_completed": len(already_completed),
|
"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
|
dispatched = 0
|
||||||
for idx, seg in pending:
|
for idx, seg in pending:
|
||||||
try:
|
try:
|
||||||
@@ -1695,15 +1485,15 @@ def sync_listing_task(
|
|||||||
"only_new": effective_only_new,
|
"only_new": effective_only_new,
|
||||||
"is_bootstrap": force_bootstrap_full_scan,
|
"is_bootstrap": force_bootstrap_full_scan,
|
||||||
},
|
},
|
||||||
queue="scraping",
|
queue=IAAI_SYNC_QUEUE,
|
||||||
)
|
)
|
||||||
dispatched += 1
|
dispatched += 1
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
|
logger.warning("Failed to dispatch segment %d", idx, exc_info=True)
|
||||||
|
|
||||||
logger.warning(
|
logger.warning(
|
||||||
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s, cycle=%s)",
|
"Parallel segments dispatched: %d/%d (already_completed=%d, bootstrap=%s)",
|
||||||
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan, cycle_id,
|
dispatched, len(segments), len(already_completed), force_bootstrap_full_scan,
|
||||||
)
|
)
|
||||||
return {
|
return {
|
||||||
"status": "success",
|
"status": "success",
|
||||||
@@ -1712,11 +1502,9 @@ def sync_listing_task(
|
|||||||
"segments_total": len(segments),
|
"segments_total": len(segments),
|
||||||
"segments_dispatched": dispatched,
|
"segments_dispatched": dispatched,
|
||||||
"segments_already_completed": len(already_completed),
|
"segments_already_completed": len(already_completed),
|
||||||
"segments_cycle_id": cycle_id,
|
|
||||||
}
|
}
|
||||||
# --- конец параллельной ветки ---
|
# --- конец параллельной ветки ---
|
||||||
|
|
||||||
precomputed_result = None
|
|
||||||
resume_from_segment = 0
|
resume_from_segment = 0
|
||||||
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
if force_bootstrap_full_scan and use_segmented and last_completed_segment is not None:
|
||||||
resume_from_segment = max(0, last_completed_segment + 1)
|
resume_from_segment = max(0, last_completed_segment + 1)
|
||||||
@@ -1734,117 +1522,6 @@ def sync_listing_task(
|
|||||||
resume_from_segment, last_completed_segment,
|
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):
|
def _progress_cb_main(stage, meta):
|
||||||
_update_task_progress(
|
_update_task_progress(
|
||||||
redis_client,
|
redis_client,
|
||||||
@@ -1877,16 +1554,13 @@ def sync_listing_task(
|
|||||||
only_new=effective_only_new,
|
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:
|
if force_bootstrap_full_scan:
|
||||||
bootstrap_completed = bool(result.get("full_scan_completed"))
|
bootstrap_completed = bool(result.get("full_scan_completed"))
|
||||||
if bootstrap_completed:
|
if bootstrap_completed:
|
||||||
_set_full_scan_done(redis_client, False if always_full_scan else True)
|
_set_full_scan_done(redis_client, False if always_full_scan else True)
|
||||||
_clear_sync_checkpoint(redis_client)
|
_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_failure_streak(redis_client)
|
||||||
_clear_bootstrap_continuation_streak(redis_client)
|
_clear_bootstrap_continuation_streak(redis_client)
|
||||||
if always_full_scan:
|
if always_full_scan:
|
||||||
@@ -1902,7 +1576,7 @@ def sync_listing_task(
|
|||||||
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
|
for key in ("cars_upserted", "images_upserted", "skipped_existing", "total_discovered")
|
||||||
) or int(listing_payload.get("vehicles_collected") or 0) > 0
|
) 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)
|
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"
|
followup_reason = "bootstrap_not_completed"
|
||||||
if anti_bot_detected:
|
if anti_bot_detected:
|
||||||
followup_delay = 180
|
followup_delay = 180
|
||||||
@@ -1922,7 +1596,6 @@ def sync_listing_task(
|
|||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
_clear_sync_checkpoint(redis_client)
|
_clear_sync_checkpoint(redis_client)
|
||||||
_clear_bootstrap_resume_target(redis_client)
|
|
||||||
|
|
||||||
summary = {
|
summary = {
|
||||||
"task_id": task_id,
|
"task_id": task_id,
|
||||||
@@ -1986,7 +1659,7 @@ def sync_listing_task(
|
|||||||
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
|
if _try_set_followup_pending(redis_client, ttl_seconds=followup_ttl):
|
||||||
try:
|
try:
|
||||||
self.app.send_task(
|
self.app.send_task(
|
||||||
"iaai_scraper.worker.tasks.sync_listing_task",
|
SYNC_LISTING_TASK_NAME,
|
||||||
kwargs={
|
kwargs={
|
||||||
"make": make,
|
"make": make,
|
||||||
"model": model,
|
"model": model,
|
||||||
@@ -1994,7 +1667,7 @@ def sync_listing_task(
|
|||||||
"limit": limit,
|
"limit": limit,
|
||||||
"only_new": only_new,
|
"only_new": only_new,
|
||||||
},
|
},
|
||||||
queue="scraping",
|
queue=IAAI_SYNC_QUEUE,
|
||||||
countdown=10,
|
countdown=10,
|
||||||
expires=followup_ttl,
|
expires=followup_ttl,
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user