cog_geographic_rollup(): push aggregation and pagination into SQL #59

Closed
opened 2026-08-09 11:21:16 -04:00 by jared · 2 comments
Owner

cog_geographic_rollup() calls cog_spending() with no limit, then does a
dplyr::left_join(relationship = "many-to-many") and filters/reorders in R. It is the
last verb with no bound on what it materializes.

Measured

Through cog-api against the full corpus (schema v7), /rollups?layer=city&category=Police&years=2022:

  • 2,990 ms before cog-api stopped converting the whole result to per-row lists before
    paginating; 1,268 ms after. The remaining ~1.3 s is inside this verb.
  • The documented worst case for that route is ~19,236 rows, all of which are built even
    when the caller asked for 1,000.

cog-api has now fixed everything it can on its side (it slices the data frame before
materializing rows), so the remaining cost is the unpaginated verb plus the R-side join.
/rollups is the slowest route in the API by an order of magnitude, and the only one
still without pushdown.

Ask

  1. limit / offset with the #39 semantics (NULL default, applied in SQL, unpaginated
    count as total_rows).
  2. Move the aggregation and the crosswalk join into SQL rather than joining in R after a
    full fetch. The many-to-many join is the part most likely to be doing avoidable work.

Additive, adoptable behind the existing formals probe.

`cog_geographic_rollup()` calls `cog_spending()` with **no limit**, then does a `dplyr::left_join(relationship = "many-to-many")` and filters/reorders in R. It is the last verb with no bound on what it materializes. ## Measured Through cog-api against the full corpus (schema v7), `/rollups?layer=city&category=Police&years=2022`: - **2,990 ms** before cog-api stopped converting the whole result to per-row lists before paginating; **1,268 ms** after. The remaining ~1.3 s is inside this verb. - The documented worst case for that route is ~19,236 rows, all of which are built even when the caller asked for 1,000. cog-api has now fixed everything it can on its side (it slices the data frame before materializing rows), so the remaining cost is the unpaginated verb plus the R-side join. `/rollups` is the slowest route in the API by an order of magnitude, and the only one still without pushdown. ## Ask 1. `limit` / `offset` with the #39 semantics (`NULL` default, applied in SQL, unpaginated count as `total_rows`). 2. Move the aggregation and the crosswalk join into SQL rather than joining in R after a full fetch. The many-to-many join is the part most likely to be doing avoidable work. Additive, adoptable behind the existing formals probe.
Author
Owner

Measured before implementing, and the premise here does not hold

This issue proposes moving the aggregation and crosswalk join into SQL, on the theory that
"the many-to-many join is the part most likely to be doing avoidable work." I profiled it
against the full production corpus (layer=city, category=Police, years=2022, 20,106
governments, 19,236 rows out):

step time
cog_spending() alone 1442 ms
cog_geographic_rollup() end to end 1163 ms
the many-to-many left_join 6 ms
.coverage_table() 7 ms

The R-side post-processing this issue targets is ~13 ms, about 1% of the total. There is
essentially nothing to win by moving it into SQL.

Pagination does not help either, for a structural reason: GROUP BY must complete
before LIMIT applies, and .coverage_table() plus the coverage = "consistent" filter
both need the complete result — computing them on a page would report wrong coverage.
Measured through cog-api, limit=10 costs 1083 ms against 1332 ms for limit=1000, so the
~19% difference is DuckDB→R transfer, not recoverable query work.

I had estimated 3–4x for this issue earlier. That estimate was wrong — it read the
limit=10/limit=1000 gap as recoverable when the aggregate is the floor.

The real fix is #58, and it is a bigger win than expected

The cohort is passed to cog_spending() as a govid vector, which .sql_lit_chr() renders
into a 301,591-character IN list embedded in 5–8 statements per call. Same aggregate,
four ways:

cohort expressed as time
IN (20,106 literals) — today 449 ms
join against a temp cohort table 99 ms
predicate on canonical_fips_xwalk — what #58 proposes 94 ms
no cohort filter at all (floor) 88 ms

A 4.8x penalty purely from the IN list, and the #58 approach lands within 7% of the
theoretical floor. Since the list is embedded in several statements per call, the saving is
plausibly larger than the 355 ms this single query shows.

Suggestion

Close or re-scope this in favour of #58, which is where the rollup win actually lives. If
limit/offset are still wanted on this verb for API symmetry, they should be added as
ergonomics with the honest note that they bound response size rather than query cost — and
they must not be allowed to reach the coverage computation.

