Metadata-Version: 2.4
Name: parafetch
Version: 0.1.0
Classifier: Programming Language :: Rust
Classifier: Programming Language :: Python :: 3
Classifier: Programming Language :: Python :: 3.12
Classifier: Programming Language :: Python :: 3.13
Classifier: Programming Language :: Python :: 3.14
Classifier: Programming Language :: Python :: Implementation :: CPython
Requires-Dist: polars>=1.0
Requires-Dist: pyarrow>=14
Requires-Dist: pandas>=2 ; extra == 'pandas'
Requires-Dist: pytest>=8 ; extra == 'test'
Requires-Dist: pandas>=2 ; extra == 'test'
Provides-Extra: pandas
Provides-Extra: test
Summary: Rust-powered parallel data access for Python: feather-fast SQLite reads, GIL-aware parallel map, and batched REST / JSON-RPC / XML-RPC fetching.
License: MIT
Requires-Python: >=3.12
Description-Content-Type: text/markdown; charset=UTF-8; variant=GFM

# parafetch

**Rust-powered parallel data access for Python 3.12+.** No Rust needed to install it.

parafetch solves three common bottlenecks in Python data pipelines:

| You have... | parafetch gives you |
|---|---|
| A big SQLite table that takes ages to load | **`read_sqlite`**: a multi-threaded SQLite reader that returns a **polars DataFrame**, about **16x faster** than `pandas.read_sql` (pandas and pyarrow output too) |
| A function that handles **one** item at a time (e.g. `get_price(isin)`) | **`parallel_map` / `parallel_dict`**: call it for thousands of items concurrently, with retries and rate limiting |
| An API that accepts **at most N IDs per call** (e.g. 50 trade IDs) | **`HttpClient`, `JsonRpcClient`, `XmlRpcClient`, `batched_map`**: send any number of IDs. They are split into batches, fetched concurrently and merged back. |

The heavy lifting (SQLite decoding, HTTP, retries, waiting, JSON/XML parsing) runs in compiled
Rust with Python's GIL released. Your Python code only sees ordinary lists, dicts and DataFrames.

---

## Contents

