Fast parallel extraction for UniVerse databases using uopy
Project description
universe-reads
Fast parallel extraction for UniVerse databases using uopy. Reads large volumes of records with 1 session by default (configurable up to N), multiple output formats, streaming, record filtering, key-range partitioning, and file introspection.
Features
- Parallel reads — configurable number of concurrent
uopysessions via multiprocessing (spawn) - Output formats — NDJSON (default), CSV, and Parquet (via optional
pyarrow) - Streaming — pipe records directly to stdout (
--stream) for Unix pipeline integration - Record filtering — server-side
WHERE/SELECTclauses pushed to UniVerse - Key-range partitioning —
prefix(ordinal) orhashstrategies for balanced worker loads - File inspection —
universe-reads inspect FILEfor quick metadata, field counts, and sample records - Enhanced metrics — per-worker tabular summary with peak/average rates after every run
- Retry with backoff — configurable exponential backoff for transient RPC failures
- Dry-run mode — validate parallelism and sinks without a UniVerse connection
Quick start
1) Install
python -m venv .venv
source .venv/bin/activate
pip install -e .
With optional dependencies:
# uopy (if available on PyPI)
pip install -e ".[uopy]"
# Parquet output support
pip install -e ".[parquet]"
# Everything
pip install -e ".[all]"
If uopy is not available on PyPI in your environment, install it using your internal method and keep the import name uopy.
For a quick run without installing the package:
PYTHONPATH=src python -m universe_reads --dry-run 10000 --count-only
2) Configure
universe-reads itself only needs:
UV_FILE(optional if you pass--file)
You provide the UniVerse session creation via --session-factory.
Default: If you omit
--session-factory, it defaults touniverse_reads.session_factory:make_session— the built-in env-based factory that readsUV_HOST,UV_USER,UV_PASSWORD, etc. from environment variables. See Built-in env-based session factory below. For most users, setting the env vars and omitting--session-factoryis the quickest way to get started.
Session factory
By default, --session-factory is set to universe_reads.session_factory:make_session,
which is a built-in factory that creates sessions from environment variables
(UV_HOST, UV_USER, UV_PASSWORD, etc.). Most users can skip
--session-factory entirely and just set the required env vars.
If you need custom session logic (e.g. different credentials, connection
pooling, or a non-standard uopy setup), you can point --session-factory
at your own callable.
This project uses multiprocessing with the spawn start method, which means worker processes must be able to import whatever they need.
The recommended way to specify a session factory is a string import spec:
"module.submodule:callable_name"
This is robust across processes and lets you reuse the same factory reference in many modules.
A) Using the built-in factory (recommended)
The simplest approach — just set environment variables and run:
export UV_HOST=myhost UV_USER=myuser UV_PASSWORD=mypass UV_ACCOUNT=myaccount
UV_FILE=READS universe-reads --count-only
No --session-factory flag needed. The built-in factory is used automatically.
B) Custom factory module
Example factory module my_uv.py:
import uopy
def make_session():
return uopy.connect(host="...", account="...", user="...", password="...")
Then use it from the CLI:
UV_FILE=READS universe-reads --session-factory my_uv:make_session --count-only
Or programmatically:
from universe_reads import run
results = run(
file_name="READS",
session_factory="my_uv:make_session",
count_only=True,
)
C) Why not pass a lambda/closure?
Avoid nested functions / lambdas for session_factory when using multiprocessing spawn.
Do NOT do:
from universe_reads import run
def make_factory():
def factory():
...
return factory
run(file_name="READS", session_factory=make_factory())
Prefer a top-level function in an importable module and pass it as a string spec.
D) Reuse across modules
If you'll use this in a bunch of places, put the spec string in one location:
# settings.py
SESSION_FACTORY = "my_uv:make_session"
Then:
from settings import SESSION_FACTORY
from universe_reads import run
run(file_name="READS", session_factory=SESSION_FACTORY)
Or use the built-in env-based factory (this is also the default, so the flag is optional):
--session-factory universe_reads.session_factory:make_session
E) Pooling kwargs are optional (factory-friendly)
The CLI/library can pass pooling hints as keyword args:
pooling_onmin_pool_sizemax_pool_size
If your factory does not accept them, they are safely ignored.
Example of a kw-accepting factory:
import uopy
def make_session(*, pooling_on=None, min_pool_size=None, max_pool_size=None):
kwargs = {"user": "...", "password": "..."}
if pooling_on:
kwargs["min_pool_size"] = min_pool_size or 1
kwargs["max_pool_size"] = max_pool_size or 10
return uopy.connect(**kwargs)
3) Run
Count-only (fastest):
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--count-only
Write NDJSON files (one per worker):
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--output-dir out
Write CSV files:
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--output-dir out \
--format csv
Write Parquet files (requires pip install universe-reads[parquet]):
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--output-dir out \
--format parquet
Stream records to stdout as NDJSON (for piping):
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--stream | jq .
Stream as CSV:
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--stream --format csv | head -100
More sessions (if you have licenses/CPU to spare):
UV_FILE=READS universe-reads \
--session-factory my_uv:make_session \
--sessions 10 \
--count-only
Dry run (no UniVerse connection) to validate parallelism:
universe-reads --dry-run 300000 --count-only
Smoke-test the session factory mechanism without needing UniVerse:
PYTHONPATH=src:. python -m tests.test_dry_run --smoke-session-factory
Output formats
| Format | Flag | Extension | Notes |
|---|---|---|---|
| NDJSON | --format ndjson |
.ndjson |
Default. One JSON object per line |
| JSONL | --format jsonl |
.ndjson |
Alias for ndjson |
| CSV | --format csv |
.csv |
Header row + data rows |
| Parquet | --format parquet |
.parquet |
Requires pyarrow>=12 |
When --output-dir is set, one file per worker is written: worker-0.ndjson, worker-1.csv, etc.
When --stream is set, records are written to stdout in real time (NDJSON or CSV only). This is ideal for piping into jq, grep, wc -l, or another process.
Record filtering
Push filtering to the UniVerse server to reduce the number of records transferred:
# WHERE clause appended to SELECT <file>
universe-reads --session-factory my_uv:make_session \
--where "WITH STATUS = 'ACTIVE'" \
--output-dir out
# Full SELECT override
universe-reads --session-factory my_uv:make_session \
--select "SELECT STUDENTS WITH GPA >= '3.0'" \
--output-dir out
Library equivalent:
results = run(
file_name="STUDENTS",
session_factory="my_uv:make_session",
where_clause="WITH STATUS = 'ACTIVE'",
)
# Or with a full SELECT override
results = run(
file_name="STUDENTS",
session_factory="my_uv:make_session",
select_override="SELECT STUDENTS WITH GPA >= '3.0'",
)
Key-range partitioning
By default, keys are distributed round-robin across workers. For skewed data, choose a partitioning strategy:
| Strategy | Flag | Description |
|---|---|---|
| prefix | --partition prefix |
Assigns keys by the ordinal value of their first character |
| hash | --partition hash |
Assigns keys by a deterministic hash for even distribution |
universe-reads --session-factory my_uv:make_session \
--partition hash \
--sessions 3 \
--output-dir out
Library equivalent:
results = run(
file_name="READS",
session_factory="my_uv:make_session",
partition="hash",
sessions=3,
)
File inspection
Quickly inspect a UniVerse file without doing a full read — see record count, field counts, and sample records:
# Human-readable summary to stderr
universe-reads inspect ATTENDANCE \
--session-factory my_uv:make_session
# Machine-readable JSON to stdout
universe-reads inspect ATTENDANCE \
--session-factory my_uv:make_session \
--json
# Control the number of samples
universe-reads inspect ATTENDANCE \
--session-factory my_uv:make_session \
--samples 10
Library equivalent:
from universe_reads import inspect
info = inspect(
file_name="ATTENDANCE",
session_factory_spec="my_uv:make_session",
sample_count=5,
)
print(info)
# {'file': 'ATTENDANCE', 'record_count': 312000, 'sample_keys': [...], ...}
Enhanced metrics
After every run, a tabular summary is printed showing per-worker stats, peak rate, and average rate:
==============================================================
Worker Stats
--------------------------------------------------------------
worker records errors rate elapsed
--------------------------------------------------------------
w-0 60,000 0 780,000/s 0.1s
w-1 60,000 0 750,000/s 0.1s
w-2 60,000 0 810,000/s 0.1s
--------------------------------------------------------------
TOTAL 180,000 0 900,000/s 0.2s
Peak worker rate: 810,000 rec/s
Avg worker rate: 780,000 rec/s
==============================================================
Access the summary programmatically:
from universe_reads.metrics import format_summary
summary = format_summary(results, wall_seconds=elapsed)
print(summary)
Connection pooling (optional)
If your uopy environment has pooling enabled (often via uopy.ini), you can keep it optional by enabling it only in your session factory.
Important caveats:
- Pooling is per Python process. Since this project uses multiprocessing for parallel sessions, each worker process has its own pool.
- Pooling is not "multi-user" by itself. A pool is typically keyed by connection parameters (host/account/user/password). Different users generally mean separate pools.
With the built-in factory you can enable pooling hints via CLI flags:
universe-reads \
--session-factory universe_reads.session_factory:make_session \
--pooling \
--min-pool-size 1 \
--max-pool-size 10 \
--count-only
Or via env vars:
export UV_POOLING_ON=1
export UV_MIN_POOL_SIZE=1
export UV_MAX_POOL_SIZE=10
universe-reads --session-factory universe_reads.session_factory:make_session --count-only
Use as a library
from universe_reads import run
# Basic parallel read (single session by default)
results = run(
file_name="READS",
session_factory="my_uv:make_session",
count_only=True,
)
# With output format, filtering, and partitioning
results = run(
file_name="ATTENDANCE",
session_factory="my_uv:make_session",
sessions=3,
output_dir="out",
output_format="csv",
where_clause="WITH STATUS = 'ACTIVE'",
partition="hash",
)
# Stream to stdout
results = run(
file_name="READS",
session_factory="my_uv:make_session",
stream=True,
output_format="ndjson",
)
# Dry-run (no UniVerse needed)
results = run(file_name="", dry_run=300000)
Arguments
CLI (universe-reads)
Read mode (default)
| Flag | Default | Description |
|---|---|---|
--file |
UV_FILE env |
UniVerse file name to read |
--session-factory |
universe_reads.session_factory:make_session |
Session factory in module.sub:callable form |
--sessions |
1 |
Number of concurrent sessions/processes |
--count-only |
off (implied when no --output-dir and no --stream) |
Count records only |
--output-dir |
unset | Write one file per worker into this directory |
--format |
ndjson |
Output format: ndjson, jsonl, csv, parquet |
--stream |
off | Stream records to stdout (NDJSON or CSV) |
--where |
unset | UniVerse WHERE/WITH clause, e.g. "WITH STATUS = 'ACTIVE'" |
--select |
unset | Full UniVerse SELECT statement override |
--partition |
unset (round-robin) | Key partitioning strategy: prefix or hash |
--batch-size |
500 |
Keys per worker IPC batch |
--progress-every |
5.0 |
Seconds between throughput logs |
--max-records |
unset | Debug cap |
--dry-run N |
unset | Generate N fake keys/records instead of connecting |
--retries |
3 |
Retries per batch/key on RPC failure (0 disables) |
--retry-backoff |
0.5 |
Base delay in seconds between retries (doubles each attempt) |
--pooling / --no-pooling |
neither | Enable/disable pooling hints |
--min-pool-size |
unset | Minimum pool size hint |
--max-pool-size |
unset | Maximum pool size hint |
Inspect sub-command
universe-reads inspect FILE [--session-factory SPEC] [--samples N] [--json]
| Flag | Default | Description |
|---|---|---|
FILE |
(required) | UniVerse file name to inspect |
--session-factory |
universe_reads.session_factory:make_session |
Session factory spec |
--samples |
3 |
Number of sample records to read |
--json |
off | Output as JSON to stdout |
Library (universe_reads.run)
run() accepts the same knobs as the CLI, as keyword args:
file_namesession_factory(spec string or top-level callable; defaults to"universe_reads.session_factory:make_session", optional whendry_runis set)sessions=1count_only=Trueoutput_dir=Noneoutput_format="ndjson"—"ndjson"|"jsonl"|"csv"|"parquet"stream=Falsewhere_clause=Noneselect_override=Nonepartition=None—"prefix"|"hash"|Nonebatch_size=500progress_every_seconds=5.0max_records=Nonedry_run=Nonedry_run_data=None—{key: record}fixtures for testingretries=3,retry_backoff=0.5pooling_on=None,min_pool_size=None,max_pool_size=None
Library (universe_reads.inspect)
from universe_reads import inspect
info = inspect(
file_name="FILE",
session_factory_spec="module:factory",
sample_count=3,
)
Returns a dict with keys: file, record_count, sample_keys, field_counts, sample_records, elapsed_seconds.
Built-in env-based session factory
If you use:
--session-factory universe_reads.session_factory:make_session
It supports these environment variables:
- Required:
UV_USER,UV_PASSWORD - Optional:
UV_HOST,UV_ACCOUNT,UV_PORT,UV_TIMEOUT,UV_SERVICE,UV_ENCODING,UV_SSL - Optional pooling hints:
UV_POOLING_ON,UV_MIN_POOL_SIZE,UV_MAX_POOL_SIZE
Testing
Run the test suite without a UniVerse connection:
PYTHONPATH=src:. python -m pytest tests/ -v
Individual test scripts can also be run directly:
# Dry-run harness (300K records, 5 workers)
PYTHONPATH=src:. python -m tests.test_dry_run
# Session factory smoke test
PYTHONPATH=src:. python -m tests.test_dry_run --smoke-session-factory
# Collector test (1M fixture records with attendance data)
PYTHONPATH=src:. python -m tests.test_collector
# Unit tests for individual modules
PYTHONPATH=src:. python tests/test_sinks.py
PYTHONPATH=src:. python tests/test_partitioning.py
PYTHONPATH=src:. python tests/test_streaming.py
PYTHONPATH=src:. python tests/test_metrics.py
PYTHONPATH=src:. python tests/test_inspect.py
PYTHONPATH=src:. python tests/test_cli_args.py
Where to plug in real uopy calls
All UniVerse file I/O calls are isolated in:
src/universe_reads/uopy_adapter.py
You manage session creation via --session-factory. The adapter implements:
iter_keys()— select list key streamingiter_keys_filtered()— filtered key streaming withWHERE/SELECTclausesread_record()/read_records()— single and batch record reads
Project details
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distribution
Built Distribution
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file universe_reads-0.2.8.tar.gz.
File metadata
- Download URL: universe_reads-0.2.8.tar.gz
- Upload date:
- Size: 43.6 kB
- Tags: Source
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
dfc5e39e514ccca121c3c5a698b49dc9696b2bb6bce80269bc6764cf62a2ee3c
|
|
| MD5 |
d9af5bf8e99f484b30ef5615de610d10
|
|
| BLAKE2b-256 |
7bff319fc61cd810ede7b8b63bc911c444dbe5f5b2e3d5c7edbfad934d01a1cf
|
Provenance
The following attestation bundles were made for universe_reads-0.2.8.tar.gz:
Publisher:
publish.yml on yangaxnkohla/universe-reads
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
universe_reads-0.2.8.tar.gz -
Subject digest:
dfc5e39e514ccca121c3c5a698b49dc9696b2bb6bce80269bc6764cf62a2ee3c - Sigstore transparency entry: 1103865738
- Sigstore integration time:
-
Permalink:
yangaxnkohla/universe-reads@a8fbdb40e81500986fa9ca7f49443309aabe58a4 -
Branch / Tag:
refs/heads/main - Owner: https://github.com/yangaxnkohla
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@a8fbdb40e81500986fa9ca7f49443309aabe58a4 -
Trigger Event:
workflow_dispatch
-
Statement type:
File details
Details for the file universe_reads-0.2.8-py3-none-any.whl.
File metadata
- Download URL: universe_reads-0.2.8-py3-none-any.whl
- Upload date:
- Size: 28.0 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? Yes
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
b6368c9d2070cb7610bcf1ce8d40625047a4ede453edcf0e2996b136a14620b2
|
|
| MD5 |
513cbd22ac61cddfec274a6c66b81284
|
|
| BLAKE2b-256 |
1bb1f3a90d31e714b779d0889d94e3ae2e3986a4c7e63b3c68b3c6df5067336a
|
Provenance
The following attestation bundles were made for universe_reads-0.2.8-py3-none-any.whl:
Publisher:
publish.yml on yangaxnkohla/universe-reads
-
Statement:
-
Statement type:
https://in-toto.io/Statement/v1 -
Predicate type:
https://docs.pypi.org/attestations/publish/v1 -
Subject name:
universe_reads-0.2.8-py3-none-any.whl -
Subject digest:
b6368c9d2070cb7610bcf1ce8d40625047a4ede453edcf0e2996b136a14620b2 - Sigstore transparency entry: 1103865837
- Sigstore integration time:
-
Permalink:
yangaxnkohla/universe-reads@a8fbdb40e81500986fa9ca7f49443309aabe58a4 -
Branch / Tag:
refs/heads/main - Owner: https://github.com/yangaxnkohla
-
Access:
private
-
Token Issuer:
https://token.actions.githubusercontent.com -
Runner Environment:
github-hosted -
Publication workflow:
publish.yml@a8fbdb40e81500986fa9ca7f49443309aabe58a4 -
Trigger Event:
workflow_dispatch
-
Statement type: