Some checks failed
CI / changes (push) Successful in 7s
CI / frontend-tests (push) Has been skipped
CI / changes (pull_request) Successful in 10s
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (push) Successful in 1m45s
CI / openapi-codegen-check (pull_request) Successful in 1m43s
CI / backend-tests (push) Failing after 8m36s
CI / backend-tests (pull_request) Failing after 8m36s
Имплементация фиксов 2-го аудита backend/app/** (после merge #1543). Воркер на файл, точечные правки. Верификация: py_compile 58/58 .py. Полностью исправлено (82). Оставлены открытыми (13): partial/needs-cross-file/needs-leha — #1569, #1590, #1593, #1606, #1609, #1617, #1633, #1635, #1637, #1638, #1640, #1642, #1650. Closes #1560 Closes #1561 Closes #1562 Closes #1563 Closes #1564 Closes #1565 Closes #1566 Closes #1567 Closes #1570 Closes #1571 Closes #1572 Closes #1573 Closes #1574 Closes #1576 Closes #1577 Closes #1578 Closes #1579 Closes #1580 Closes #1581 Closes #1582 Closes #1583 Closes #1584 Closes #1585 Closes #1586 Closes #1587 Closes #1588 Closes #1589 Closes #1591 Closes #1592 Closes #1594 Closes #1595 Closes #1596 Closes #1597 Closes #1598 Closes #1599 Closes #1600 Closes #1601 Closes #1602 Closes #1603 Closes #1604 Closes #1605 Closes #1607 Closes #1608 Closes #1610 Closes #1611 Closes #1612 Closes #1613 Closes #1614 Closes #1615 Closes #1616 Closes #1618 Closes #1619 Closes #1620 Closes #1621 Closes #1622 Closes #1623 Closes #1624 Closes #1625 Closes #1626 Closes #1627 Closes #1628 Closes #1629 Closes #1630 Closes #1631 Closes #1632 Closes #1634 Closes #1636 Closes #1639 Closes #1641 Closes #1643 Closes #1644 Closes #1645 Closes #1646 Closes #1647 Closes #1648 Closes #1649 Closes #1651 Closes #1652 Closes #1653 Closes #1654 Closes #1655 Closes #1656
88 lines
3.6 KiB
Python
88 lines
3.6 KiB
Python
"""Refresh helper for mv_quarter_price_index (Issue #762).
|
|
|
|
mv_quarter_price_index depends on mv_quarter_price_per_m2, so both must be
|
|
refreshed in order:
|
|
1. mv_quarter_price_per_m2 (source MV — deals aggregation)
|
|
2. mv_quarter_price_index (derived MV — price index normalised to city median)
|
|
|
|
Scheduled via Celery beat hardcoded entry in workers/beat_schedule.py.
|
|
Cadence: monthly on the 5th at 05:00 MSK (after refresh_ekb_districts_medians
|
|
which runs at 04:00 MSK on the same day — deals data is already settled).
|
|
|
|
Usage example (manual, via psql-connected shell or admin endpoint):
|
|
from sqlalchemy.orm import Session
|
|
from app.services.site_finder.quarter_price_index_refresh import refresh_quarter_price_index
|
|
|
|
count = refresh_quarter_price_index(db)
|
|
# logs: "mv_quarter_price_index refreshed: 1972 rows"
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
|
|
from sqlalchemy import text
|
|
from sqlalchemy.exc import DatabaseError
|
|
from sqlalchemy.orm import Session
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
|
|
def _refresh_mv(db: Session, mv_name: str, *, concurrently: bool) -> None:
|
|
"""Run REFRESH MATERIALIZED VIEW [CONCURRENTLY] <mv_name>.
|
|
|
|
Falls back to non-concurrent on the known "cannot refresh concurrently"
|
|
error (MV empty or no UNIQUE index — should not happen in prod, but
|
|
provides a safe recovery path for first-run edge cases).
|
|
"""
|
|
try:
|
|
if concurrently:
|
|
db.execute(text(f"REFRESH MATERIALIZED VIEW CONCURRENTLY {mv_name}"))
|
|
else:
|
|
db.execute(text(f"REFRESH MATERIALIZED VIEW {mv_name}"))
|
|
db.commit()
|
|
except DatabaseError as e:
|
|
# PostgreSQL emits "CONCURRENTLY cannot be used when the materialized
|
|
# view ... is not populated" (matview.c, SQLSTATE 55000), which psycopg3
|
|
# surfaces as InternalError (a DatabaseError sibling of OperationalError).
|
|
if concurrently and "concurrently cannot be used" in str(e).lower():
|
|
logger.warning(
|
|
"%s: CONCURRENTLY failed (MV likely not populated), "
|
|
"falling back to non-concurrent refresh",
|
|
mv_name,
|
|
)
|
|
db.rollback()
|
|
db.execute(text(f"REFRESH MATERIALIZED VIEW {mv_name}"))
|
|
db.commit()
|
|
else:
|
|
raise
|
|
|
|
|
|
def refresh_quarter_price_index(db: Session, *, concurrently: bool = True) -> int:
|
|
"""Refresh mv_quarter_price_per_m2 then mv_quarter_price_index (in order).
|
|
|
|
mv_quarter_price_index reads from mv_quarter_price_per_m2, so the source MV
|
|
must be refreshed first. Both refreshes happen in the same call.
|
|
|
|
Args:
|
|
db: SQLAlchemy Session (sync).
|
|
concurrently: When True, uses REFRESH CONCURRENTLY for both MVs —
|
|
non-blocking (readers continue). Requires the respective UNIQUE
|
|
indexes (mv_quarter_price_pk on source, mv_quarter_price_index_uq
|
|
on derived) and both MVs to be already populated.
|
|
Pass False only for first populate or after MV recreation.
|
|
|
|
Returns:
|
|
Row count of mv_quarter_price_index after refresh (for observability).
|
|
"""
|
|
# Step 1: refresh source MV (deals aggregation layer)
|
|
_refresh_mv(db, "mv_quarter_price_per_m2", concurrently=concurrently)
|
|
logger.info("mv_quarter_price_per_m2 refreshed (chain step 1/2)")
|
|
|
|
# Step 2: refresh derived MV (price index, reads from source)
|
|
_refresh_mv(db, "mv_quarter_price_index", concurrently=concurrently)
|
|
|
|
row = db.execute(text("SELECT COUNT(*) FROM mv_quarter_price_index")).first()
|
|
count = int(row[0]) if row else 0
|
|
logger.info("mv_quarter_price_index refreshed: %d rows (chain step 2/2)", count)
|
|
return count
|