diff options
| author | Christian Cleberg <[email protected]> | 2026-04-11 22:56:22 -0500 |
|---|---|---|
| committer | Christian Cleberg <[email protected]> | 2026-04-11 22:56:22 -0500 |
| commit | 1c8d0bdd0a0a46d7400b0827aec88a787f5a04c3 (patch) | |
| tree | 516a47d260e9a1a003110dd5b20f26dffe013536 | |
| parent | 3489797460f9f8d19d37806da3b99c6b0baf4cf0 (diff) | |
| download | hutch-stats-1c8d0bdd0a0a46d7400b0827aec88a787f5a04c3.tar.gz hutch-stats-1c8d0bdd0a0a46d7400b0827aec88a787f5a04c3.tar.bz2 hutch-stats-1c8d0bdd0a0a46d7400b0827aec88a787f5a04c3.zip | |
Reduce SQLite lock contention during actor enqueue
| -rw-r--r-- | src/srht_contrib/config.py | 1 | ||||
| -rw-r--r-- | src/srht_contrib/db.py | 23 | ||||
| -rw-r--r-- | src/srht_contrib/scripts/enqueue_actors.py | 27 |
3 files changed, 39 insertions, 12 deletions
diff --git a/src/srht_contrib/config.py b/src/srht_contrib/config.py index 2cab56e..3a7af2c 100644 --- a/src/srht_contrib/config.py +++ b/src/srht_contrib/config.py @@ -37,6 +37,7 @@ class Settings(BaseSettings): poll_interval_seconds: int = Field(default=900, alias="POLL_INTERVAL_SECONDS") sync_overlap_hours: int = Field(default=1, alias="SYNC_OVERLAP_HOURS") srht_request_delay_seconds: float = Field(default=0.5, alias="SRHT_REQUEST_DELAY_SECONDS") + sqlite_busy_timeout_seconds: float = Field(default=30.0, alias="SQLITE_BUSY_TIMEOUT_SECONDS") discovery_batch_size: int = Field(default=5, alias="DISCOVERY_BATCH_SIZE") indexed_actor_repoll_seconds: int = Field(default=21600, alias="INDEXED_ACTOR_REPOLL_SECONDS") discovery_error_backoff_seconds: int = Field(default=3600, alias="DISCOVERY_ERROR_BACKOFF_SECONDS") diff --git a/src/srht_contrib/db.py b/src/srht_contrib/db.py index d96de57..7750df6 100644 --- a/src/srht_contrib/db.py +++ b/src/srht_contrib/db.py @@ -3,7 +3,7 @@ from __future__ import annotations from collections.abc import Generator from fastapi import HTTPException, Request, status -from sqlalchemy import Engine, create_engine, text +from sqlalchemy import Engine, create_engine, event, text from sqlalchemy.pool import StaticPool from sqlalchemy.orm import Session, declarative_base, sessionmaker @@ -13,11 +13,28 @@ Base = declarative_base() def make_engine(settings: Settings) -> Engine: - connect_args = {"check_same_thread": False} if settings.database_url.startswith("sqlite") else {} + connect_args = {} + if settings.database_url.startswith("sqlite"): + connect_args = { + "check_same_thread": False, + "timeout": settings.sqlite_busy_timeout_seconds, + } engine_kwargs = {"future": True, "connect_args": connect_args} if settings.database_url in {"sqlite://", "sqlite:///:memory:"}: engine_kwargs["poolclass"] = StaticPool - return create_engine(settings.database_url, **engine_kwargs) + engine = create_engine(settings.database_url, **engine_kwargs) + + if settings.database_url.startswith("sqlite"): + @event.listens_for(engine, "connect") + def _configure_sqlite(dbapi_connection, connection_record) -> None: # type: ignore[unused-ignore] + cursor = dbapi_connection.cursor() + cursor.execute(f"PRAGMA busy_timeout = {int(settings.sqlite_busy_timeout_seconds * 1000)}") + if settings.database_url not in {"sqlite://", "sqlite:///:memory:"}: + cursor.execute("PRAGMA journal_mode = WAL") + cursor.execute("PRAGMA synchronous = NORMAL") + cursor.close() + + return engine def make_session_factory(settings: Settings) -> sessionmaker[Session]: diff --git a/src/srht_contrib/scripts/enqueue_actors.py b/src/srht_contrib/scripts/enqueue_actors.py index 6735881..f6999af 100644 --- a/src/srht_contrib/scripts/enqueue_actors.py +++ b/src/srht_contrib/scripts/enqueue_actors.py @@ -34,6 +34,9 @@ def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: queued_at = start_at or datetime.now(tz=UTC) inserted = 0 + batch_size = 250 + queued_in_batch = 0 + with session_factory() as db: for index, actor in enumerate(usernames): next_poll_after = queued_at + timedelta(seconds=index * stagger_seconds) @@ -49,15 +52,21 @@ def enqueue_actors(username_file: Path, *, stagger_seconds: int = 300, start_at: ) db.add(tracked_actor) inserted += 1 - continue - - tracked_actor.is_active = True - if tracked_actor.queued_for_discovery_at is None: - tracked_actor.queued_for_discovery_at = queued_at - if tracked_actor.last_polled_at is None and tracked_actor.discovery_state != "indexed": - tracked_actor.discovery_state = "queued" - tracked_actor.next_poll_after = next_poll_after - db.commit() + else: + tracked_actor.is_active = True + if tracked_actor.queued_for_discovery_at is None: + tracked_actor.queued_for_discovery_at = queued_at + if tracked_actor.last_polled_at is None and tracked_actor.discovery_state != "indexed": + tracked_actor.discovery_state = "queued" + tracked_actor.next_poll_after = next_poll_after + + queued_in_batch += 1 + if queued_in_batch >= batch_size: + db.commit() + queued_in_batch = 0 + + if queued_in_batch: + db.commit() return inserted |
