Skip to content

PIPELINE-4277: Add GCS Parquet compaction pipeline - #9

Merged
tomaslink merged 3 commits into
mainfrom
feature/compact-parquet
Jul 14, 2026
Merged

tomaslink merged 3 commits into
mainfrom
feature/compact-parquet

Conversation

@tomaslink

@tomaslink tomaslink commented Jul 2, 2026 •

Copy link
Copy Markdown
Collaborator

https://globalfishingwatch.atlassian.net/browse/PIPELINE-4277


Summary

Adds two new pipelines for working with hive-partitioned Parquet files on GCS: compact-parquet and benchmark-parquet.

compact-parquet

Compacts small Parquet files within each date partition into larger files using DuckDB. Operates in two modes:

  • Swap mode (default): writes compacted output to an auto-generated staging sibling path, deletes the originals, then copies the staged files back so the external table path never changes. If the process is interrupted between the delete and copy steps, the next run detects the staging files and resumes automatically.
  • Copy mode (--gcs-staging-path): writes compacted output to a separate path, leaving source files untouched — useful when you want both uncompacted and compacted versions available via separate external tables.

DuckDB memory and thread count are configurable via --memory-limit and --threads.

benchmark-parquet

Runs configurable queries against one or more BigQuery tables and prints a comparison of bytes processed, slot milliseconds, and wall-clock time. Supports SELECT * and an hourly aggregation query, with query cache disabled so each run reflects real scan cost. Designed to compare native BigQuery tables against Parquet external tables.

Tests

Both pipelines have 100% test coverage. Compact-parquet tests cover the full compaction flow, swap/copy modes, interrupted swap resume, leftover staging cleanup, and skip logic — without patching private methods.

@tomaslink tomaslink self-assigned this Jul 2, 2026
@tomaslink tomaslink changed the title Add compact-parquet pipeline and CLI command PIPELINE-4277: Add GCS Parquet compaction pipeline Jul 2, 2026
@tomaslink
tomaslink force-pushed the feature/compact-parquet branch 7 times, most recently from c98c67a to d9bfeea Compare July 3, 2026 15:01
@tomaslink
tomaslink changed the base branch from feature/bq-export to main July 5, 2026 15:10
Compacts small hive-partitioned Parquet files on GCS into larger files
using a staging-based swap to keep the original path intact. Adds
pyarrow and gcsfs dependencies.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
@tomaslink
tomaslink force-pushed the feature/compact-parquet branch 2 times, most recently from 0baf2c6 to 6b5e32c Compare July 6, 2026 18:27
@tomaslink
tomaslink requested a review from andres-arana July 6, 2026 18:27
@tomaslink
tomaslink force-pushed the feature/compact-parquet branch 4 times, most recently from bc7757c to c6e3548 Compare July 7, 2026 23:20

@andres-arana andres-arana 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.

I think the big one is the potential dangerous delete. Let's brainstorm a solution for that here. Everything else is minor (but easy to fix).

Comment thread src/gfw/ops/cli/commands/compact_parquet.py
Comment thread src/gfw/ops/pipelines/compact_parquet/main.py Outdated
Comment thread src/gfw/ops/pipelines/compact_parquet/main.py Outdated
Comment thread requirements.txt Outdated
@tomaslink
tomaslink force-pushed the feature/compact-parquet branch from c6e3548 to 09fa16a Compare July 10, 2026 13:49
Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
@tomaslink
tomaslink force-pushed the feature/compact-parquet branch from 09fa16a to b11e11e Compare July 10, 2026 14:01
A crash partway through deleting original blobs left the live partition
with a partial remnant that the old resume check ("staging present, source
empty") couldn't distinguish from a fresh partition, causing the correct
staged replacement to be wiped and only the survivors recompacted.

Record the exact list of originals to delete in a manifest before starting
the delete, so a resumed run can finish from that authoritative list
regardless of how far the previous attempt got. Raise instead of guessing
when staging is non-empty, source is empty, and no manifest exists, since
that state is unverifiable and should not occur post-migration.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
@tomaslink
tomaslink merged commit 62d102d into main Jul 14, 2026
3 checks passed
@tomaslink
tomaslink deleted the feature/compact-parquet branch July 14, 2026 01:15
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