|
11 | 11 | DEFAULT_OVERLAP = dt.timedelta(days=2) |
12 | 12 |
|
13 | 13 |
|
14 | | -def prune_pg(dsn: str, cutoff: dt.datetime) -> tuple[int, int]: |
15 | | - """Delete changesets and their changeset_stats older than cutoff; child rows first for the FK. The |
16 | | - DSN is interpolated into ATTACH, so it must be trusted.""" |
| 14 | +def _attach(dsn: str) -> duckdb.DuckDBPyConnection: |
17 | 15 | conn = duckdb.connect() |
18 | 16 | conn.execute("INSTALL postgres") |
19 | 17 | conn.execute("LOAD postgres") |
20 | 18 | conn.execute(f"ATTACH '{dsn.replace(chr(39), chr(39) * 2)}' AS pg (TYPE postgres)") |
| 19 | + return conn |
| 20 | + |
| 21 | + |
| 22 | +def _pg_execute(conn: duckdb.DuckDBPyConnection, sql: str) -> None: |
| 23 | + """Run one statement natively on the attached Postgres, so a bulk DELETE is a single indexed |
| 24 | + statement server-side instead of DuckDB's per-row ctid batches.""" |
| 25 | + conn.execute(f"CALL postgres_execute('pg', $osmsg_stmt${sql}$osmsg_stmt$)") |
| 26 | + |
| 27 | + |
| 28 | +def prune_pg(dsn: str, cutoff: dt.datetime) -> tuple[int, int]: |
| 29 | + """Delete changesets and their changeset_stats older than cutoff; child rows first for the FK. The |
| 30 | + DSN is interpolated into ATTACH, so it must be trusted. Counting and deleting use separate |
| 31 | + connections because a read pins the connection read-only, which would block the native deletes.""" |
21 | 32 | iso = cutoff.astimezone(dt.UTC).isoformat() |
22 | | - old_cs = f"SELECT changeset_id FROM pg.changesets WHERE created_at < TIMESTAMPTZ '{iso}'" |
23 | | - stats_row = conn.execute(f"SELECT count(*) FROM pg.changeset_stats WHERE changeset_id IN ({old_cs})").fetchone() |
24 | | - cs_row = conn.execute(f"SELECT count(*) FROM pg.changesets WHERE created_at < TIMESTAMPTZ '{iso}'").fetchone() |
| 33 | + older = f"created_at < TIMESTAMPTZ '{iso}'" |
| 34 | + |
| 35 | + reader = _attach(dsn) |
| 36 | + stats_row = reader.execute( |
| 37 | + "SELECT count(*) FROM pg.changeset_stats s " |
| 38 | + f"WHERE EXISTS (SELECT 1 FROM pg.changesets c WHERE c.changeset_id = s.changeset_id AND c.{older})" |
| 39 | + ).fetchone() |
| 40 | + cs_row = reader.execute(f"SELECT count(*) FROM pg.changesets WHERE {older}").fetchone() |
| 41 | + reader.close() |
25 | 42 | stats_n = stats_row[0] if stats_row else 0 |
26 | 43 | cs_n = cs_row[0] if cs_row else 0 |
| 44 | + |
27 | 45 | if cs_n: |
28 | | - conn.execute(f"DELETE FROM pg.changeset_stats WHERE changeset_id IN ({old_cs})") |
29 | | - conn.execute(f"DELETE FROM pg.changesets WHERE created_at < TIMESTAMPTZ '{iso}'") |
30 | | - conn.close() |
| 46 | + writer = _attach(dsn) |
| 47 | + _pg_execute( |
| 48 | + writer, |
| 49 | + f"DELETE FROM changeset_stats s USING changesets c WHERE s.changeset_id = c.changeset_id AND c.{older}", |
| 50 | + ) |
| 51 | + _pg_execute(writer, f"DELETE FROM changesets WHERE {older}") |
| 52 | + writer.close() |
31 | 53 | return stats_n, cs_n |
32 | 54 |
|
33 | 55 |
|
|
0 commit comments