Compare commits
3 Commits
eba3750f15
...
ec49ed9bc5
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
ec49ed9bc5 | ||
|
|
bc5c3698c3 | ||
|
|
12b313ad52 |
@@ -26,21 +26,28 @@ def list_cars(
|
|||||||
# Список автомобилей с пагинацией и фильтрами
|
# Список автомобилей с пагинацией и фильтрами
|
||||||
with persistence.session_scope() as session:
|
with persistence.session_scope() as session:
|
||||||
query = select(Car)
|
query = select(Car)
|
||||||
|
count_query = select(func.count(Car.id))
|
||||||
|
|
||||||
if brand:
|
if brand:
|
||||||
escaped_brand = brand.replace("%", r"\%").replace("_", r"\_")
|
escaped_brand = brand.replace("%", r"\%").replace("_", r"\_")
|
||||||
query = query.where(Car.brand.ilike(f"%{escaped_brand}%", escape="\\"))
|
cond = Car.brand.ilike(f"%{escaped_brand}%", escape="\\")
|
||||||
|
query = query.where(cond)
|
||||||
|
count_query = count_query.where(cond)
|
||||||
if model:
|
if model:
|
||||||
escaped_model = model.replace("%", r"\%").replace("_", r"\_")
|
escaped_model = model.replace("%", r"\%").replace("_", r"\_")
|
||||||
query = query.where(Car.model.ilike(f"%{escaped_model}%", escape="\\"))
|
cond = Car.model.ilike(f"%{escaped_model}%", escape="\\")
|
||||||
|
query = query.where(cond)
|
||||||
|
count_query = count_query.where(cond)
|
||||||
if year_min is not None:
|
if year_min is not None:
|
||||||
query = query.where(Car.year >= year_min)
|
query = query.where(Car.year >= year_min)
|
||||||
|
count_query = count_query.where(Car.year >= year_min)
|
||||||
if year_max is not None:
|
if year_max is not None:
|
||||||
query = query.where(Car.year <= year_max)
|
query = query.where(Car.year <= year_max)
|
||||||
|
count_query = count_query.where(Car.year <= year_max)
|
||||||
if is_sold is not None:
|
if is_sold is not None:
|
||||||
query = query.where(Car.is_sold == is_sold)
|
query = query.where(Car.is_sold == is_sold)
|
||||||
|
count_query = count_query.where(Car.is_sold == is_sold)
|
||||||
|
|
||||||
count_query = select(func.count()).select_from(query.subquery())
|
|
||||||
total = session.execute(count_query).scalar() or 0
|
total = session.execute(count_query).scalar() or 0
|
||||||
|
|
||||||
offset = (page - 1) * per_page
|
offset = (page - 1) * per_page
|
||||||
|
|||||||
@@ -919,6 +919,41 @@ class EncarScraper:
|
|||||||
excluded_brands: set[str] | None = None,
|
excluded_brands: set[str] | None = None,
|
||||||
runtime_filters: FiltersConfig | None = None,
|
runtime_filters: FiltersConfig | None = None,
|
||||||
probe_all_photos: bool = False,
|
probe_all_photos: bool = False,
|
||||||
|
) -> dict[str, Any]:
|
||||||
|
"""Публичная обёртка: гарантирует закрытие HTTP pool даже при ошибках."""
|
||||||
|
try:
|
||||||
|
return self._sync_listing_impl(
|
||||||
|
limit=limit,
|
||||||
|
filters=filters,
|
||||||
|
lane=lane,
|
||||||
|
only_new=only_new,
|
||||||
|
batch_size=batch_size,
|
||||||
|
redis_client=redis_client,
|
||||||
|
allowed_brands=allowed_brands,
|
||||||
|
excluded_brands=excluded_brands,
|
||||||
|
runtime_filters=runtime_filters,
|
||||||
|
probe_all_photos=probe_all_photos,
|
||||||
|
)
|
||||||
|
finally:
|
||||||
|
if self._batch_pool is not None:
|
||||||
|
try:
|
||||||
|
self._batch_pool.close()
|
||||||
|
except Exception:
|
||||||
|
logger.debug("Failed to close batch pool", exc_info=True)
|
||||||
|
self._batch_pool = None
|
||||||
|
|
||||||
|
def _sync_listing_impl(
|
||||||
|
self,
|
||||||
|
limit: int | None = None,
|
||||||
|
filters: EncarFilters | None = None,
|
||||||
|
lane: str = "encar",
|
||||||
|
only_new: bool = False,
|
||||||
|
batch_size: int = 1000,
|
||||||
|
redis_client: Any | None = None,
|
||||||
|
allowed_brands: set[str] | None = None,
|
||||||
|
excluded_brands: set[str] | None = None,
|
||||||
|
runtime_filters: FiltersConfig | None = None,
|
||||||
|
probe_all_photos: bool = False,
|
||||||
) -> dict[str, Any]:
|
) -> dict[str, Any]:
|
||||||
"""Полная синхронизация листинга Encar.
|
"""Полная синхронизация листинга Encar.
|
||||||
|
|
||||||
@@ -988,11 +1023,21 @@ class EncarScraper:
|
|||||||
max_consecutive_errors = 5
|
max_consecutive_errors = 5
|
||||||
consecutive_zero_new = 0
|
consecutive_zero_new = 0
|
||||||
max_consecutive_zero_new = 5
|
max_consecutive_zero_new = 5
|
||||||
|
# Защита от бесконечной пагинации: Encar API отдаёт максимум
|
||||||
|
# ~10k записей на query, шарды строятся ≤9500 → не более ~10
|
||||||
|
# страниц при page_size=1000. Жёсткий cap в 50 — паранойя.
|
||||||
|
max_pages_per_shard = 50
|
||||||
|
|
||||||
shard_query = shard_filters.build_query()
|
shard_query = shard_filters.build_query()
|
||||||
shard_collected = 0
|
shard_collected = 0
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
|
if page >= max_pages_per_shard:
|
||||||
|
logger.warning(
|
||||||
|
"Shard %d hit hard page cap (%d), moving to next",
|
||||||
|
shard_idx, max_pages_per_shard,
|
||||||
|
)
|
||||||
|
break
|
||||||
offset = page * page_size
|
offset = page * page_size
|
||||||
try:
|
try:
|
||||||
response = self._fetch_listing_page(
|
response = self._fetch_listing_page(
|
||||||
@@ -1152,11 +1197,6 @@ class EncarScraper:
|
|||||||
|
|
||||||
self._clear_checkpoint(redis_client)
|
self._clear_checkpoint(redis_client)
|
||||||
|
|
||||||
# Закрываем pool batch API
|
|
||||||
if self._batch_pool is not None:
|
|
||||||
self._batch_pool.close()
|
|
||||||
self._batch_pool = None
|
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"Full sync complete: %d shards, %d collected, %d synced, %d failed, %d marked sold",
|
"Full sync complete: %d shards, %d collected, %d synced, %d failed, %d marked sold",
|
||||||
len(shards), items_collected, synced, failed, marked_sold,
|
len(shards), items_collected, synced, failed, marked_sold,
|
||||||
|
|||||||
@@ -46,7 +46,10 @@ class PersistenceService:
|
|||||||
try:
|
try:
|
||||||
Base.metadata.create_all(self.engine)
|
Base.metadata.create_all(self.engine)
|
||||||
except Exception:
|
except Exception:
|
||||||
logger.debug("create_tables skipped (schema already exists)")
|
# Не глотаем: таблицы могут уже существовать (ок), либо есть проблема
|
||||||
|
# доступа — пусть вышестоящий код видит её при первой операции, но
|
||||||
|
# сигнализируем в warning чтобы упростить диагностику.
|
||||||
|
logger.warning("create_tables failed (continuing — schema may already exist)", exc_info=True)
|
||||||
|
|
||||||
@contextmanager
|
@contextmanager
|
||||||
def session_scope(self) -> Iterator[Session]:
|
def session_scope(self) -> Iterator[Session]:
|
||||||
|
|||||||
Reference in New Issue
Block a user