API reference¶
Every signature, parameter table and accepted value below is rendered from the
docstrings of the version named in the header, so it changes with the code. In
Python, help(sus.load) shows the same text for the version you have installed.
Which years, ufs and months each dataset takes is in
Bases e argumentos. To install, follow the
installation guide.
For researchers¶
Four calls cover an analysis from download to citation:
import omnisus as sus
dados = sus.load("sim_obitos", years=[2023], ufs=["RR"]) # rows, codes as published
dados = sus.label("sim_obitos", dados, columns=["sexo", "racacor"]) # + sexo_rotulo, racacor_rotulo
sus.check_columns("sim_obitos", dados) # can I trust each column?
with sus.LakeReader() as lake:
print(sus.cite(lake, dataset="sim_obitos").text) # files, SHA-256, snapshot
loaddownloads what DATASUS publishes into the lake (where it lives) and returns the rows. A second call downloads nothing. When the files are validated sources it also adds harmonised categories (idade_anos_completos,sexo_categoria,<date>_data); otherwise it warns which were left out (see ADR 0003).labelputs the dictionary's label next to each code; a code the dictionary does not know gets no label.check_columnsreports, per column, empties, unlabelled codes and date ranges.
Download what DATASUS publishes for these scopes and return the rows.
The one-call path for a script or notebook cell. It imports into the lake (a
second call downloads nothing new) and returns the requested scopes' rows with
codes as DATASUS published them. When every file comes from a validated source
(see describe_dataset(dataset)["analytics"]["validated_sources"]) it also adds
the harmonised categories: idade_anos_completos, idade_status,
sexo_categoria, sexo_status and <date>_data columns. Otherwise it
leaves them out and warns once. Readable labels come from :func:label.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
DATASUS FTP dataset name, e.g. |
required |
years
|
Iterable[int]
|
Years, e.g. |
required |
ufs
|
Sequence[str] | None
|
UF abbreviations, e.g. |
None
|
months
|
Iterable[int] | None
|
Months 1-12 of monthly datasets (SIH, SIA, CNES), e.g. |
None
|
target
|
str | None
|
DuckLake target; |
None
|
policy
|
ImportPolicy
|
|
'skip_same'
|
Returns:
| Type | Description |
|---|---|
DataFrame
|
A polars DataFrame with one row per record of the requested scopes. |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
RuntimeError
|
a scope failed to import; the message names each one. |
LookupError
|
DATASUS publishes none of the requested scopes. |
Warns:
| Type | Description |
|---|---|
UserWarning
|
harmonised categories were left out, and why. |
Examples:
>>> import omnisus as sus
>>> dados = sus.load("sim_obitos", years=[2023], ufs=["RR"])
>>> dados.height
3311
>>> dados = sus.label("sim_obitos", dados, columns=["sexo"])
Source code in src/omnisus/__init__.py
252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 | |
Add the dictionary's label next to each coded column, keeping the code.
Each coded column c gains c_rotulo right after it. A code the
dictionary does not know gets a null label, and so do a null and a blank the
map does not declare: nothing is guessed, and the code stays in c.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
df
|
DataFrame
|
Rows of that dataset, e.g. from :func: |
required |
columns
|
Sequence[str] | None
|
Columns to label, e.g. |
None
|
Returns:
| Type | Description |
|---|---|
DataFrame
|
|
Raises:
| Type | Description |
|---|---|
ValueError
|
a name in |
FileNotFoundError
|
|
Examples:
>>> import polars as pl
>>> import omnisus as sus
>>> dados = pl.DataFrame({"sexo": ["1", "2", " "], "idade": ["435", "450", "401"]})
>>> sus.label("sim_obitos", dados)
shape: (3, 3)
┌──────┬─────────────┬───────┐
│ sexo ┆ sexo_rotulo ┆ idade │
│ --- ┆ --- ┆ --- │
│ str ┆ str ┆ str │
╞══════╪═════════════╪═══════╡
│ 1 ┆ Masculino ┆ 435 │
│ 2 ┆ Feminino ┆ 450 │
│ ┆ null ┆ 401 │
└──────┴─────────────┴───────┘
Source code in src/omnisus/transforms/dictionaries.py
161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 | |
Report, per column, the dictionary's description and what the rows show.
Answers "can I trust this column?" before an analysis: how much is empty,
which published codes have no label, and whether dates parse and fall in a
plausible range (impossible dates such as 1899-12-30 parse, so look at
date_min). A value counts as empty when it is null or an empty string.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
df
|
DataFrame
|
Rows of that dataset, e.g. from :func: |
required |
Returns:
| Type | Description |
|---|---|
DataFrame
|
One row per column of |
DataFrame
|
|
DataFrame
|
|
DataFrame
|
dictionary" |
DataFrame
|
|
DataFrame
|
label); |
DataFrame
|
label) and |
DataFrame
|
|
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
|
Examples:
SIH publishes homonimo 2, which no DATASUS table labels:
>>> import polars as pl
>>> import omnisus as sus
>>> dados = pl.DataFrame({"homonimo": ["0", "2", ""], "nasc": ["19850320", "18991230", ""]})
>>> relatorio = sus.check_columns("sih_aih_reduzida", dados)
>>> relatorio.select("column", "pct_empty", "unlabelled_codes", "date_min").rows()
[('homonimo', 33.3, ['2'], None), ('nasc', 33.3, [], datetime.date(1899, 12, 30))]
Source code in src/omnisus/transforms/columns.py
41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 | |
The Portuguese citation of what a lake holds: files, hashes, snapshot and run.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
lake
|
Lake | LakeReader
|
An open :class: |
required |
dataset
|
str | None
|
Cite only this dataset, e.g. |
None
|
snapshot_id
|
int | None
|
The snapshot your analysis read; |
None
|
run_id
|
str | None
|
Cite only this run's publications. |
None
|
accessed
|
date | None
|
Access date to print; |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
Citation
|
class: |
Raises:
| Type | Description |
|---|---|
LookupError
|
|
Examples:
>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
... print(sus.cite(lake, dataset="sim_obitos").text)
Source code in src/omnisus/research.py
290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 | |
Discovery¶
Ask the server what exists before deciding what to import.
The scopes DATASUS publishes now: what an import can actually get.
Pass the result to :func:~omnisus.import_dataset; unlike
:func:~omnisus.scopes_for, it plans only files the server lists. Files of
other datasets sharing a directory are skipped.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
years
|
Iterable[int] | None
|
Keep these years, e.g. |
None
|
ufs
|
Sequence[str] | None
|
Keep these UFs, e.g. |
None
|
months
|
Iterable[int] | None
|
Keep these months 1-12; |
None
|
refresh
|
bool
|
|
False
|
Returns:
| Type | Description |
|---|---|
list[ScopeKey]
|
The scopes, ordered by year, UF and month. |
Raises:
| Type | Description |
|---|---|
ValueError
|
as :func: |
FtpUnavailable
|
the server did not answer after retries. |
Examples:
>>> import omnisus as sus
>>> sus.available("sim_obitos", years=[2023], ufs=["RR"])
[ScopeKey(uf='RR', ano=2023, mes=None)]
Source code in src/omnisus/sources/datasus_ftp/inventory.py
403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 | |
The scopes DATASUS publishes now, and whether each is final or preliminary.
One cached listing per directory. A scope found in two directories is a server inconsistency and raises; it is never resolved by preference.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
years
|
Iterable[int] | None
|
Keep these years, e.g. |
None
|
ufs
|
Sequence[str] | None
|
Keep these UFs, e.g. |
None
|
months
|
Iterable[int] | None
|
Keep these months 1-12; |
None
|
refresh
|
bool
|
|
False
|
Returns:
| Type | Description |
|---|---|
dict[ScopeKey, Release]
|
|
Raises:
| Type | Description |
|---|---|
ValueError
|
|
FtpUnavailable
|
the server did not answer after retries. |
Examples:
>>> import omnisus as sus
>>> sus.available_releases("sim_obitos", years=[2024], ufs=["RR"])
{ScopeKey(uf='RR', ano=2024, mes=None): 'prelim'}
Source code in src/omnisus/sources/datasus_ftp/inventory.py
351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 | |
List any DATASUS FTP path, including datasets the package does not curate.
The open-world counterpart to :func:available: it reaches other SINAN agravos,
CIHA or PCE. Eager, so a huge tree is listed whole; call
omnisus.sources.datasus_ftp.inventory.crawl for a lazy walk.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str
|
Absolute FTP path, e.g. |
required |
depth
|
int
|
How many directory levels to list; |
1
|
refresh
|
bool
|
|
False
|
Returns:
| Name | Type | Description |
|---|---|---|
One |
list[FtpEntry]
|
class: |
Raises:
| Type | Description |
|---|---|
FtpPathNotFound
|
the server has no such path. |
FtpUnavailable
|
the server did not answer after retries. |
Examples:
>>> import omnisus as sus
>>> entries = sus.browse("/dissemin/publicos/SINAN/DADOS/FINAIS")
>>> entries[0].name
'ACBIBR06.dbc'
Source code in src/omnisus/__init__.py
125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 | |
One line of a DATASUS FTP directory listing.
Source code in src/omnisus/sources/datasus_ftp/inventory.py
54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 | |
modified
instance-attribute
¶
Server-reported mtime. Detects DATASUS republishing a file we ingested.
name
instance-attribute
¶
File or directory name, spaces preserved.
parent
instance-attribute
¶
Absolute path of the containing directory.
path
instance-attribute
¶
Absolute remote path.
size_bytes
instance-attribute
¶
0 for directories. Exceeds 32 bits in the wild.
Releases¶
Some datasets publish final and preliminary files under the same names in two
directories (a row's prelim_dir). available_releases reports which
directory each scope came from; outdated compares the files a lake has
published against what the server lists today — by path, and by size or
server mtime when recorded — and returns the scopes that differ (a moved
directory, or a same-name republish), to be re-imported with
import_dataset(..., policy="replace"). See
reprocessing and maintenance.
Scopes whose imported files differ from what the server lists now.
Different means another path (a preliminary year that became final, or a moved directory), or the same path with another size or server time (DATASUS republished it). Publications without a recorded source are always reported. Scopes the server no longer lists are a withdrawal, not returned here. Read-only.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
lake
|
Lake | LakeReader
|
An open :class: |
required |
refresh
|
bool
|
|
False
|
Returns:
| Type | Description |
|---|---|
list[ScopeKey]
|
The scopes to re-import with :func: |
list[ScopeKey]
|
|
Raises:
| Type | Description |
|---|---|
ValueError
|
|
FtpUnavailable
|
the server did not answer after retries. |
Examples:
>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
... stale = sus.outdated("sim_obitos", lake=lake)
>>> sus.import_dataset(
... "sim_obitos", scopes=stale, policy="replace", run_id="refresh-2026-09"
... )
Source code in src/omnisus/__init__.py
532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 | |
Planning¶
Planning is composition: build a list of scopes any way you like and hand it
to import_dataset. There is no planner flag on the Python API.
Every scope (UF x year [x month]) of a dataset, without asking the server.
Planning is composition: pass the result, or any other list of scopes, to
:func:import_dataset. Scopes DATASUS does not publish are later reported as
skipped; :func:available plans only what the server lists.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
years
|
Iterable[int]
|
Years, e.g. |
required |
ufs
|
Sequence[str] | None
|
UF abbreviations, e.g. |
None
|
months
|
Iterable[int] | None
|
Months 1-12 for monthly datasets (SIH, SIA, CNES); |
None
|
Returns:
| Type | Description |
|---|---|
list[ScopeKey]
|
The scopes, ordered by year, then UF, then month. |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> import omnisus as sus
>>> sus.scopes_for("sih_aih_reduzida", years=[2024], ufs=["RR"], months=[1, 2])
[ScopeKey(uf='RR', ano=2024, mes=1), ScopeKey(uf='RR', ano=2024, mes=2)]
>>> sus.scopes_for("sinan_chagas", years=[2023])
[ScopeKey(uf=None, ano=2023, mes=None)]
Source code in src/omnisus/__init__.py
67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 | |
Importing¶
The FTP importers return ImportReport; inspect report.failed, report.skipped
and report.ok. import_ibge_populacao returns list[ImportResult], while
import_cnes_master returns the number of records written.
import_research is the researcher door: it requires run_id and defaults to
policy="skip_same". It refuses append. import_dataset remains the operator
API and still appends unless a policy is set.
ImportAbortedError interrupts an FTP run when it cannot safely continue.
Inspect its report for determined outcomes and unresolved for
(input_index, ScopeKey) pairs before retrying. Imports append data unless an
explicit replay policy is selected. FTP imports accept append (default),
skip_same, error_if_exists and replace. See
reprocessing and maintenance for legacy
scope restrictions, run IDs and byte budgets.
Every import function also runs from inside a notebook cell, where an event loop
is already running. Importing cnes_estabelecimentos by any function refreshes the
aux_cnes view.
Import the given scopes of a DATASUS FTP dataset into the lake (operator door).
The import iterates scopes and never asks where they came from:
:func:scopes_for plans blindly and :func:available asks the server first.
Every run lists the server once and imports what it lists; callers never choose
a release or file name. Importing cnes_estabelecimentos also refreshes the
aux_cnes view. Works inside a notebook's running event loop.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
scopes
|
Sequence[ScopeKey]
|
The scopes to import, e.g. from :func: |
required |
target
|
str | None
|
DuckLake target such as |
None
|
concurrency
|
int
|
Downloads in flight. DATASUS FTP is a shared public server; raise it only with reason. |
DEFAULT_CONCURRENCY
|
batch_size
|
int
|
Scopes committed in one DuckLake transaction, hence one snapshot. |
DEFAULT_BATCH_SIZE
|
policy
|
ImportPolicy
|
What to do with a scope already in the lake: |
'append'
|
run_id
|
str | None
|
Your identifier for this run, chosen before it starts, to reconcile an
interrupted commit through |
None
|
max_payload_bytes
|
int
|
Largest compressed download accepted per file (512 MiB). |
DEFAULT_MAX_PAYLOAD_BYTES
|
max_inflight_bytes
|
int
|
Compressed bytes held at once before parsing (1 GiB). Neither bounds total memory: a decompressed DBF is held whole. |
DEFAULT_MAX_INFLIGHT_BYTES
|
Returns:
| Name | Type | Description |
|---|---|---|
One |
ImportReport
|
class: |
ImportReport
|
does not list it) or |
|
ImportReport
|
|
Raises:
| Type | Description |
|---|---|
ImportAbortedError
|
the lake transaction failed mid-run; its |
ValueError
|
|
Examples:
>>> import omnisus as sus
>>> scopes = sus.available("sim_obitos", years=[2023], ufs=["RR"])
>>> report = sus.import_dataset("sim_obitos", scopes=scopes)
>>> [(str(o.scope), o.status) for o in report.outcomes]
[('RR_2023', 'ok')]
Source code in src/omnisus/__init__.py
160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 | |
Import for a citable run: you choose the run_id before anything downloads.
Same as :func:~omnisus.import_dataset, but run_id is required, so an
interrupted run is reconciled through Lake.publications(run_id=...), and the
default policy never duplicates rows on a retry.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
scopes
|
Sequence[ScopeKey]
|
The scopes to import, e.g. from :func: |
required |
run_id
|
str
|
Your non-empty identifier for this run, e.g. |
required |
target
|
str | None
|
DuckLake target; |
None
|
concurrency
|
int
|
Downloads in flight; see :func: |
DEFAULT_CONCURRENCY
|
batch_size
|
int
|
Scopes committed per DuckLake transaction. |
DEFAULT_BATCH_SIZE
|
policy
|
ImportPolicy
|
|
'skip_same'
|
max_payload_bytes
|
int
|
Largest compressed download accepted per file (512 MiB). |
DEFAULT_MAX_PAYLOAD_BYTES
|
max_inflight_bytes
|
int
|
Compressed bytes held at once before parsing (1 GiB). |
DEFAULT_MAX_INFLIGHT_BYTES
|
Returns:
| Name | Type | Description |
|---|---|---|
The |
ImportReport
|
class: |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
ImportAbortedError
|
as :func: |
Examples:
>>> import omnisus as sus
>>> scopes = sus.scopes_for("sim_obitos", years=[2023], ufs=["RR"])
>>> sus.import_research("sim_obitos", scopes=scopes, run_id="cap2")
Source code in src/omnisus/research.py
66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 | |
Import IBGE population by municipality: a census or the year's latest estimate.
Each population edition is its own publication, identified by the returned
publication_id in ibge_population_manifest; it never appears in
Lake.publications(). Reconcile an interrupted run through that manifest.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
years
|
Iterable[int]
|
Years, e.g. |
required |
census
|
bool
|
Required, with no default, so nobody gets an estimate thinking it is
the census. |
required |
target
|
str | None
|
DuckLake target; |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
One |
list[ImportResult]
|
class: |
Raises:
| Type | Description |
|---|---|
ValueError
|
the year has no such edition. |
TypeError
|
|
Examples:
>>> import omnisus as sus
>>> (censo,) = sus.import_ibge_populacao(years=[2022], census=True)
>>> censo.rows
5570
Source code in src/omnisus/__init__.py
374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 | |
Fetch CNES establishment names from the public API into cnes_master.
The CNES-ST DBF has no establishment names; this pulls them from
apidadosabertos.saude.gov.br and joins them into the aux_cnes view. It
is an idempotent upsert, not a publication: no run_id, no policy, no
manifest row. To finish an interrupted run, run it again.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
codes
|
Sequence[str] | None
|
7-digit CNES codes, e.g. |
None
|
target
|
str | None
|
DuckLake target; |
None
|
concurrency
|
int
|
API requests in flight (default 5). |
5
|
only_missing
|
bool
|
|
True
|
progress
|
Callable[[int, int], None] | None
|
Optional |
None
|
Returns:
| Type | Description |
|---|---|
int
|
The number of records written in this run. |
Raises:
| Type | Description |
|---|---|
ValueError
|
a code is not seven ASCII digits, or the API returned conflicting records for one code. |
Examples:
>>> import omnisus as sus
>>> sus.import_cnes_master(codes=["2789590"])
1
Source code in src/omnisus/__init__.py
482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 | |
Import SIGTAP procedures, one competência per year x month.
Each competência is one national monthly publication from
ftp2.datasus.gov.br/public/sistemas/tup/downloads into
aux_sigtap_procedimentos. Nothing joins at import: join on
co_procedimento and the record's competência yourself.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
years
|
Iterable[int]
|
Years, e.g. |
required |
months
|
Iterable[int] | None
|
Months 1-12, e.g. |
None
|
target
|
str | None
|
DuckLake target; |
None
|
policy
|
ImportPolicy
|
|
'skip_same'
|
Returns:
| Name | Type | Description |
|---|---|---|
One |
ImportReport
|
class: |
ImportReport
|
|
Examples:
>>> import omnisus as sus
>>> report = sus.import_sigtap(years=[2024], months=[1])
>>> [(str(o.scope), o.status) for o in report.outcomes]
[('national_2024_01', 'ok')]
Source code in src/omnisus/__init__.py
422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 | |
Every SIGTAP competência the server lists now, oldest first.
Returns:
| Type | Description |
|---|---|
list[tuple[int, int]]
|
|
Raises:
| Type | Description |
|---|---|
FtpUnavailable
|
the server did not answer after retries. |
Examples:
>>> import omnisus as sus
>>> sus.available_sigtap()[:2]
[(2008, 1), (2008, 2)]
Source code in src/omnisus/__init__.py
463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 | |
Vocabularies¶
The bootstrap tables (aux_uf, aux_municipios, aux_cid10, aux_ocupacoes,
aux_paises) come from DATASUS's SIM/CID10/TABELAS; the CBO 2002 codes in
aux_ocupacoes come from CBO2002.CNV in SIM/CID10/TAB/OBITOS_CID10_TAB.zip. Each
file is hashed in the source registry. Nothing joins at import: reference_join_sql writes the LEFT JOIN a
field's dictionary declares. See vocabularies.
The LEFT JOIN from a coded field to the vocabulary its dictionary declares.
The declaration is the field's foreignKeys entry; an x-join-rule adds a
condition. Nothing joins at import: paste the result after
FROM lake.<dataset> AS <alias>.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
field
|
str
|
A field that declares a reference, e.g. |
required |
alias
|
str
|
Alias of the dataset table in your query (default |
'd'
|
lake_alias
|
str
|
Alias the lake is attached as (default |
'lake'
|
Returns:
| Type | Description |
|---|---|
str
|
SQL text; the vocabulary is aliased |
Raises:
| Type | Description |
|---|---|
ValueError
|
the field declares no reference, or more than one. |
FileNotFoundError
|
|
Examples:
>>> import omnisus as sus
>>> sus.reference_join_sql("sim_obitos", "causabas")
'LEFT JOIN "lake"."aux_cid10" AS "ref_causabas" ON d."causabas" = "ref_causabas"."codigo"'
Source code in src/omnisus/research.py
185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 | |
Dictionaries, labels and harmonised categories¶
Each dataset ships a dictionary: fields, labels, code maps and audited analytical
rules. display_row is the per-row door to the same lookup as label.
analytical_projection returns the SQL behind load's harmonised categories, for
queries you write yourself; take the schema and the sources from the snapshot you
query. See consumption.
A dataset's packaged metadata: fields, sources, analytical rules and a digest.
Works offline. schema.fields keeps the authored definitions;
fields holds one self-contained document per column with its review state.
Neither describes a live lake: use DESCRIBE at the snapshot you query.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
A new JSON-compatible dict on every call, with |
dict[str, Any]
|
|
dict[str, Any]
|
categories. |
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
|
Examples:
>>> import omnisus as sus
>>> metadata = sus.describe_dataset("sim_obitos")
>>> sorted(metadata)[:3]
['analytics', 'dataset', 'dictionary_version']
>>> len(metadata["metadata_hash"])
64
Source code in src/omnisus/metadata.py
147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 | |
Replace each code in one row by its dictionary label, without mutating the input.
The per-row door to the same lookup as :func:label: a coded field gets its
label, or None when the dictionary does not know the code; every other
field is returned as published.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
row
|
Mapping[str, Any]
|
One record as |
required |
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
A new dict with the same keys. |
Raises:
| Type | Description |
|---|---|
FileNotFoundError
|
|
Examples:
>>> import omnisus as sus
>>> sus.display_row("sim_obitos", {"sexo": "2", "dtobito": "01012023"})
{'sexo': 'Feminino', 'dtobito': '01012023'}
Source code in src/omnisus/transforms/dictionaries.py
136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 | |
SQL for the harmonised categories of a dataset, enabled only for validated sources.
Opens no lake and changes no stored record: it returns expressions to put in a
SELECT. :func:~omnisus.load already does this for you. Each category is
enabled only when every source in scopes is listed, by scope, release and
SHA-256, in the rule's validated_sources (ADR 0003); otherwise it is
reported in unavailable with the reason. Take the schema and the sources from
the same snapshot as the rows you select.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str
|
Dataset name, e.g. |
required |
observed_schema
|
Mapping[str, str]
|
|
required |
scopes
|
Sequence[SourceContext]
|
The :class: |
required |
rule_version
|
str | None
|
Require this rules edition, e.g. |
None
|
Returns:
| Type | Description |
|---|---|
AnalyticalProjection
|
|
AnalyticalProjection
|
|
AnalyticalProjection
|
|
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> import omnisus as sus
>>> dorr2023 = sus.SourceContext(
... sus.ScopeKey("RR", 2023),
... "final",
... "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
... )
>>> projection = sus.analytical_projection(
... "sim_obitos", observed_schema={"sexo": "VARCHAR"}, scopes=[dorr2023]
... )
>>> [column.name for column in projection.columns]
['sexo_categoria', 'sexo_status']
Source code in src/omnisus/transforms/analytics.py
130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 | |
Source identity, read from the same snapshot as the queried records.
Source code in src/omnisus/transforms/analytics.py
19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 | |
from_publication(row)
classmethod
¶
The source identity recorded by one Lake.publications() row.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
row
|
Mapping[str, Any]
|
A publication row with |
required |
Returns:
| Type | Description |
|---|---|
SourceContext
|
Scope, release (from the file's directory) and aggregate SHA-256. |
Raises:
| Type | Description |
|---|---|
ValueError
|
the row lacks one of those keys or has no supported scope. |
Examples:
>>> import omnisus as sus
>>> sus.SourceContext.from_publication({
... "dataset": "sim_obitos",
... "scope_json": '{"ano": 2023, "uf": "RR"}',
... "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
... "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
... }).release
'final'
Source code in src/omnisus/transforms/analytics.py
27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 | |
Results¶
ImportResult.bytes_written measures the temporary staging Parquet file, not
final lake storage growth. snapshot_id=None can mean an uncommitted result or
an unavailable snapshot ID; it does not alone establish whether a write committed.
Managed FTP results also carry run_id, batch_id and publication_id; IBGE
results carry their canonical publication_id.
Per-scope outcomes of one import run.
Normal FTP completion reports each requested input position. An ImportAbortedError carries a partial report of determined outcomes and separately identifies unresolved positions; inspect before retrying.
Inspect :attr:failed — never the report's truthiness. An empty run and a
run where everything failed are different facts, and no falsy sentinel
stands in for either.
Source code in src/omnisus/sources/_base.py
133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 | |
rows
property
¶
Total rows ingested across successful scopes.
What happened to one scope in an import run.
skipped and failed are different facts and must not be collapsed.
skipped means the listing does not have the scope (not_listed), it
is outside the row's coverage (outside_coverage), or the same source is
already published (unchanged). failed is worth retrying.
Source code in src/omnisus/sources/_base.py
111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 | |
Outcome of one import_scope call.
Source code in src/omnisus/sources/_base.py
90 91 92 93 94 95 96 97 98 99 100 | |
Identifies one slice of a dataset (e.g., SP/2024 or MG/2024/01).
mes is None for yearly datasets, set for monthly (SIH).
Source code in src/omnisus/sources/_base.py
72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 | |
uf
instance-attribute
¶
None denotes a national source, never an artificial record-level UF.
What :func:delete_scope removed: data rows, and the manifest rows retired with them.
Source code in src/omnisus/lake/publication.py
50 51 52 53 54 55 | |
The lake¶
Lake.local() and LakeReader() without a target open the default lake, in the
folder set by set_lake_dir (where the lake lives).
Pass a target such as ducklake:/path/omnisus.ducklake for another lake.
Local handles enforce a cooperative writer lock for their lifetime.
Use Lake.transaction() for managed writes; raw SQL
transaction control is outside this contract. Managed transactions cannot nest.
To read, open a LakeReader on the same target: it attaches the catalog
read-only, takes no lock, creates nothing and sets no option, so it runs
alongside an import. Pass snapshot_id to pin the session; without it every
statement reads the latest committed snapshot.
import omnisus as sus
with sus.LakeReader() as reader:
latest = reader.snapshots()[-1]["snapshot_id"]
rows = reader.connect().execute("SELECT count(*) FROM lake.sim_obitos").fetchone()
with sus.LakeReader(snapshot_id=latest) as reader:
... # every statement here sees exactly that snapshot
import omnisus as sus
import polars as pl
with sus.Lake.local() as lake:
with lake.transaction() as receipt:
result = lake.ingest("example", pl.DataFrame({"id": [1]}).lazy())
assert result.snapshot_id is None
assert receipt.committed
print(receipt.snapshot_id, result.rows)
The receipt's snapshot_id may remain None after a successful commit when no new snapshot was created or the lookup was unavailable.
See Architecture for rollback and recovery boundaries.
Lake.publications(run_id=...) reads the durable source-publication manifest;
each row carries scope, the ScopeKey the package wrote (None for a shape
this version does not write), alongside the raw scope_json.
Lake.attempts(run_id=...) reads separately recorded known failures.
Lake.ingest_parquet appends a staging file directly. Lake.publish_scope adds
scope validation, source identity and replay policy to that write.
Lake.delete_scope(table, scope) removes one source scope and retires every
publication within it in the same transaction; a yearly scope on a monthly
table covers all its months. import_ibge_populacao and import_cnes_master do not
take part in this manifest — see their docstrings for how each reconciles.
Lake.optimize(table) merges adjacent files. Lake.expire_snapshots and
Lake.cleanup_files take older_than as a timezone-aware datetime and default
to dry_run=True. Each returns a list of result dictionaries.
Keep the default lake in path for the rest of this session.
Every call without target (load, LakeReader, import_*) then reads
and writes the lake there. On Colab, point it at the mounted Google Drive so the
lake survives the runtime. It sets $OMNISUS_DATA_DIR, which scripts can set
directly instead. The folder is created by the first import.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
path
|
str | PathLike[str]
|
The lake's folder; |
required |
Returns:
| Type | Description |
|---|---|
Path
|
The absolute folder. |
Examples:
>>> import omnisus as sus
>>> sus.set_lake_dir("/content/drive/MyDrive/omnisus")
PosixPath('/content/drive/MyDrive/omnisus')
Source code in src/omnisus/lake/catalog.py
81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 | |
Bases: Session
Writer handle on a DuckLake.
Use :meth:local with a ducklake: target or :meth:cloud for a
Postgres-backed catalog. Managed ingestion requires one writer per lake,
including when the catalog is hosted in Postgres. To only read, open a
:class:~omnisus.lake.session.LakeReader instead: it takes no lock.
Source code in src/omnisus/lake/operations.py
33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 | |
bootstrap_auxiliares()
¶
Load the aux_* vocabularies (UF, municipalities, CID-10, CBO, countries).
They come from the package's bootstrap zip, each file hashed in the source registry. Re-running replaces the tables' contents.
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.bootstrap_auxiliares()
Source code in src/omnisus/lake/operations.py
293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 | |
cleanup_files(*, older_than, dry_run=True)
¶
Delete data files no snapshot needs, older than a cutoff; simulates by default.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
older_than
|
datetime
|
Timezone-aware |
required |
dry_run
|
bool
|
|
True
|
Returns:
| Type | Description |
|---|---|
list[dict]
|
The files deleted, or that would be. |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> from datetime import UTC, datetime
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.cleanup_files(older_than=datetime(2026, 1, 1, tzinfo=UTC))
Source code in src/omnisus/lake/operations.py
264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 | |
close()
¶
Close the connection and release the writer lock.
Examples:
>>> import omnisus as sus
>>> lake = sus.Lake.local()
>>> lake.close()
Source code in src/omnisus/lake/operations.py
620 621 622 623 624 625 626 627 628 629 630 631 632 633 | |
cloud(*, catalog, storage)
classmethod
¶
Open a lake whose catalog is in PostgreSQL and data in object storage.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
catalog
|
str
|
|
required |
storage
|
str
|
Data path, e.g. |
required |
Returns:
| Type | Description |
|---|---|
Self
|
An open :class: |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
CatalogAttachError
|
DuckDB could not attach the catalog. |
Examples:
>>> import omnisus as sus
>>> lake = sus.Lake.cloud(
... catalog="postgresql://u:p@db/lake", storage="s3://bucket/lake"
... )
Source code in src/omnisus/lake/operations.py
96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 | |
delete_scope(table, scope)
¶
Delete one scope's rows and retire its publications, in one transaction.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
table
|
str
|
Dataset table, e.g. |
required |
scope
|
ScopeKey
|
The scope to delete, e.g. |
required |
Returns:
| Type | Description |
|---|---|
DeletionResult
|
Rows deleted and publications retired. |
Raises:
| Type | Description |
|---|---|
ValueError
|
the table's national/state shape does not match |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.delete_scope("sim_obitos", sus.ScopeKey("RR", 2023))
Source code in src/omnisus/lake/operations.py
598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 | |
ensure_aux_cnes_view()
¶
Create or refresh aux_cnes: one row per CNES establishment.
Every import of cnes_estabelecimentos already calls this. The view joins
the latest competência of cnes_estabelecimentos with cnes_master
names: cnes (7 digits), nome (NULL until
:func:~omnisus.import_cnes_master runs), tp_unid, codufmun and
yyyymm_max.
Returns:
| Type | Description |
|---|---|
bool
|
|
Raises:
| Type | Description |
|---|---|
RuntimeError
|
the handle is closed or unusable. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.ensure_aux_cnes_view()
True
Source code in src/omnisus/lake/operations.py
329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 | |
expire_snapshots(*, older_than, dry_run=True)
¶
Forget snapshots older than a cutoff; only simulates unless told otherwise.
Expired snapshots can no longer be read or cited.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
older_than
|
datetime
|
Timezone-aware |
required |
dry_run
|
bool
|
|
True
|
Returns:
| Type | Description |
|---|---|
list[dict]
|
The snapshots expired, or that would be. |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> from datetime import UTC, datetime
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.expire_snapshots(older_than=datetime(2026, 1, 1, tzinfo=UTC))
Source code in src/omnisus/lake/operations.py
233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 | |
ingest(table, lazyframe, *, partition_by=())
¶
Append a polars LazyFrame to a table, without a publication record.
Outside :meth:transaction it commits before returning; inside, the result's
snapshot_id stays None until the context commits.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
table
|
str
|
Destination table, e.g. |
required |
lazyframe
|
LazyFrame
|
The rows to materialize and insert. |
required |
partition_by
|
tuple[str, ...]
|
Partition columns, applied once when the table is created. |
()
|
Returns:
| Type | Description |
|---|---|
ImportResult
|
Rows inserted, staging bytes and duration. |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
the handle is closed or unusable. |
Examples:
>>> import polars as pl
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.ingest("minha_tabela", pl.LazyFrame({"x": [1, 2]}))
Source code in src/omnisus/lake/operations.py
482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 | |
ingest_parquet(table, staging, *, partition_by=())
¶
Append a validated Parquet staging file to a table, without rewriting it.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
table
|
str
|
Destination table, e.g. |
required |
staging
|
str | Path
|
Path of the Parquet file. |
required |
partition_by
|
tuple[str, ...]
|
Partition columns, applied once when the table is created. |
()
|
Returns:
| Type | Description |
|---|---|
ImportResult
|
Rows inserted, the staging file's size and duration. |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
the handle is closed or unusable. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.ingest_parquet("minha_tabela", "staging.parquet")
Source code in src/omnisus/lake/operations.py
527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 | |
local(target=None)
classmethod
¶
Open, or create, a lake on local disk for writing.
One writer per lake: a second handle opening it for writing, in this process
or another, fails at once with WriterBusyError. Use
:class:~omnisus.LakeReader to read while another process writes.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
target
|
str | None
|
|
None
|
Returns:
| Type | Description |
|---|---|
Self
|
An open :class: |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
WriterBusyError
|
another handle already has this lake open for writing
( |
CatalogAttachError
|
DuckDB could not attach the catalog. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.bootstrap_auxiliares()
Source code in src/omnisus/lake/operations.py
66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 | |
optimize(table)
¶
Merge a table's small data files, keeping every historical snapshot.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
table
|
str
|
Table name, e.g. |
required |
Returns:
| Type | Description |
|---|---|
list[dict]
|
One dict per maintenance step DuckLake reports. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.optimize("sim_obitos")
Source code in src/omnisus/lake/operations.py
215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 | |
publish_scope(table, staging, **kwargs)
¶
Publish one scope from a staging file, recording its source and policy.
The low-level step :func:~omnisus.import_dataset runs per scope.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
table
|
str
|
Dataset table, e.g. |
required |
staging
|
Path
|
Path of the validated Parquet staging file. |
required |
**kwargs
|
Any
|
|
{}
|
Returns:
| Type | Description |
|---|---|
ImportResult | None
|
The result, or |
Raises:
| Type | Description |
|---|---|
ValueError
|
the staging file does not match the scope or the table. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake:
... lake.publish_scope("sim_obitos", "staging.parquet", scope=scope, ...)
Source code in src/omnisus/lake/operations.py
571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 | |
transaction()
¶
Group writes into one DuckLake transaction, hence one snapshot.
Yields:
| Type | Description |
|---|---|
TransactionReceipt
|
A receipt whose |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
already inside a transaction on this handle. |
TransactionStateError
|
the transaction could not begin or end cleanly; the handle becomes unusable, so close it and inspect the catalog. |
Examples:
>>> import omnisus as sus
>>> with sus.Lake.local() as lake, lake.transaction() as receipt:
... lake.ensure_aux_cnes_view()
Source code in src/omnisus/lake/operations.py
131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 | |
Bases: Session
Read-only session on an existing lake.
target uses the same grammar as :meth:Lake.local —
ducklake:./omnisus.ducklake or ducklake:postgresql://…?storage=… —
and the storage part is ignored: the catalog records its own data path.
Nothing is created and no option is set, so a reader never takes the
writer lock and never blocks a writer; DuckLake refuses writes on the
connection itself. snapshot_id pins the session to one snapshot
(snapshots()[-1]["snapshot_id"] is the latest); without it every
statement reads the latest committed snapshot.
Raises :class:~omnisus.lake.CatalogAttachError when the catalog does
not exist or the pinned snapshot is unknown.
Source code in src/omnisus/lake/session.py
184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 | |
snapshot_id = snapshot_id
instance-attribute
¶
The pinned snapshot, or None for a session that reads the latest one.
Research citations and joins¶
import_research and cite are above. cite reads the publication manifest (or
ibge_population_manifest) and returns Portuguese text matching the
reproducibility guide. latest_snapshot_id
is the newest catalog snapshot. Municipality helpers take the leftmost 6 or 7
digits; they do not pad and they do not rewrite stored columns.
Structured citation plus the Portuguese paragraph the guide uses.
Source code in src/omnisus/research.py
28 29 30 31 32 33 34 35 36 37 | |
Format publication rows you already have, e.g. kept in a JSON file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
publications
|
Sequence[Mapping[str, Any]]
|
Rows of |
required |
snapshot_id
|
int
|
The lake snapshot the analysis read. |
required |
dataset
|
str | None
|
Dataset name, used when a row lacks it; |
None
|
run_id
|
str | None
|
The run to name; |
None
|
accessed
|
date | None
|
Access date to print; |
None
|
Returns:
| Name | Type | Description |
|---|---|---|
A |
Citation
|
class: |
Examples:
>>> from datetime import date
>>> import omnisus as sus
>>> row = {
... "dataset": "sim_obitos",
... "source_uri": "ftp://ftp.datasus.gov.br/dissemin/publicos/SIM/CID10/DORES/DORR2023.dbc",
... "source_sha256": "15b5203507161b7c35f9c69a52c955bdac133629549099b85a88e03a9a45baf0",
... "run_id": "cap2",
... }
>>> print(sus.citation_from_publications([row], snapshot_id=7, accessed=date(2026, 9, 23)).text)
SIM — Declarações de Óbito (sim_obitos), arquivo DORR2023.dbc (...), SHA-256 15b5...baf0, acessado em 2026-09-23 pelo DATASUS. Importado com omnisus ..., lake snapshot 7, execução cap2.
Source code in src/omnisus/research.py
229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 | |
The newest snapshot of the lake, to pin a query or a citation to it.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
lake
|
Lake | LakeReader
|
An open :class: |
required |
Returns:
| Type | Description |
|---|---|
int
|
The snapshot id; pass it to |
Raises:
| Type | Description |
|---|---|
LookupError
|
the lake has no snapshots yet. |
Examples:
>>> import omnisus as sus
>>> with sus.LakeReader() as lake:
... snapshot = sus.latest_snapshot_id(lake)
Source code in src/omnisus/research.py
40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 | |
The leftmost digits of an IBGE municipality code, to join bases on them.
DATASUS bases use 6 digits and IBGE uses 7 (with a check digit). This keeps the
leftmost digits and never pads or invents a check digit.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
value
|
object
|
A municipality code, e.g. |
required |
digits
|
int
|
|
6
|
Returns:
| Type | Description |
|---|---|
str | None
|
The key as text, or |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> import omnisus as sus
>>> sus.municipality_join_key("3550308")
'355030'
>>> sus.municipality_join_key(" ") is None
True
Source code in src/omnisus/research.py
130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 | |
SQL for :func:municipality_join_key on a lake column.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
column
|
str
|
Column name, e.g. |
required |
digits
|
int
|
|
6
|
Returns:
| Type | Description |
|---|---|
str
|
A DuckDB expression to use in |
Raises:
| Type | Description |
|---|---|
ValueError
|
|
Examples:
>>> import omnisus as sus
>>> sus.municipality_join_key_sql("codmunres")
'left(trim(CAST("codmunres" AS VARCHAR)), 6)'
Source code in src/omnisus/research.py
162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 | |
Registry¶
datasets() lists every curated FTP dataset; products() covers every importer
family except SIGTAP: the FTP datasets and the two families outside the registry
(ibge_populacao, cnes_master). It states, per family, the scope fields, accepted
policies, how an interrupted run is reconciled and whether available() applies. Year rules for IBGE remain in
omnisus.sources.ibge.products (CENSUS_YEARS, ESTIMATE_UNAVAILABLE_YEARS);
an estimate is importable only as its latest edition, so there is no floor year.
Identity and location of one DATASUS-FTP dataset.
Holds only identity and location. Anything behavioural belongs in an importer module, never as a flag here.
Source code in src/omnisus/sources/datasus_ftp/datasets.py
31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 | |
cadence
instance-attribute
¶
How DATASUS publishes files. Decides filename shape (upstream fact).
coverage
instance-attribute
¶
(first, last) published; last=None means ongoing.
Provisional until the Tier 3 probe validates it against the server.
dictionary = None
class-attribute
instance-attribute
¶
Frictionless YAML. None -> packaged dicionarios/<name>.yaml.
ftp_dir
instance-attribute
¶
Directory on ftp.datasus.gov.br holding this dataset's files.
geography = 'state'
class-attribute
instance-attribute
¶
Source coverage. National SINAN filenames use PREFIX + BR + YY.
name
instance-attribute
¶
Registry key = lake table = YAML stem = CLI name (sim_obitos).
national_code = 'BR'
class-attribute
instance-attribute
¶
What a national file name carries between prefix and year: BR for SINAN
(CHAGBR23.dbc), nothing for the SIM subsets (DOEXT23.dbc).
partition_by
instance-attribute
¶
How the lake table is laid out (our storage policy).
Applied by Lake.ingest via SET PARTITIONED BY when the table is
created, so files land under ano=…/uf=…/. Each scope is already one
(uf, ano[, mes]), so partitioning costs nothing at write time and buys
query pruning.
prefix
instance-attribute
¶
DATASUS filename prefix: DO, DN, RD, BI, ATD…
prelim_dir = None
class-attribute
instance-attribute
¶
Directory where DATASUS publishes this dataset's preliminary files,
when it has one. Same filenames as ftp_dir; a file is in one or the
other, never both.
year_digits = 4
class-attribute
instance-attribute
¶
Year width in a yearly state file name: 4 (DORR2023.dbc) or 2 (SIM CID-9's
DORRR79.DBC). Monthly and national names always carry two digits.
directories()
¶
Every FTP directory this dataset is published in, final first.
Returns:
| Type | Description |
|---|---|
dict[Release, str]
|
|
dict[Release, str]
|
files elsewhere. |
Examples:
>>> import omnisus as sus
>>> sus.resolve("sim_obitos").directories()["prelim"]
'/dissemin/publicos/SIM/PRELIM/DORES'
Source code in src/omnisus/sources/datasus_ftp/datasets.py
88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 | |
The :class:Dataset for a name; a Dataset value passes through.
A value passing through untouched is how an uncurated dataset reaches the pipeline.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
dataset
|
str | Dataset
|
Dataset name, e.g. |
required |
Returns:
| Type | Description |
|---|---|
Dataset
|
The registry row: name, FTP directories, geography, monthly or yearly. |
Raises:
| Type | Description |
|---|---|
ValueError
|
the name is not in the registry. |
Examples:
>>> import omnisus as sus
>>> d = sus.resolve("sih_aih_reduzida")
>>> (d.geography, d.monthly)
('state', True)
Source code in src/omnisus/sources/datasus_ftp/datasets.py
178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 | |
Every DATASUS FTP dataset the package curates, in registry order.
Returns:
| Type | Description |
|---|---|
Dataset
|
class: |
...
|
and the other functions accept. |
Examples:
>>> import omnisus as sus
>>> "sim_obitos" in [d.name for d in sus.datasets()]
True
Source code in src/omnisus/products.py
45 46 47 48 49 50 51 52 53 54 55 56 57 | |
Every importer family except SIGTAP, stated once: the FTP registry, IBGE and CNES master.
Facts, not forms. Year rules for IBGE stay in sources.ibge.products
(CENSUS_YEARS, ESTIMATE_UNAVAILABLE_YEARS); an estimate is importable
only as its latest edition, checked live at import, so there is no floor year.
Product
dataclass
¶
One importer family and what it supports.
Source code in src/omnisus/products.py
19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 | |
dataset
instance-attribute
¶
The FTP registry row, or None for IBGE and CNES master.
inventory
instance-attribute
¶
Whether :func:omnisus.available can list the source's scopes.
name
instance-attribute
¶
Registry key or importer name: sim_obitos, ibge_populacao, cnes_master.
policies
instance-attribute
¶
Accepted ImportPolicy values; families outside the manifest append only.
reconcile_by
instance-attribute
¶
How an interrupted run is reconciled: Lake.publications(run_id=), the
returned publication_id in ibge_population_manifest, or a rerun
(CNES master is an idempotent upsert).
scope_fields
instance-attribute
¶
What identifies one import unit: ("uf", "ano"[, "mes"]), ("ano",)
for national datasets, ("product", "ano") for IBGE, empty for CNES master.
datasets()
¶
Every DATASUS FTP dataset the package curates, in registry order.
Returns:
| Type | Description |
|---|---|
Dataset
|
class: |
...
|
and the other functions accept. |
Examples:
>>> import omnisus as sus
>>> "sim_obitos" in [d.name for d in sus.datasets()]
True
Source code in src/omnisus/products.py
45 46 47 48 49 50 51 52 53 54 55 56 57 | |
products()
¶
Every importer family except SIGTAP: the FTP datasets, IBGE population, CNES master.
Returns:
| Name | Type | Description |
|---|---|---|
One |
Product
|
class: |
...
|
policies and how an interrupted run is reconciled. |
Examples:
>>> import omnisus as sus
>>> [(p.name, p.reconcile_by) for p in sus.products()[-2:]]
[('ibge_populacao', 'publication_id'), ('cnes_master', 'rerun')]
Source code in src/omnisus/products.py
60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 | |
One importer family and what it supports.
Source code in src/omnisus/products.py
19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 | |
dataset
instance-attribute
¶
The FTP registry row, or None for IBGE and CNES master.
inventory
instance-attribute
¶
Whether :func:omnisus.available can list the source's scopes.
name
instance-attribute
¶
Registry key or importer name: sim_obitos, ibge_populacao, cnes_master.
policies
instance-attribute
¶
Accepted ImportPolicy values; families outside the manifest append only.
reconcile_by
instance-attribute
¶
How an interrupted run is reconciled: Lake.publications(run_id=), the
returned publication_id in ibge_population_manifest, or a rerun
(CNES master is an idempotent upsert).
scope_fields
instance-attribute
¶
What identifies one import unit: ("uf", "ano"[, "mes"]), ("ano",)
for national datasets, ("product", "ano") for IBGE, empty for CNES master.
Errors¶
Lake transaction errors are imported from omnisus.lake. A
CommitOutcomeUnknown means the COMMIT raised and the handle is unusable;
inspect the catalog before retrying. The FTP runner wraps transaction state
failures in the top-level ImportAbortedError with partial progress.
CatalogAttachError means the catalog could not be opened at all: .stage
says which statement failed, and for a remote catalog the message never
carries the connection string.
Bases: RuntimeError
Partial progress is known, but the remaining inputs need inspection.
Source code in src/omnisus/sources/_base.py
167 168 169 170 171 172 173 174 175 176 177 178 179 | |
Bases: RuntimeError
The DuckLake catalog could not be opened, so there is no lake to use.
stage names the statement that failed: "install" (the ducklake
extension), "attach" (the catalog itself) or "set_option"; the
message adds the DuckDB error class. For a remote catalog the DuckDB error
is not chained — its text can echo the connection string — so stage and
class are the whole diagnostic. For a local catalog it is chained, so the
file-level reason stays visible.
Source code in src/omnisus/lake/connection.py
15 16 17 18 19 20 21 22 23 24 25 26 27 28 | |
Bases: RuntimeError
The handle cannot safely continue its managed transaction workflow.
Source code in src/omnisus/lake/_transactions.py
10 11 | |
Bases: TransactionStateError
COMMIT raised; do not infer rollback or retry the write automatically.
Source code in src/omnisus/lake/_transactions.py
14 15 | |
Bases: Exception
The remote directory does not exist, or access was denied (550).
Terminal — never retried. Distinct from an existing but empty directory,
which returns an empty :class:Listing.
Source code in src/omnisus/sources/datasus_ftp/inventory.py
42 43 44 45 46 47 | |