Skip to content

hoglake sink: pyhoglake 1.3.7 writer path — in-memory prepared appends, read_snapshot kept, zero sink-side reads per flush - #148

Merged
jghoman merged 2 commits into
mainfrom
jakob/pyhoglake-1.3.7
Oct 1, 2026
Merged

jghoman merged 2 commits into
mainfrom
jakob/pyhoglake-1.3.7

Conversation

@jghoman

@jghoman jghoman commented Oct 1, 2026

Copy link
Copy Markdown
Collaborator

Why

hoglake 1.3.6/1.3.7 moved the writer path (hoglake #234, #237, #239, #255): in-memory prepared appends with concurrent single-request uploads, a cached table shape whose read_snapshot_id rides the commit as read_snapshot so the server answers "did DDL touch this table since my read" instead of the client re-reading it, typed re-prepare refusals (ddl_since_read_snapshot, table_recreated, a 410 for a basis under the expiry floor), and info(totals=False) to skip the live-totals scan. millpond still wrote each partition's parquet to a temp dir, made two table GETs per flush with totals (a count and two sums over ~15M file rows on prod-us), and stripped read_snapshot from the payload. The server's HOGLAKE_REFUSE_BLIND_PARTITIONED_APPENDS flag will refuse payloads without a basis once the fleet sends one.

What

  • pyhoglake[fast-upload]>=1.3.7. Objects ≤ 8 MiB go up as one PutObject through boto3 instead of pyarrow's three-request multipart; uploads fan out on a thread pool (capped at the group count, 64 on the events writer). uv lock; the extra lands in the default dependency set, so the Dockerfile is unchanged (verified with its exact uv sync command).
  • _prepare → prepare_append_tables. No temp files; each partition's Arrow table encodes to a buffer and uploads from it. Orphan accounting on exceptions (uploaded_files / uploaded_uris) unchanged. Variant destinations: millpond cannot create one and columns_to_arrow_schema already refuses one, so no fallback path.
  • read_snapshot stays on the payload. The pre-commit table GET and _check_destination_still_ours are gone. Column drift is refused at prepare by the encoded-footer compare and self-heals in one retry with zero orphans. DDL that lands between the basis and the commit comes back as the typed 409/410, handled by a new _destination_moved arm: discard the payload (orphans counted), drop the cached shape, raise retryable so main.py rebuilds. commit_prepared(payload, table=...) so pyhoglake invalidates its own cache on a re-prepare refusal; without it the rebuild could resend the refused basis.
  • Zero sink-side reads per steady-state flush. The sink keeps the last TableInfo and re-reads only on first use, reset_caches(), a destination-moved refusal, the alignment self-heal, or when the copy is older than retention/4. That TTL is the correctness guard the review found missing: pyhoglake refreshes its own writer cache at retention/2, and once it has, a same-arity partition-spec change (identity → bucket) is outside the conflict window; a sink still computing values under the old spec would have published mis-partitioned files that the server stamps with the live spec_id. The sink's reads re-seed pyhoglake's cache, so with the shorter TTL the sink can never be the staler of the two. The periodic re-read also runs _reconcile_specs, so the "config vs live spec is fatal" tripwire fires on a long-lived pod, not only at start. totals=False on every remaining info call. Drop+recreate stays closed by expected_table_uuid → table_recreated.
  • Liveness budget. A rebuild (re-encode + re-upload the fanout) is now a routine per-attempt cost during a producer rollout. _hoglake_worst_case_flush_s charges a 20 s per-attempt prepare allowance (basis in the comment, flagged as an over-estimate); with it, 8 × 45 s models at ~634 s against the 480 s budget, so HOGLAKE_MAX_RETRY_COUNT defaults to 6 (~429 s, 51 s margin). A values file pinning 8 explicitly is refused at startup by the existing budget check; no charts values file sets it today (charts/millpond/values.yaml leaves maxRetryCount empty).
  • millpond_hoglake_orphaned_files_total gains a reason label (prepare_failed, commit_refused, ddl_since_read_snapshot, table_recreated, read_snapshot_expired, already_published, superseded, shutdown), so rollout churn and a broken pod stop looking alike. Summing over reason is the old series. No dashboard queries the bare series.
  • PYHOGLAKE_UPLOAD_CONCURRENCY validated at config.load() like every other knob. HOGLAKE_S3_REGION should be set explicitly even under IRSA: boto3 resolves region differently from Arrow's SDK (documented).
  • Integration stack pin was never applied. tests/hoglake_stack/stack.py resolved DEFAULT_IMAGE for docker pull only; compose fell back to :latest, 13 days stale locally, which the first run exposed (a 1.3.7 refusal arriving as a bare commit_conflict). The pin now reaches compose, both session fixtures assert the booted image via docker inspect, and the compose default is ${HOGLAKE_SERVER_IMAGE:?…} so a bare docker compose up cannot resurrect :latest.

Behaviour changes to know before promoting

  • A concurrent add_column now refuses the in-flight flush instead of being invisible to it. Each racing DDL costs one rebuild and one flush's orphaned objects (~270 on the events writer); refusals per flush are bounded by distinct new columns, not by pod count (a loser's alter mints no change row). The reason label makes this visible.
  • Memory: a flush holds ~3 Arrow-sized copies (main.py's consolidated batch across the retry loop, the aligned cast, the materialised partition groups) plus the encoded buffers. Not new, now stated correctly; ~1 GiB is the prod-us values figure, the code default is 100 MiB.

Tests

  • Unit: 1,315 passed, 1 xfailed (1,289 before the review fixes). hoglake sink file: TestZeroSinkSideReadsPerFlush, TestTypedDestinationMovedRefusals, TestTheCachedShapeHasATtl (9 tests, injected clock; removing the TTL fails 3, raising it to retention/2 fails 4), the rebuild-under-the-new-spec test, totals=False asserted at every call site by the harness (all six sites mutation-verified), per-flush orphan (reason, count) assertions. The harness change exposed and fixed one vacuous pre-existing test (test_conflict_on_add_reresolves_and_proceeds never attempted the add).
  • Integration against the real 1.3.7 server image: 49 passed in 119 s, including TestOneRequestPerFlush (five steady-state flushes produce exactly five POST .../commit/prepared on the wire, every table GET carries totals=false), per-flush exactly-once in the concurrent-writers test, the mid-flight re-spec refusal, the warm-cache recreate costing exactly one flush of uploads, and the boto3 transport's orphan accounting with a thread-safe Nth-failure gate.
  • E2E (Kafka → main.py → hoglake): 6 passed. DuckLake integration: 65 passed. ruff check and ruff format --check clean on tracked files.
  • Reviewed by two independent passes (principal SWE on the sink, lead QE on the tests with 14 mutations) before the fixes above.

Deploy

  • Image via the usual release workflow; promote with HOGLAKE_S3_REGION set (it already is in charts/millpond).
  • Watch after the roll: millpond_hoglake_orphaned_files_total{reason=...} and hoglake_commits_total{result="ddl_since_read_snapshot"} on the server; hoglake_blind_partitioned_appends_total should drop to zero for millpond writers, after which the server flag can flip.

…s, read_snapshot kept, zero sink-side reads per flush

Pin pyhoglake[fast-upload]>=1.3.7. _prepare encodes each partition's
Arrow table to a buffer and uploads it through prepare_append_tables
(single-request PutObject for small objects, concurrent fan-out)
instead of writing temp files for prepare_append_files.

The payload keeps read_snapshot. The pre-commit table GET and
_check_destination_still_ours are gone: the server's typed refusals
(ddl_since_read_snapshot, table_recreated, the 410 for an expired
basis) are handled by discarding the prepared payload, dropping the
cached shape and raising retryable so main.py rebuilds. pyhoglake is
told to invalidate its own cache on the same refusal.

The sink keeps the last TableInfo and re-reads it only on first use,
reset, a destination-moved refusal, the alignment self-heal, or when
it is older than retention/4 — strictly shorter than pyhoglake's own
retention/2 refresh, so the sink can never be the staler of the two.
Without that TTL a same-arity partition-spec change landed outside
the conflict window after pyhoglake refreshed, and files computed
under the old spec would have been accepted and stamped with the
live spec_id. The periodic re-read also runs _reconcile_specs. Every
remaining info call sends totals=False.

The liveness budget charges a 20 s per-attempt prepare allowance, so
HOGLAKE_MAX_RETRY_COUNT defaults to 6. The orphan counter gains a
reason label. PYHOGLAKE_UPLOAD_CONCURRENCY is validated at startup.

The integration stack's server-image pin never reached compose and
the suites ran against a stale :latest; the pin is now passed in,
asserted via docker inspect, and the compose default requires it.
@jghoman
jghoman force-pushed the jakob/pyhoglake-1.3.7 branch from b8423bb to 920f4c8 Compare October 1, 2026 17:56
@jghoman
jghoman requested a balanced review from Copilot October 1, 2026 18:09

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Copilot review overview

🟡 Changes recommended

Cache and receipt handling can compromise publication correctness, and the retry budget understates liveness exposure.

Review effort: Balanced
Findings: 4 High severity · 2 Medium severity

Open (6)
What changed in this PR

Upgrades the Hoglake sink to pyhoglake 1.3.7’s in-memory prepared writer path, reducing upload overhead and steady-state catalog reads.

Changes:

  • Uses buffered, concurrent uploads while preserving read_snapshot.
  • Adds cache expiry, typed-refusal recovery, labeled orphan accounting, and revised retry validation.
  • Expands tests and verifies the integration server image.
File Description
uv.lock Locks pyhoglake 1.3.7 with fast-upload dependencies.
tests/​unit/​test_hoglake.py Tests cached shapes, refusal recovery, and orphan accounting.
tests/​unit/​test_config.py Tests upload concurrency and retry budgets.
tests/​integration/​test_hoglake_integration.py Exercises real-server publication and request behavior.
tests/​hoglake_stack/​stack.py Updates the server digest and adds image inspection.
tests/​hoglake_stack/​docker-compose.yaml Requires an explicit server image.
tests/​e2e/​test_hoglake_e2e.py Verifies the running server image.
README.md Documents writer behavior and deployment considerations.
pyproject.toml Raises the client dependency floor and enables fast uploads.
millpond/​metrics.py Adds orphan-reason labels.
millpond/​hoglake.py Implements buffered preparation and cached-shape recovery.
millpond/​config.py Validates concurrency and revises retry budgeting.
AGENT.md Updates writer and dependency documentation.

💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.

Comment thread millpond/hoglake.py
Comment thread millpond/hoglake.py
Comment thread millpond/hoglake.py Outdated
Comment thread millpond/hoglake.py
Comment thread millpond/config.py
Comment thread tests/integration/test_hoglake_integration.py Outdated
…sal, full reset on table_recreated, re-reconcile adopted specs, runtime flush deadline

The sink's cached shape must never outlive pyhoglake's writer cache:
pyhoglake invalidates on any non-retryable refusal whose basis it
cached, on 410 and on re_prepare, so the sink now drops its shape at
the top of every HoglakeError path out of commit_prepared, the
reused-key accept included. A table_recreated refusal resets the
handle and uuid as well, so a retry without reset_caches() rebuilds
against the new incarnation instead of re-presenting the dead one.
A shape adopted through schema evolution or the alignment self-heal
is reconciled against config whenever its spec identity differs from
the last verdict, with no table read.

main._write_with_retry stops retrying when the next attempt cannot
finish inside the liveness budget (450 s with the margin), re-raises
the last error and counts errors_total{type="flush_deadline_exceeded"};
the startup arithmetic stays as the model, the deadline is the
enforcement.

The wire test selects table reads by endpoint, which exposed one read
the old filter hid: pyhoglake's Namespace.table() takes no totals
parameter, so a resolve pays the live-totals scan once per resolve.
The receipt-before-incarnation ordering on replay is documented as
the intended exactly-once-by-key semantics.
@jghoman
jghoman merged commit cbb1ba0 into main Oct 1, 2026
17 checks passed
@jghoman
jghoman deleted the jakob/pyhoglake-1.3.7 branch October 1, 2026 18:49
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants