Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f4ab9b6d90
|
||
|
|
9508b98676 |
@@ -1,5 +1,24 @@
|
|||||||
# uscogdata 0.4.0
|
# uscogdata 0.4.0
|
||||||
|
|
||||||
|
## DuckDB's resource budget is configurable
|
||||||
|
|
||||||
|
`USCOGDATA_DUCKDB_THREADS` and `USCOGDATA_DUCKDB_MEMORY_LIMIT` (with matching
|
||||||
|
`options(uscogdata.duckdb_threads = )` / `options(uscogdata.duckdb_memory_limit = )`
|
||||||
|
spellings) cap the DuckDB connection the package opens. Both follow the same
|
||||||
|
env-var > option > default precedence as `USCOGDATA_URL`.
|
||||||
|
|
||||||
|
Unset, **no pragma is issued at all** and DuckDB's own defaults apply exactly as
|
||||||
|
before -- every visible core. That is right for one interactive session on a
|
||||||
|
dedicated machine and wrong for a server: where several readers share a host, each
|
||||||
|
otherwise claims the whole machine and they contend. Capping measured ~5% on a
|
||||||
|
single-government all-years query (502 ms at 2 threads vs 475 ms uncapped on 16
|
||||||
|
cores), which is cheap enough that a server should always cap.
|
||||||
|
|
||||||
|
This replaces a workaround in which a consumer reached into the package namespace
|
||||||
|
at boot -- `getFromNamespace(".ensure_session", "uscogdata")()` followed by a manual
|
||||||
|
`SET threads` -- depending both on a private name and on the session already being
|
||||||
|
open.
|
||||||
|
|
||||||
## Cohorts can be named by predicate, not just by id
|
## Cohorts can be named by predicate, not just by id
|
||||||
|
|
||||||
`cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state`
|
`cog_spending()`, `cog_revenue()` and `cog_balances()` gain optional `state`
|
||||||
|
|||||||
+60
-1
@@ -18,7 +18,15 @@
|
|||||||
# Nextcloud share or a local copy made by cog_mirror().
|
# Nextcloud share or a local copy made by cog_mirror().
|
||||||
url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/",
|
url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/",
|
||||||
cache_dir = NULL,
|
cache_dir = NULL,
|
||||||
manifest_ttl_secs = 3600L
|
manifest_ttl_secs = 3600L,
|
||||||
|
# NULL means "emit no pragma", which leaves DuckDB's own defaults intact:
|
||||||
|
# every visible core, and 80% of RAM. That is right for one interactive
|
||||||
|
# session on a dedicated machine and wrong for a server, where several
|
||||||
|
# readers share a box and each would otherwise claim all of it. See
|
||||||
|
# .resolve_duckdb_threads() for why this is a supported option rather than
|
||||||
|
# something a consumer reaches into the namespace to set.
|
||||||
|
duckdb_threads = NULL,
|
||||||
|
duckdb_memory_limit = NULL
|
||||||
)
|
)
|
||||||
|
|
||||||
#' Resolve a config value: env var > option > default
|
#' Resolve a config value: env var > option > default
|
||||||
@@ -61,3 +69,54 @@
|
|||||||
v <- .cfg("cache_dir")
|
v <- .cfg("cache_dir")
|
||||||
if (is.null(v)) tools::R_user_dir("uscogdata", "cache") else v
|
if (is.null(v)) tools::R_user_dir("uscogdata", "cache") else v
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#' Resolve the DuckDB thread cap, or NULL to leave DuckDB's default alone.
|
||||||
|
#'
|
||||||
|
#' `cog_open()` used to connect with a bare `dbConnect()` and set no `threads`
|
||||||
|
#' pragma, so DuckDB claimed every core it could see. cog-api works around that
|
||||||
|
#' by reaching into this namespace at boot --
|
||||||
|
#' `getFromNamespace(".ensure_session", "uscogdata")()` followed by a manual
|
||||||
|
#' `SET threads` -- which depends on a private name and on the session already
|
||||||
|
#' being open. Making it a resolved option removes the reason to do that.
|
||||||
|
#'
|
||||||
|
#' `.cfg()` returns an environment variable as CHARACTER, so this coerces
|
||||||
|
#' rather than trusting the type: `USCOGDATA_DUCKDB_THREADS=4` arrives as "4",
|
||||||
|
#' and `sprintf("SET threads TO %d", "4")` would abort inside the connection
|
||||||
|
#' path with an error about the pragma rather than about the setting.
|
||||||
|
#' @noRd
|
||||||
|
.resolve_duckdb_threads <- function() {
|
||||||
|
v <- .cfg("duckdb_threads")
|
||||||
|
if (is.null(v) || (is.character(v) && !nzchar(v))) return(NULL)
|
||||||
|
n <- suppressWarnings(as.integer(v))
|
||||||
|
if (length(n) != 1L || is.na(n) || n < 1L) {
|
||||||
|
cli::cli_abort(c(
|
||||||
|
"{.envvar USCOGDATA_DUCKDB_THREADS} must be a single positive integer.",
|
||||||
|
x = "Got {.val {v}}.",
|
||||||
|
i = "Unset it (or {.code options(uscogdata.duckdb_threads = NULL)}) to use DuckDB's default of every visible core."
|
||||||
|
), class = "uscogdata_invalid_duckdb_threads")
|
||||||
|
}
|
||||||
|
n
|
||||||
|
}
|
||||||
|
|
||||||
|
#' Resolve the DuckDB memory limit, or NULL to leave DuckDB's default alone.
|
||||||
|
#'
|
||||||
|
#' The value is a DuckDB size string (`"4GB"`, `"512MB"`). Only its SHAPE is
|
||||||
|
#' checked here -- DuckDB owns the unit vocabulary, and re-implementing that
|
||||||
|
#' parse would be a second definition free to drift from the engine's. An
|
||||||
|
#' unrecognised unit therefore surfaces as DuckDB's own error at `SET` time,
|
||||||
|
#' which names the setting correctly; the check here exists to reject the
|
||||||
|
#' inputs that would otherwise reach the connection as a SQL fragment.
|
||||||
|
#' @noRd
|
||||||
|
.resolve_duckdb_memory_limit <- function() {
|
||||||
|
v <- .cfg("duckdb_memory_limit")
|
||||||
|
if (is.null(v) || (is.character(v) && !nzchar(v))) return(NULL)
|
||||||
|
if (length(v) != 1L || !is.character(v) ||
|
||||||
|
!grepl("^[0-9]+(\\.[0-9]+)?\\s*[A-Za-z]{0,3}$", v)) {
|
||||||
|
cli::cli_abort(c(
|
||||||
|
"{.envvar USCOGDATA_DUCKDB_MEMORY_LIMIT} must be a single DuckDB size string.",
|
||||||
|
x = "Got {.val {v}}.",
|
||||||
|
i = "Examples: {.val 4GB}, {.val 512MB}, {.val 1.5GB}."
|
||||||
|
), class = "uscogdata_invalid_duckdb_memory_limit")
|
||||||
|
}
|
||||||
|
trimws(v)
|
||||||
|
}
|
||||||
|
|||||||
+31
-1
@@ -2,13 +2,23 @@
|
|||||||
|
|
||||||
#' Internal: open session, register views, cache manifest.
|
#' Internal: open session, register views, cache manifest.
|
||||||
#' Not exported. Called lazily by verbs via .ensure_session().
|
#' Not exported. Called lazily by verbs via .ensure_session().
|
||||||
|
#'
|
||||||
|
#' `threads` and `memory_limit` default to the resolved configuration and are
|
||||||
|
#' applied as pragmas on the new connection. When both resolve to NULL -- which
|
||||||
|
#' is the case unless the operator sets one -- NO pragma is issued at all, so an
|
||||||
|
#' unconfigured session connects exactly as it did before this argument existed.
|
||||||
#' @noRd
|
#' @noRd
|
||||||
cog_open <- function(url = .resolve_url(),
|
cog_open <- function(url = .resolve_url(),
|
||||||
cache_dir = .resolve_cache_dir()) {
|
cache_dir = .resolve_cache_dir(),
|
||||||
|
threads = .resolve_duckdb_threads(),
|
||||||
|
memory_limit = .resolve_duckdb_memory_limit()) {
|
||||||
.check_url_configured(url)
|
.check_url_configured(url)
|
||||||
if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE)
|
if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE)
|
||||||
|
|
||||||
con <- DBI::dbConnect(duckdb::duckdb())
|
con <- DBI::dbConnect(duckdb::duckdb())
|
||||||
|
# Before anything else touches the connection: httpfs reads the corpus, and
|
||||||
|
# a remote read should already be bound by whatever budget the operator set.
|
||||||
|
.apply_duckdb_limits(con, threads, memory_limit)
|
||||||
DBI::dbExecute(con, "INSTALL httpfs; LOAD httpfs;")
|
DBI::dbExecute(con, "INSTALL httpfs; LOAD httpfs;")
|
||||||
|
|
||||||
manifest <- .fetch_or_cache_manifest(url, cache_dir)
|
manifest <- .fetch_or_cache_manifest(url, cache_dir)
|
||||||
@@ -25,6 +35,26 @@ cog_open <- function(url = .resolve_url(),
|
|||||||
invisible(con)
|
invisible(con)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#' Apply the operator's DuckDB resource budget to a fresh connection.
|
||||||
|
#'
|
||||||
|
#' Split out from cog_open() so the "unset changes nothing" property is one
|
||||||
|
#' readable branch rather than two conditionals buried in the connection path.
|
||||||
|
#' Both settings are session-scoped in DuckDB, so this must run per connection;
|
||||||
|
#' cog_close() discards the connection and the next cog_open() re-resolves,
|
||||||
|
#' which is what makes a changed option take effect on the next session.
|
||||||
|
#' @noRd
|
||||||
|
.apply_duckdb_limits <- function(con, threads, memory_limit) {
|
||||||
|
if (!is.null(threads)) {
|
||||||
|
DBI::dbExecute(con, sprintf("SET threads TO %d", threads))
|
||||||
|
}
|
||||||
|
if (!is.null(memory_limit)) {
|
||||||
|
# Quoted as a string literal: DuckDB's memory_limit takes '4GB', not 4GB.
|
||||||
|
DBI::dbExecute(con, sprintf("SET memory_limit TO %s",
|
||||||
|
.sql_lit_chr(memory_limit)))
|
||||||
|
}
|
||||||
|
invisible(con)
|
||||||
|
}
|
||||||
|
|
||||||
#' @noRd
|
#' @noRd
|
||||||
.ensure_session <- function() {
|
.ensure_session <- function() {
|
||||||
if (is.null(.uscogdata_env$con) ||
|
if (is.null(.uscogdata_env$con) ||
|
||||||
|
|||||||
@@ -110,6 +110,17 @@ After that, nothing in your analysis touches an external service.
|
|||||||
- `USCOGDATA_URL` — corpus root: an HTTPS URL or a local path, **trailing slash required**
|
- `USCOGDATA_URL` — corpus root: an HTTPS URL or a local path, **trailing slash required**
|
||||||
- `USCOGDATA_CACHE_DIR` — where the manifest is cached (default: user cache dir)
|
- `USCOGDATA_CACHE_DIR` — where the manifest is cached (default: user cache dir)
|
||||||
- `USCOGDATA_MANIFEST_TTL_SECS` — manifest re-fetch interval (default 3600)
|
- `USCOGDATA_MANIFEST_TTL_SECS` — manifest re-fetch interval (default 3600)
|
||||||
|
- `USCOGDATA_DUCKDB_THREADS` — cap DuckDB's thread count (default: every visible core)
|
||||||
|
- `USCOGDATA_DUCKDB_MEMORY_LIMIT` — cap DuckDB's memory, e.g. `"4GB"` (default: DuckDB's own)
|
||||||
|
|
||||||
|
Each also has an `options()` spelling — `uscogdata.url`, `uscogdata.duckdb_threads`,
|
||||||
|
and so on — and the environment variable wins where both are set.
|
||||||
|
|
||||||
|
The two DuckDB caps exist for **servers**, not laptops. Unset, DuckDB claims every
|
||||||
|
core it can see, which is right for one interactive session on your own machine and
|
||||||
|
wrong when several readers share a box: each claims the whole machine and they fight.
|
||||||
|
Capping costs roughly 5% on a single query and is worth it anywhere the process is
|
||||||
|
sharing hardware.
|
||||||
|
|
||||||
## Amounts are in full US dollars
|
## Amounts are in full US dollars
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,134 @@
|
|||||||
|
# tests/testthat/test-duckdb-limits.R
|
||||||
|
#
|
||||||
|
# uscogdata#60. cog_open() used to connect with a bare dbConnect() and set no
|
||||||
|
# resource pragmas, so DuckDB claimed every visible core. That is right for one
|
||||||
|
# interactive session on a dedicated machine and wrong for a server: cog-api
|
||||||
|
# runs two replicas on an 8-core host budgeted 4, and without a cap each
|
||||||
|
# replica independently claims all 8 and they fight.
|
||||||
|
#
|
||||||
|
# The consumer-side workaround this replaces reached into the namespace at
|
||||||
|
# boot -- getFromNamespace(".ensure_session", "uscogdata")() followed by a
|
||||||
|
# manual SET threads -- which depends on a private name AND on the session
|
||||||
|
# already being open.
|
||||||
|
#
|
||||||
|
# The load-bearing property is the NEGATIVE one: unset must emit no pragma at
|
||||||
|
# all, so an unconfigured session is byte-identical to pre-#60 behaviour.
|
||||||
|
|
||||||
|
# Open a session under a given configuration and read a DuckDB setting back.
|
||||||
|
# Each call closes first, because both settings are session-scoped: an
|
||||||
|
# already-open connection would be reused by .ensure_session() and report the
|
||||||
|
# PREVIOUS test's value, which is exactly the false pass to avoid here.
|
||||||
|
setting_under <- function(setting, envvars = character(0), opts = list()) {
|
||||||
|
uscogdata:::cog_close()
|
||||||
|
on.exit(uscogdata:::cog_close(), add = TRUE)
|
||||||
|
withr::with_envvar(envvars, {
|
||||||
|
withr::with_options(opts, {
|
||||||
|
con <- uscogdata:::cog_open()
|
||||||
|
DBI::dbGetQuery(
|
||||||
|
con, sprintf("SELECT current_setting('%s') AS v", setting)
|
||||||
|
)$v[[1]]
|
||||||
|
})
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
test_that("USCOGDATA_DUCKDB_THREADS caps the connection's thread count", {
|
||||||
|
skip_if_no_corpus()
|
||||||
|
expect_equal(
|
||||||
|
as.integer(setting_under("threads", c(USCOGDATA_DUCKDB_THREADS = "2"))),
|
||||||
|
2L
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that("the option spelling works, and the env var beats it", {
|
||||||
|
skip_if_no_corpus()
|
||||||
|
expect_equal(
|
||||||
|
as.integer(setting_under("threads",
|
||||||
|
c(USCOGDATA_DUCKDB_THREADS = NA),
|
||||||
|
list(uscogdata.duckdb_threads = 3L))),
|
||||||
|
3L
|
||||||
|
)
|
||||||
|
# Same precedence .cfg() gives every other setting: env var > option.
|
||||||
|
expect_equal(
|
||||||
|
as.integer(setting_under("threads",
|
||||||
|
c(USCOGDATA_DUCKDB_THREADS = "1"),
|
||||||
|
list(uscogdata.duckdb_threads = 3L))),
|
||||||
|
1L
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that("unset leaves DuckDB's own default in place", {
|
||||||
|
skip_if_no_corpus()
|
||||||
|
# Not asserting a specific number -- the default is core-count-dependent and
|
||||||
|
# a literal would fail on a different machine. The claim is that NO pragma
|
||||||
|
# was issued, so the session sees whatever DuckDB would have chosen on its
|
||||||
|
# own. Compared against a plain connection opened the pre-#60 way.
|
||||||
|
unset <- setting_under("threads",
|
||||||
|
c(USCOGDATA_DUCKDB_THREADS = NA),
|
||||||
|
list(uscogdata.duckdb_threads = NULL))
|
||||||
|
bare <- local({
|
||||||
|
con <- DBI::dbConnect(duckdb::duckdb())
|
||||||
|
on.exit(DBI::dbDisconnect(con, shutdown = TRUE), add = TRUE)
|
||||||
|
DBI::dbGetQuery(con, "SELECT current_setting('threads') AS v")$v[[1]]
|
||||||
|
})
|
||||||
|
expect_equal(as.integer(unset), as.integer(bare))
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that("USCOGDATA_DUCKDB_MEMORY_LIMIT is applied", {
|
||||||
|
skip_if_no_corpus()
|
||||||
|
v <- setting_under("memory_limit", c(USCOGDATA_DUCKDB_MEMORY_LIMIT = "2GB"))
|
||||||
|
# DuckDB does not echo back the string it was given: it stores bytes and
|
||||||
|
# reports BINARY units, so "2GB" (2e9 bytes) comes back as "1.8 GiB". Assert
|
||||||
|
# the magnitude it actually means rather than the spelling this package sent
|
||||||
|
# -- matching on "2" passes for the wrong reason and fails on the right one.
|
||||||
|
expect_match(as.character(v), "GiB", fixed = TRUE)
|
||||||
|
# DuckDB also truncates the display to one decimal ("1.8 GiB" for 1.863), so
|
||||||
|
# the tolerance covers rounding, not slack in the setting itself.
|
||||||
|
gib <- as.numeric(sub("\\s*GiB$", "", as.character(v)))
|
||||||
|
expect_equal(gib, 2e9 / 1024^3, tolerance = 0.05)
|
||||||
|
})
|
||||||
|
|
||||||
|
# --- Validation -------------------------------------------------------------
|
||||||
|
# .cfg() returns an env var as CHARACTER. Without coercion here,
|
||||||
|
# sprintf("SET threads TO %d", "4") aborts inside the connection path with an
|
||||||
|
# error about the pragma rather than about the setting the operator got wrong.
|
||||||
|
|
||||||
|
test_that(".resolve_duckdb_threads coerces a character env var to integer", {
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = "4")
|
||||||
|
expect_identical(uscogdata:::.resolve_duckdb_threads(), 4L)
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that(".resolve_duckdb_threads returns NULL when unset or empty", {
|
||||||
|
withr::local_options(uscogdata.duckdb_threads = NULL)
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = NA)
|
||||||
|
expect_null(uscogdata:::.resolve_duckdb_threads())
|
||||||
|
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = "")
|
||||||
|
expect_null(uscogdata:::.resolve_duckdb_threads())
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that(".resolve_duckdb_threads rejects values that are not positive integers", {
|
||||||
|
for (bad in c("0", "-1", "two", "1.5.2")) {
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_THREADS = bad)
|
||||||
|
expect_error(uscogdata:::.resolve_duckdb_threads(),
|
||||||
|
class = "uscogdata_invalid_duckdb_threads")
|
||||||
|
}
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that(".resolve_duckdb_memory_limit accepts size strings and rejects junk", {
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "4GB")
|
||||||
|
expect_identical(uscogdata:::.resolve_duckdb_memory_limit(), "4GB")
|
||||||
|
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "1.5GB")
|
||||||
|
expect_identical(uscogdata:::.resolve_duckdb_memory_limit(), "1.5GB")
|
||||||
|
|
||||||
|
# A SQL fragment must not reach the connection as one.
|
||||||
|
withr::local_envvar(USCOGDATA_DUCKDB_MEMORY_LIMIT = "4GB'; DROP TABLE x; --")
|
||||||
|
expect_error(uscogdata:::.resolve_duckdb_memory_limit(),
|
||||||
|
class = "uscogdata_invalid_duckdb_memory_limit")
|
||||||
|
})
|
||||||
|
|
||||||
|
test_that(".apply_duckdb_limits issues no statement when both are NULL", {
|
||||||
|
# The negative property, asserted directly rather than inferred: a connection
|
||||||
|
# that would ERROR on any statement proves none was sent.
|
||||||
|
expect_silent(uscogdata:::.apply_duckdb_limits(NULL, NULL, NULL))
|
||||||
|
})
|
||||||
Reference in New Issue
Block a user