Metadata-Version: 2.4
Name: quickhouse
Version: 0.3.5
Classifier: Development Status :: 3 - Alpha
Classifier: Intended Audience :: Developers
Classifier: License :: OSI Approved :: MIT License
Classifier: Operating System :: OS Independent
Classifier: Programming Language :: Rust
Classifier: Programming Language :: Python :: 3
Classifier: Topic :: Database
Classifier: Topic :: Software Development :: Libraries
Requires-Dist: tqdm>=4.60 ; extra == 'progress'
Requires-Dist: pytest>=7 ; extra == 'test'
Requires-Dist: psycopg[binary]>=3.1 ; extra == 'test'
Requires-Dist: pymysql>=1.1 ; extra == 'test'
Requires-Dist: clickhouse-connect>=0.7 ; extra == 'test'
Requires-Dist: tqdm>=4.60 ; extra == 'test'
Requires-Dist: boto3>=1.34 ; extra == 'test'
Requires-Dist: pyarrow>=14 ; extra == 'test'
Provides-Extra: progress
Provides-Extra: test
License-File: LICENSE
Summary: Fast PostgreSQL/MySQL/BigQuery -> ClickHouse ETL with a Rust engine (parallel, bounded-memory, Arrow-based).
Keywords: etl,postgresql,mysql,bigquery,clickhouse,arrow,rust,data-engineering
Author-email: M Mirza Fahmi <mr.sm1l3yz@gmail.com>
License: MIT
Requires-Python: >=3.9
Description-Content-Type: text/markdown; charset=UTF-8; variant=GFM
Project-URL: Homepage, https://github.com/mmirzafahmi/quickhouse
Project-URL: Issues, https://github.com/mmirzafahmi/quickhouse/issues
Project-URL: Repository, https://github.com/mmirzafahmi/quickhouse

# quickhouse

**Move tables from PostgreSQL, MySQL, or BigQuery into ClickHouse or BigQuery — fast, in one function call.**

quickhouse is a small, typed Python API on top of a native Rust engine. You
hand it a source, a destination, and a table name; it figures out the schema,
creates the destination table, streams the rows across in parallel, and keeps
memory flat the whole way. The heavy lifting never touches Python objects —
each database's native wire protocol flows straight into Apache Arrow and out
the other side.

```python
import quickhouse

src = quickhouse.Postgres("postgresql://user:pw@localhost:5432/shop")
dst = quickhouse.ClickHouse("http://localhost:8123", database="analytics")

result = quickhouse.sync(src, dst, dest_table="orders",
                         source_table="orders", key=["id"])
print(result)   # rows_read, rows_written, bytes_written, duration_secs, new_watermark
```

## Why quickhouse

- **It's fast.** Rows are decoded straight off the wire into Arrow, in Rust —
  no per-row Python, no intermediate DataFrame. Tables are split into ranges
  and read in parallel, and decoding overlaps uploading. On a laptop-class box
  a 1M-row, 20-column full refresh runs at **hundreds of thousands of rows per
  second** while peak memory stays flat (under ~180 MB) no matter how much you
  parallelize. Reproduce it with `python benchmarks/bench_transfer.py`.

- **It's one function call.** `sync()` replaces the cursor loop, manual
  batching, retry logic, and `CREATE TABLE` you'd otherwise write by hand.
  Defaults handle table creation, type mapping, parallelism, and batching, and
  a typed stub gives you autocomplete on every argument.

- **It's safe with real, messy data.** Full refreshes swap in atomically, so a
  crash never leaves a half-written table. Incremental syncs are idempotent —
  safe to re-run or retry. Transient network blips retry automatically. And
  legacy quirks like MySQL zero-dates or out-of-range timestamps are coerced to
  `NULL` with a warning instead of aborting the run.

- **It's gentle on a small production database.** Set `read_max_rows_per_sec`
  and quickhouse paces the read to that aggregate rate across all partitions;
  because `COPY`/streaming results only produce as fast as the client consumes,
  the source scan itself backs off — you're throttling the database's work, not
  just your own. Combine it with `parallelism=1` (one connection), incremental
  mode (only new rows), and a `statement_timeout_secs`, point it at a read
  replica, and a bulk export stops competing with your app. The Postgres
  connection also shows up as `application_name = 'quickhouse'` in
  `pg_stat_activity`, so a DBA can see and kill it.

- **There's nothing to stand up.** `pip install quickhouse` and you're done —
  no JVM, no Spark cluster, no separate service. It's an ordinary Python
  dependency that runs wherever your jobs already run: cron, Airflow, Dagster,
  a Lambda, or a plain script.

## Install