Method: scripts/bench.sh in Civilytics/cog-api, plus direct DuckDB timings against the
production corpus.

## Measured before implementing, and the premise here does not hold This issue proposes moving the aggregation and crosswalk join into SQL, on the theory that "the many-to-many join is the part most likely to be doing avoidable work." I profiled it against the full production corpus (`layer=city`, `category=Police`, `years=2022`, 20,106 governments, 19,236 rows out): | step | time | |---|---:| | `cog_spending()` alone | **1442 ms** | | `cog_geographic_rollup()` end to end | 1163 ms | | the many-to-many `left_join` | **6 ms** | | `.coverage_table()` | **7 ms** | The R-side post-processing this issue targets is **~13 ms, about 1% of the total**. There is essentially nothing to win by moving it into SQL. **Pagination does not help either**, for a structural reason: `GROUP BY` must complete before `LIMIT` applies, and `.coverage_table()` plus the `coverage = "consistent"` filter both need the *complete* result — computing them on a page would report wrong coverage. Measured through cog-api, `limit=10` costs 1083 ms against 1332 ms for `limit=1000`, so the ~19% difference is DuckDB→R transfer, not recoverable query work. I had estimated 3–4x for this issue earlier. **That estimate was wrong** — it read the limit=10/limit=1000 gap as recoverable when the aggregate is the floor. ## The real fix is #58, and it is a bigger win than expected The cohort is passed to `cog_spending()` as a `govid` vector, which `.sql_lit_chr()` renders into a **301,591-character** `IN` list embedded in 5–8 statements per call. Same aggregate, four ways: | cohort expressed as | time | |---|---:| | `IN (20,106 literals)` — today | **449 ms** | | join against a temp cohort table | 99 ms | | predicate on `canonical_fips_xwalk` — **what #58 proposes** | **94 ms** | | no cohort filter at all (floor) | 88 ms | **A 4.8x penalty purely from the IN list**, and the #58 approach lands within 7% of the theoretical floor. Since the list is embedded in several statements per call, the saving is plausibly larger than the 355 ms this single query shows. ## Suggestion Close or re-scope this in favour of #58, which is where the rollup win actually lives. If `limit`/`offset` are still wanted on this verb for API symmetry, they should be added as ergonomics with the honest note that they bound response size rather than query cost — and they must not be allowed to reach the coverage computation. Method: `scripts/bench.sh` in `Civilytics/cog-api`, plus direct DuckDB timings against the production corpus.
Author
Owner

Closing in favour of #58, which is where the rollup win actually is.

Profiling (detail in the comment above) showed this issue targets ~1% of the runtime: the
many-to-many left_join is 6 ms and .coverage_table() 7 ms, against 1442 ms inside
cog_spending(). Pagination cannot help either — GROUP BY completes before LIMIT
applies, and both the coverage table and coverage = "consistent" need the full result, so
a page-scoped computation would report wrong coverage.

The cost is the cohort predicate, not the join: expressing the same 20,106-government cohort
as a crosswalk predicate instead of a 301,591-character IN list takes 94 ms instead of
449 ms
. That is #58.

Reopen if limit/offset are wanted here purely for API symmetry — but they would bound
response size, not query cost, and must be kept away from the coverage computation.

Closing in favour of **#58**, which is where the rollup win actually is. Profiling (detail in the comment above) showed this issue targets ~1% of the runtime: the many-to-many `left_join` is 6 ms and `.coverage_table()` 7 ms, against 1442 ms inside `cog_spending()`. Pagination cannot help either — `GROUP BY` completes before `LIMIT` applies, and both the coverage table and `coverage = "consistent"` need the full result, so a page-scoped computation would report wrong coverage. The cost is the cohort predicate, not the join: expressing the same 20,106-government cohort as a crosswalk predicate instead of a 301,591-character `IN` list takes **94 ms instead of 449 ms**. That is #58. Reopen if `limit`/`offset` are wanted here purely for API symmetry — but they would bound response size, not query cost, and must be kept away from the coverage computation.
jared closed this issue 2026-08-09 13:48:35 -04:00
Sign in to join this conversation.
1 Participants
Notifications
Due Date
No due date set.
Dependencies

No dependencies set.

Reference: Civilytics/uscogdata#59