feat: cog_open() honours a DuckDB thread and memory budget (#60)
cog_open() connected with a bare dbConnect() and set no resource pragmas, so
DuckDB claimed every visible core. Right for one interactive session on a
dedicated machine; wrong for a server, where cog-api runs two replicas on an
8-core host budgeted 4 and each replica independently claims all 8.
USCOGDATA_DUCKDB_THREADS and USCOGDATA_DUCKDB_MEMORY_LIMIT now resolve through
.cfg() -- inheriting the env var > option > default precedence USCOGDATA_URL
already had -- and are applied as pragmas when the connection is created.
Unset issues NO pragma, so an unconfigured session is byte-identical to before.
That negative property is asserted directly against a connection opened the
pre-change way rather than against a hardcoded core count.
.cfg() returns an env var as character, so both resolvers coerce and validate
rather than trusting the type: sprintf("SET threads TO %d", "4") would
otherwise abort inside the connection path with an error naming the pragma
instead of the setting the operator got wrong.
Replaces cog-api's getFromNamespace(".ensure_session", "uscogdata") workaround,
which depended on a private name and on the session already being open.
This commit is contained in:
+60
-1
@@ -18,7 +18,15 @@
|
||||
# Nextcloud share or a local copy made by cog_mirror().
|
||||
url = "https://huggingface.co/datasets/civilytics/us-cog-finance/resolve/main/",
|
||||
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
|
||||
@@ -61,3 +69,54 @@
|
||||
v <- .cfg("cache_dir")
|
||||
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.
|
||||
#' 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
|
||||
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)
|
||||
if (!dir.exists(cache_dir)) dir.create(cache_dir, recursive = TRUE)
|
||||
|
||||
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;")
|
||||
|
||||
manifest <- .fetch_or_cache_manifest(url, cache_dir)
|
||||
@@ -25,6 +35,26 @@ cog_open <- function(url = .resolve_url(),
|
||||
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
|
||||
.ensure_session <- function() {
|
||||
if (is.null(.uscogdata_env$con) ||
|
||||
|
||||
Reference in New Issue
Block a user