```bash
pip install quickhouse
pip install "quickhouse[progress]"   # adds a ready-made tqdm progress bar
```

Prebuilt wheels ship for Python 3.9+ on Linux, macOS (Intel + Apple Silicon),
and Windows (x86_64) — no Rust toolchain needed. Building from source is only
for development; see [CONTRIBUTING.md](CONTRIBUTING.md).

## Using it

A fuller call, with the options you'll reach for most:

```python
import quickhouse as qh

src = qh.Postgres("postgresql://user:pw@localhost:5432/shop")
dst = qh.ClickHouse("http://localhost:8123", database="analytics")

qh.sync(
    src, dst,
    dest_table="orders",
    source_table="orders",        # or source_query="SELECT ..."
    mode="incremental",           # or "full"
    watermark="updated_at",       # required for incremental
    key=["id"],                   # dedup key / ORDER BY
    parallelism=8,
    exclude=["internal_notes"],
    rename={"amount": "amt"},
    on_progress=lambda p: print(f"{p.rows_written:,} rows @ {p.rows_per_sec:,.0f}/s"),
)
```

### Sources and destinations

Pick a source and a destination by constructing the matching object —
everything else about `sync()` stays the same:

```python
# sources
qh.Postgres("postgresql://user:pw@host:5432/db")
qh.MySQL("mysql://user:pw@host:3306/db", require_tls=True)
qh.BigQuery("my-gcp-project")                       # source_table="dataset.table"

# destinations
qh.ClickHouse("http://host:8123", database="analytics")
qh.BigQuery("my-gcp-project", dataset_id="analytics")
```

BigQuery authenticates with a service-account key (`credentials_file=...`) or
Application Default Credentials. As a **destination** it also takes
`write_method`: the default `"insert_all"` (simple, proven) or the opt-in
`"storage_write"` (the gRPC Storage Write API — free and higher-throughput).

A ClickHouse destination can also archive every synced batch to S3 as a data
lake — a secondary, best-effort-free backup independent of ClickHouse's own
retention:

```python
qh.ClickHouse(
    "http://host:8123", database="analytics",
    archive=qh.S3Archive(bucket="my-data-lake", prefix="quickhouse"),
)
```

This streams Parquet — one file per parallel partition, never fully buffered
in memory — to `s3://{bucket}/{prefix}/{dest_table}/dt=<date>/run=<id>/
part-<partition>.parquet`, a Hive-style layout directly queryable by Athena,
Spark, or DuckDB. Credentials fall back to the standard AWS chain (env vars,
IAM role) unless overridden; pass `endpoint=` for an S3-compatible service
like MinIO. A persistent upload failure fails the whole `sync()` call, same as
a ClickHouse insert failure. Storage/request costs are billed by AWS as usual
(free on a self-hosted MinIO).

The DDL knobs (`engine`, `partition_by`, `order_by`, `primary_key`, `key`) are
interpreted per destination — for ClickHouse they shape the `MergeTree`
DDL; for BigQuery they map to partitioning and clustering. quickhouse creates
the table for you (`create_if_missing=True` by default) with a sensible schema
derived from the source.

### Full vs. incremental

**Full** reloads the whole table into a staging table, then swaps it into place
atomically — a crash mid-run never leaves the destination partial. For a
BigQuery destination that swap runs as a query (a billed scan of the staged
data), not a free copy job — BigQuery's copy jobs can silently skip rows still
sitting in a table's streaming buffer, so a real query is what keeps this
correct rather than just fast. One accepted tradeoff on the ClickHouse path:
an insert retried after a lost acknowledgment (not after a crash — the
transfer is still running) can duplicate one batch's rows in the staging
table, since `mode="full"` has no engine-level dedup like `ReplacingMergeTree`
— rare, and harmless for `key`-based incremental syncs, but worth knowing if
you see an unexpected small over-count on a full-refresh right after a
transient network blip.

**Incremental** tracks a high-water mark (the `watermark` column) in a small
state table in the destination and copies only newer rows. Updated rows are
deduplicated on `key` — via ClickHouse's `ReplacingMergeTree`, or a `MERGE`
upsert on BigQuery (where `key` is therefore required). Re-running with no new
data does nothing.

Both modes stage through a per-run-unique table (`{dest}_quickhouse_tmp_<id>`)
that's dropped when the run finishes, including on failure. The unique name is
what makes rapid re-runs and whole-call retries safe on BigQuery, whose
streaming ingestion rejects writes into a table recently recreated under the
same name.

For daily syncs that need to catch late-arriving or edited rows, set
`lookback_seconds` to re-scan a trailing window (e.g. `3 * 86400` for the last
three days) — the dedup above keeps that overlap from creating duplicates.

### Watching progress and diagnosing failures

`on_progress` is a plain callback you can point at anything; `qh.progress_bar()`
wraps [tqdm](https://github.com/tqdm/tqdm) for a ready-made bar. Every `sync()`
also logs each step to stderr (`RUST_LOG=quickhouse_core=debug` for the actual
SQL).

When something goes wrong, `sync()` raises a `RuntimeError` written to be
actionable on its own: it names the table involved, and for a bad config or an
unmappable column it says exactly what's wrong and how to fix it (e.g.
`exclude=` the column or cast it in a `source_query`). Underlying database
errors are surfaced verbatim rather than wrapped in something generic.

### Full parameter list

| Parameter | Meaning |
| --- | --- |
| `source_table` / `source_query` | Read a whole table, or a custom `SELECT` (one required) |
| `dest_table` | Destination table name |
| `mode` | `"full"` or `"incremental"` |
| `watermark` | Monotonic column for incremental (e.g. `updated_at`); ignored in full mode |
| `lookback_seconds` | Re-scan a trailing window of the watermark to catch late/edited rows; `0` disables (default) |
| `key` | Dedup key (required for BigQuery incremental) |
| `create_if_missing` | Auto-create the destination table (default `True`) |
| `engine`, `order_by`, `partition_by`, `primary_key` | DDL knobs, interpreted per destination |
| `parallelism` | Concurrent read streams |
| `batch_rows` / `batch_bytes` | Per-batch size knobs (rows, or estimated bytes) |
| `max_memory_bytes` | Hard ceiling on total in-flight memory; decoding blocks when hit (default 512 MiB, `0` = unbounded) |
| `read_max_rows_per_sec` | Cap the aggregate source read rate to be gentle on a small DB; `None` = unlimited (default). Postgres/MySQL only |
| `type_overrides` | Force a destination column type, e.g. `{"qty": "Decimal(18, 3)"}` |
| `rename`, `include`, `exclude` | Column renames and allow/deny lists |
| `on_progress` | Progress callback |

## How types are mapped

quickhouse maps each source type to a sensible destination type automatically:
integers to integers, floats to floats, text/JSON/UUID to strings, dates and
timestamps across as-is, and booleans preserved. A few deliberate choices worth
knowing:

- **Arbitrary-precision decimals** (`numeric`/`DECIMAL`/`NUMERIC`) default to
  `Float64`, since precision can't be recovered from the type alone — pin an
  exact type with `type_overrides` (e.g. `"Decimal(18, 2)"`, `P <= 38`) and the
  value is decoded exactly (no `Float64` round-trip), not just declared with
  the right destination type. A value that doesn't fit the declared precision,
  or is NaN/Infinity (PostgreSQL `numeric` only), coerces to `NULL` with a
  warning, same as the out-of-range-date handling below. `P > 38`
  (`Decimal256`) isn't supported yet and is rejected as a config error up
  front, rather than silently falling back to `Float64`.
- **`TIME`** columns transfer as canonical text into a `String` column
  (ClickHouse has no time-of-day type).
- **MySQL `DATETIME`/`TIMESTAMP`** map to a UTC-aware timestamp (BigQuery
  `TIMESTAMP`, ClickHouse `DateTime64(6, 'UTC')`) — the wall-clock value is read
  as UTC, matching how a `TIMESTAMP` column expects it. To land a column as a
  naive BigQuery `DATETIME` instead, opt out per-column with
  `type_overrides={"col": "DATETIME"}` — that flips the actual encoding, not
  just the declared type. (PostgreSQL keeps the distinction natively:
  `timestamptz` → UTC-aware, `timestamp` → naive.)
- **Out-of-range dates** (and MySQL zero-dates like `0000-00-00`) coerce to
  `NULL` with a warning rather than failing the transfer.
- **Nullable** source columns stay nullable in the destination.

Arrays and composite (`RECORD`/`STRUCT`) types aren't supported yet.

## Limitations

- **mTLS** (client-certificate auth) isn't supported; server TLS is, including
  an extra CA file via `ca_cert_file=...` for providers like AWS RDS.
- **Array / composite types** aren't mapped yet.
- **BigQuery as a source** reads through a single connection — `parallelism`
  becomes a server-side hint rather than true client-side fan-out (a limitation
  of the underlying crate's read API).
- **No CLI yet**, and CDC / custom transforms are future work.

## Contributing

Bug reports, new source/type mappings, and PRs are welcome — see
[CONTRIBUTING.md](CONTRIBUTING.md) for build steps, tests, and layout.

## License

MIT