- [Installation](#installation)
- [Quick start](#quick-start)
- [1. Reading SQLite fast — `read_sqlite`](#1-reading-sqlite-fast)
- [2. Parallelising one-at-a-time functions — `parallel_map`](#2-parallelising-one-at-a-time-functions)
- [3. Fetching from size-limited APIs — REST, JSON-RPC, XML-RPC](#3-fetching-from-size-limited-apis)
- [Error handling](#error-handling)
- [How parafetch deals with the GIL](#how-parafetch-deals-with-the-gil)
- [API reference](#api-reference)
- [FAQ & troubleshooting](#faq--troubleshooting)

---

## Installation

```bash
pip install parafetch
```

Requirements:

- **64-bit CPython 3.12 or newer**: 3.12, 3.13, 3.14, and so on. One wheel covers all of them.
- Windows (x64), Linux (x86_64 / aarch64) or macOS (Intel / Apple silicon).
- Installs `polars` and `pyarrow` automatically, both prebuilt wheels. `pandas` is optional: `pip install "parafetch[pandas]"`.

Wheels are precompiled. You do **not** need Rust, a C compiler, Visual Studio or the VC++
redistributable. SQLite is built into the package, and HTTPS uses the operating system's TLS and
certificate store, so corporate root certificates work automatically on Windows and macOS.

### Corporate / restricted machines

Check your interpreter first:

```powershell
python -c "import sys, platform; print(sys.version, platform.architecture()[0])"
# needs: 3.12 or newer, and 64bit
```

To make sure pip never tries to compile anything, install binaries only:

```powershell
pip install --only-binary=:all: parafetch
```

**Offline machines:** download the wheels on a connected machine, copy the folder over, and install from it:

```powershell
pip download parafetch --only-binary=:all: --platform win_amd64 --python-version 3.12 -d wheelhouse
pip install --no-index --find-links wheelhouse parafetch
```

If pip reports *"no matching distribution found"*, the interpreter is not 64-bit CPython 3.12 or newer.

---

## Quick start

```python
import parafetch as pf

# 1. Load a large SQLite table
df = pf.read_sqlite("market.db", "prices")              # polars DataFrame
pdf = pf.read_sqlite_pandas("market.db", "prices")      # pandas DataFrame

# 2. Call a one-ISIN-at-a-time function for many ISINs at once
from mylib import get_price
prices = pf.parallel_dict(get_price, isins, workers=32)        # {isin: price}

# 3. Fetch 10,000 trades from an API that only accepts 50 IDs per request
api = pf.HttpClient("https://trades.example.com/api", bearer_token=TOKEN)
trades = api.fetch_batched(trade_ids, batch_size=50, path="/trades", id_param="ids")
```

---

## 1. Reading SQLite fast

```python
pf.read_sqlite(path, table=None, *, query=None, columns=None, where=None, params=None,
               partition_on=None, partitions=None, schema_overrides=None,
               batch_size=131072, immutable=False) -> polars.DataFrame
```

`read_sqlite_pandas(...)` (pandas DataFrame) and `read_sqlite_arrow(...)` (pyarrow Table) take
exactly the same arguments.

### Read a whole table

```python
import parafetch as pf

df    = pf.read_sqlite("market.db", "prices")          # polars.DataFrame
pdf   = pf.read_sqlite_pandas("market.db", "prices")   # pandas.DataFrame
table = pf.read_sqlite_arrow("market.db", "prices")    # pyarrow.Table
```

### Pick columns and filter rows

```python
df = pf.read_sqlite(
    "market.db", "prices",
    columns=["isin", "px", "asof"],
    where="asof >= ? AND ccy = ?",
    params=["2024-01-01", "EUR"],
)
```

`where` is a SQL expression added to the query. Use `params` placeholders for values rather than
formatting them into the string.

### Run any SQL query

```python
sql = """
    SELECT t.id, t.isin, t.qty, p.px
    FROM trades t JOIN prices p ON p.isin = t.isin
    WHERE t.book = :book
"""
df = pf.read_sqlite("market.db", query=sql, params={"book": "EQ-LON"}, partition_on="id")
```

- `params` takes a list for `?` placeholders or a dict for `:name` placeholders.
- **Parallelism for queries:** give `partition_on`, an integer column that the query returns (ideally
  indexed, such as a primary key). Without it, a query is read on a single thread, which is still
  much faster than `pandas.read_sql`.

### How the speed-up works

- **Parallel:** for a table read, the rows are split into ranges of SQLite's `rowid`. Each range is
  read by its own thread through its own read-only connection. `partitions` sets how many threads
  (default: number of CPUs, max 16). Partitions are concatenated in order, so rows come back in
  `rowid` order.
- **No per-row Python objects:** values are written straight into Arrow column buffers, and the
  result is handed to polars (or pyarrow) without copying, just like reading a Feather file.
  `read_sqlite` doesn't rechunk, so the DataFrame keeps one chunk per batch. Call `df.rechunk()` if
  you need contiguous columns.
- **Tuned connections:** each connection is read-only with memory-mapped I/O and a large page cache.
  `immutable=True` additionally skips file locking. Use it only for files that nothing writes to
  during the read, such as a nightly snapshot.

Benchmark (2,000,000 rows × 6 columns, 8-CPU Windows laptop, `benchmarks/bench_sqlite.py`):

| method | time | |
|---|---|---|
| `sqlite3` + `pandas.read_sql` | 10.7 s | 1x |
| `pf.read_sqlite(...)`: polars, parallel | 0.65 s | **16x** |
| `pf.read_sqlite_arrow(...)`: pyarrow, parallel | 0.70 s | 15x |
| `pf.read_sqlite_pandas(...)`: pandas, parallel | 0.81 s | 13x |

### Column types

SQLite columns don't have strict types, so parafetch decides each column's type like this:

| declared SQLite type contains | polars | pyarrow |
|---|---|---|
| `INT` | `Int64` | `int64` |
| `CHAR`, `CLOB`, `TEXT` | `String` | `string` |
| `REAL`, `FLOA`, `DOUB` | `Float64` | `float64` |
| `BLOB` | `Binary` | `binary` |
| `BOOL` | `Boolean` | `bool` |
| anything else (`NUMERIC`, `DATE`, no type, expressions) | inferred from the values |

If a column holds mixed values, it is **widened, never truncated**: int and float become `float64`;
numbers and text become `string`. An all-NULL column keeps its declared type.

To force a type, use `schema_overrides`. Values that can't be converted raise an error instead of
silently becoming NULL:

```python
pf.read_sqlite("market.db", "prices", schema_overrides={"isin": "str", "qty": "float64", "active": "bool"})
# allowed types: "int64", "float64", "str", "bytes", "bool"
```

Dates are stored as text in SQLite and come back as strings. Convert them after reading, e.g.
`df.with_columns(pl.col("asof").str.to_date())` in polars, or `pd.to_datetime(pdf["asof"])` in pandas.

### Limitations

- Tables declared `WITHOUT ROWID` are read on a single thread unless you pass `partition_on`.
- `:memory:` databases cannot be read, because each thread opens its own connection to the file.

---

## 2. Parallelising one-at-a-time functions

```python
pf.parallel_map(func, items, *, workers=8, mode="auto", retries=0, backoff=0.5, max_backoff=30.0,
                rate_limit=None, retry_on=None, star=False, return_exceptions=False,
                fail_fast=False, chunksize=None) -> list
pf.parallel_dict(func, keys, **same_options) -> dict
```

### Basic use

```python
from pricing import get_price            # get_price("US0378331005") -> 189.3

isins = ["US0378331005", "US5949181045", ...]   # thousands

prices = pf.parallel_map(get_price, isins, workers=32)      # list, same order as isins
prices = pf.parallel_dict(get_price, isins, workers=32)     # {isin: price}
```

### Retries, back-off and rate limits

```python
prices = pf.parallel_dict(
    get_price, isins,
    workers=32,
    retries=3,                                  # up to 3 retries per item
    backoff=0.5, max_backoff=10,                # 0.5s, 1s, 2s ... with ±25% jitter
    retry_on=(ConnectionError, TimeoutError),   # only retry these (default: any Exception)
    rate_limit=200,                             # at most 200 calls/second in total
)
```

### Functions with several arguments

```python
def get_fx(ccy_from, ccy_to, date): ...

rates = pf.parallel_map(get_fx, [("EUR", "USD", d) for d in dates], star=True)  # calls get_fx(*args)
```

### Async functions

`async def` functions are detected automatically and run on an event loop, with at most `workers`
running at once. This also works inside Jupyter or a running asyncio application.

```python
async def get_price_async(isin): ...

prices = pf.parallel_map(get_price_async, isins, workers=100)
```

### CPU-heavy pure-Python functions

Threads can't speed up pure-Python number crunching because of the GIL. Use separate processes, or on
Python 3.14+, sub-interpreters:

```python
# my_models.py
def value_option(params): ...   # pure Python, CPU-bound

# main.py
import parafetch as pf
from my_models import value_option

if __name__ == "__main__":       # required on Windows for process mode
    values = pf.parallel_map(value_option, scenarios, mode="process", workers=8)
    # Python 3.14+: mode="interpreter" (lighter than processes)
```

For `process` and `interpreter` modes, the function must be defined at module level (not a lambda or
nested function), and its arguments and results must be picklable.

### Choosing `workers`

- I/O-bound functions (DB, HTTP, services): 16–64 is typical. The limit is what the server tolerates.
- CPU-bound with `mode="process"`: the number of CPU cores.

---

## 3. Fetching from size-limited APIs

When an endpoint takes at most N IDs per call, the clients below:

1. split your IDs into batches of `batch_size`,
2. send all batches **concurrently** (up to `concurrency` requests at a time) from Rust, with
   connection pooling and keep-alive,
3. retry failed requests (HTTP 408/425/429/500/502/503/504 and network errors) with exponential
   back-off, honoring `Retry-After`,
4. parse the responses (JSON or XML-RPC) outside the GIL,
5. return one flat list: per-batch lists are concatenated and per-batch dicts are merged, in input order.

### REST — `HttpClient`

```python
api = pf.HttpClient(
    "https://trades.example.com/api/v2",
    bearer_token=TOKEN,       # or basic_auth=("user", "pwd") or headers={"X-API-Key": KEY}
    concurrency=16,           # parallel requests
    rate_limit=50,            # requests/second (optional)
    retries=3,
    timeout=60,
)
```

**IDs in a JSON body.** Every `"$ids"` string in the body is replaced by the batch's list:

```python
trades = api.fetch_batched(
    trade_ids, batch_size=50,
    method="POST", path="/trades/search",
    json={"tradeIds": "$ids", "includeLegs": True},
    result_key="data.trades",      # where the list sits in each response
)
```

**IDs in the query string:**

```python
api.fetch_batched(trade_ids, path="/trades", id_param="ids")                     # ?ids=1,2,3
api.fetch_batched(trade_ids, path="/trades", id_param="id", id_style="repeat")   # ?id=1&id=2&id=3
```

**IDs in the URL path.** `{ids}` is replaced by the comma-joined batch:

```python
api.fetch_batched(trade_ids, path="/trades/{ids}")                    # /trades/1,2,3
api.fetch_batched(isins, batch_size=1, path="/prices/{ids}")          # one ID per call
```

**Extracting the payload.** `result_key` is a dotted path (`"data.trades"`, `"items.0.rows"`) or a
function, e.g. `result_key=lambda resp: resp["data"]["trades"]`.

**Building the body yourself:**

```python
api.fetch_batched(trade_ids, method="POST", path="/query",
                  json=lambda ids: {"filter": {"id": {"$in": ids}}, "limit": len(ids)})
```

**Other requests:**

```python
api.request("GET", "/status")                                     # single request -> parsed JSON
api.get_many(f"/prices/{isin}" for isin in isins)                 # many GETs, concurrently
api.request_many([                                                # arbitrary mix
    {"method": "GET",  "path": "/prices", "params": {"isin": "US0378331005"}},
    {"method": "POST", "path": "/fx", "json": {"pair": "EURUSD"}},
])
```

Each `request_many` entry takes `method`, `path` (or a full URL), `params`, `json`, `data`,
`headers` and `parse` (`"json"`, `"text"`, `"bytes"` or `"auto"`).

### JSON-RPC — `JsonRpcClient`

```python
rpc = pf.JsonRpcClient("https://rpc.example.com/jsonrpc", basic_auth=("svc", PWD), concurrency=16)

rpc.call("getPrice", {"isin": "US0378331005"})                         # single call

# Method that accepts up to 50 IDs:
trades = rpc.call_batched("getTrades", trade_ids, batch_size=50)                   # params=[[ids]]
trades = rpc.call_batched("getTrades", trade_ids, params={"tradeIds": "$ids"})     # named params

# One call per item:
prices = rpc.call_many("getPrice", [{"isin": i} for i in isins])

# Fewer HTTP round-trips: pack 100 calls into each request (JSON-RPC batch arrays)
prices = rpc.call_many("getPrice", [{"isin": i} for i in isins], calls_per_request=100)
```

JSON-RPC 2.0 is the default. Use `JsonRpcClient(url, version="1.0")` for 1.0 servers. Server-side
errors raise `pf.JsonRpcError`, with `.code`, `.data` and `.error` attributes.

### XML-RPC — `XmlRpcClient`

```python
xr = pf.XmlRpcClient("https://legacy.example.com/RPC2", basic_auth=("svc", PWD))

xr.call("pricing.get", "US0378331005")                           # like ServerProxy().pricing.get(...)
trades = xr.call_batched("trades.get", trade_ids, batch_size=50) # trades.get([id1, ..., id50])
trades = xr.call_batched("trades.get", trade_ids, params=("$ids", "2024-06-30"))  # extra args
prices = xr.call_many("pricing.get", isins)                      # one call per ISIN
prices = xr.call_many("pricing.get", isins, multicall_size=100)  # system.multicall: 100 per request
```

- Supports every XML-RPC type: int/i4/i8, boolean, string, double, dateTime.iso8601 (returned as
  `datetime.datetime`), base64 (returned as `bytes`), array, struct and nil.
- Server faults raise the standard `xmlrpc.client.Fault`, as with the stdlib `ServerProxy`.

### Connection options (all three clients)

| option | default | meaning |
|---|---|---|
| `headers` | `None` | extra headers for every request |
| `bearer_token` | `None` | sends `Authorization: Bearer <token>` |
| `basic_auth` | `None` | `("user", "password")` |
| `concurrency` | `16` | maximum requests in flight |
| `rate_limit` | `None` | maximum requests per second (shared by all calls on the client) |
| `retries` | `3` | retries per request |
| `backoff`, `max_backoff` | `0.5`, `30` | exponential back-off in seconds, with jitter |
| `retry_statuses` | `408, 425, 429, 500, 502, 503, 504` | HTTP statuses that are retried |
| `timeout`, `connect_timeout` | `60`, `10` | seconds (`None` disables the timeout) |
| `verify` | `True` | set `False` to skip TLS certificate checks (not recommended) |
| `ca_cert` | `None` | path to a PEM bundle with extra trusted CAs |
| `proxy` | `None` | e.g. `"http://proxy.corp:8080"`. `HTTP(S)_PROXY` env vars are used automatically. |
| `user_agent` | `parafetch/<version>` | |

Create a client once and reuse it. It keeps connections open between calls.

Requests are retried on network errors, which can re-send a POST. That's fine for read/fetch APIs.
For endpoints that change data, set `retries=0`.

### Using your own Python client or SDK

If you already have a client function that takes a list of IDs, use `batched_map`:

```python
trades = pf.batched_map(sdk.get_trades, trade_ids, batch_size=50, workers=16, retries=2)
```

This works with any callable: database stored procedures, vendor SDKs, `requests`-based helpers.
Results are flattened in the same way.

---

## Error handling

parafetch never throws away work that succeeded.

**`parallel_map` / `parallel_dict` / `call_many`** raise `pf.ParallelError` if any item fails:

```python
try:
    prices = pf.parallel_map(get_price, isins, retries=2)
except pf.ParallelError as e:
    print(e)                      # "3 of 5000 items failed; first failure (item 17): TimeoutError: ..."
    e.results                     # full-length list, None where an item failed
    e.errors                      # {index: exception}
    e.failed_items                # the inputs that failed, ready to retry
    raise e.__cause__             # the first underlying exception
```

**`fetch_batched` / `call_batched` / `batched_map`** raise `pf.BatchError`:

```python
try:
    trades = api.fetch_batched(trade_ids, path="/trades", id_param="ids")
except pf.BatchError as e:
    trades = e.results            # data from the batches that succeeded
    retry  = e.failed_ids         # IDs from failed batches
    e.errors                      # [(ids, message, http_status), ...]
```

**Other options:**

- `return_exceptions=True`: no exception is raised. Each failed position in the result list holds its
  exception object instead (per item, or per batch for batched calls).
- `fail_fast=True` (`parallel_map`): stop starting new items after the first failure. Skipped items
  are reported as `concurrent.futures.CancelledError`.
- **Ctrl+C** works: running work is cancelled and `KeyboardInterrupt` is raised.

---

## How parafetch deals with the GIL

The GIL lets only one thread run Python code at a time. parafetch keeps Python code to a minimum on
the hot path:

| feature | what runs without the GIL |
|---|---|
| `read_sqlite` | everything: SQLite I/O, decoding and Arrow building, across N threads. Python only receives the finished Arrow table. |
| HTTP / JSON-RPC / XML-RPC clients | everything: connections, TLS, waiting, retries, rate limiting, JSON/XML parsing. Python builds the request bodies before the call and receives the parsed results after. |
| `parallel_map(mode="thread")` | scheduling, rate-limit waits and retry sleeps. Your function holds the GIL only while it executes Python code. When it does I/O (DB, network), CPython releases the GIL, so other workers run meanwhile. |
| `mode="process"` / `"interpreter"` | each worker has its own GIL, so pure-Python CPU work runs truly in parallel |

Rule of thumb: **waiting on I/O → `thread`** (the default), **async code → `async`**,
**CPU-heavy Python → `process`** (or `interpreter` on 3.14+).

---

## API reference

### SQLite

| function | returns |
|---|---|
| `read_sqlite(path, table=None, *, query=None, columns=None, where=None, params=None, partition_on=None, partitions=None, schema_overrides=None, batch_size=131072, immutable=False)` | `polars.DataFrame` |
| `read_sqlite_pandas(path, table=None, **kwargs)` | `pandas.DataFrame` (needs pandas) |
| `read_sqlite_arrow(path, table=None, **kwargs)` | `pyarrow.Table` |

Pass exactly one of `table` or `query`. `columns` and `where` apply to `table` reads only.

### Parallel execution

| function | description |
|---|---|
| `parallel_map(func, items, **opts)` | `[func(x) for x in items]`, run concurrently, in input order |
| `parallel_dict(func, keys, **opts)` | `{k: func(k) for k in keys}`, run concurrently |
| `batched_map(func, items, *, batch_size=50, flatten=True, **opts)` | `func(batch)` per batch, results flattened |
| `chunked(items, size)` | split into lists of at most `size` items |

`opts`: `workers`, `mode` (`"auto"`, `"thread"`, `"async"`, `"process"`, `"interpreter"`), `retries`, `backoff`,
`max_backoff`, `rate_limit`, `retry_on`, `star`, `return_exceptions`, `fail_fast`, `chunksize`.

### HTTP clients

| method | description |
|---|---|
| `HttpClient(base_url="", **conn)` | REST client |
| `.fetch_batched(ids, *, batch_size=50, method="GET", path="", id_param=None, id_style="csv", params=None, json=None, headers=None, result_key=None, flatten=True, parse="json", return_exceptions=False)` | batched ID fetch |
| `.request(method="GET", path="", *, params=None, json=None, data=None, headers=None, parse="json")` | one request |
| `.request_many(requests, *, parse="json", return_exceptions=False, raw=False)` | many requests concurrently |
| `.get_many(paths, *, params=None)` | many GETs concurrently |
| `JsonRpcClient(url, *, version="2.0", **conn)` | JSON-RPC client |
| `.call(method, params=None)` | one call |
| `.call_many(method, params_list, *, calls_per_request=1, return_exceptions=False)` | one call per params entry |
| `.call_batched(method, ids, *, batch_size=50, params=("$ids",), result_key=None, flatten=True, return_exceptions=False)` | batched ID call |
| `XmlRpcClient(url, *, allow_none=True, **conn)` | XML-RPC client |
| `.call(method, *params)` | one call |
| `.call_many(method, params_list, *, multicall_size=1, return_exceptions=False)` | one call per entry (tuple = positional args) |
| `.call_batched(method, ids, *, batch_size=50, params=("$ids",), result_key=None, flatten=True, return_exceptions=False)` | batched ID call |

`**conn` are the [connection options](#connection-options-all-three-clients).
`request_many(..., raw=True)` returns per-request dicts: `ok`, `status`, `data`, `error`, `attempts`, `elapsed`.

### Exceptions

| exception | raised by | useful attributes |
|---|---|---|
| `ParallelError` | `parallel_map`, `parallel_dict`, `call_many` | `results`, `errors`, `failed_items`, `__cause__` |
| `BatchError` | `fetch_batched`, `call_batched`, `batched_map` | `results`, `failed_ids`, `errors` |
| `JsonRpcError` | `JsonRpcClient.call` | `code`, `data`, `error` |
| `xmlrpc.client.Fault` | `XmlRpcClient.call` | `faultCode`, `faultString` |

---

## FAQ & troubleshooting

**pip tries to build from source, or asks for Rust.**
Your interpreter has no matching wheel: it's 32-bit, older than 3.12, or not CPython. Check it with
the command under [Corporate / restricted machines](#corporate--restricted-machines), and install
with `--only-binary=:all:`.

**`parallel_map` isn't faster than a loop.**
Your function is probably CPU-bound pure Python, so the GIL serialises it. Use `mode="process"`. If
the function is I/O-bound, increase `workers`.

**`mode="process"` fails with a pickling error, or hangs on Windows.**
Define the function at module level in an importable file, and call parafetch under
`if __name__ == "__main__":`.

**`SSL: certificate verify failed` against an internal API.**
On Windows and macOS the OS certificate store is used, so company CAs installed there are trusted.
Otherwise, pass `ca_cert="corp-ca.pem"`.

**The server returns 429 Too Many Requests.**
Lower `concurrency` or set `rate_limit`. 429 responses are retried automatically and honor
`Retry-After`.

**`read_sqlite` says a column "mixes Blob and Int values".**
That column contains both binary and numeric data. Use `schema_overrides={"col": "bytes"}` (or `"str"`).

**Can I read the database while another process writes to it?**
Yes, with the default settings. Readers use normal SQLite locking (WAL mode is ideal). Don't use
`immutable=True` in that case.

---

License: MIT

