Repository navigation
Compact parquet per hour partition to preserve external table hive partitioning - #11
Merged
Merged
Conversation
tomaslink
force-pushed
the
feature/compact-per-hour-partition
branch
4 times, most recently
from
July 17, 2026 21:19
db619da to
4e5e6f5
Compare
…rtitioning Compaction previously collapsed the event_hour= partition level when merging a date's parquet files, leaving compacted dates at a different partition depth than uncompacted ones. BigQuery requires every file under the external table's URI prefix to share the same partition-key structure, whether auto-detected or explicitly declared, so this broke the pipe-nmea-parsed external table once compaction ran on a date. Compactor now discovers event_hour= subpartitions per date and always preserves them when found, compacting each hour independently and writing output back under its own event_hour=HH/ folder. This isn't config-gated: a stale or wrong caller declaration can never silently collapse a real hour partition. The new hourly flag is a narrow assertion rather than a mode switch — for event_sources expected to always be hour-partitioned, it fails loudly if a date has no hour subpartitions at all, instead of silently falling back to flat compaction. Compaction units are modeled explicitly as DailyCompactionUnit and HourlyCompactionUnit (gfw/ops/pipelines/compact_parquet/units.py) — siblings under a shared CompactionUnit base rather than one inheriting the other, since an hourly unit isn't a specialization of a daily one. Bump version to 0.4.0. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
tomaslink
force-pushed
the
feature/compact-per-hour-partition
branch
from
July 17, 2026 22:01
4e5e6f5 to
113aa4f
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
https://globalfishingwatch.atlassian.net/browse/PIPELINE-4331
Summary
Compaction was collapsing the
event_hour=partition level when merging a date's parquet files down to one, leaving compacted dates at a different hive-partition depth (event_source=/event_date=) than uncompacted ones (event_source=/event_date=/event_hour=). Thepipe-nmea-parsedexternal table is hive-partitioned over this path, and BigQuery requires every file under the table's URI prefix to match the same partition-key path structure — whether that structure is auto-detected or explicitly declared. Once compaction stopped writing anevent_hour=directory for a date, files for that date no longer matched the partition layout the rest of the table used, breaking the external table for that range regardless of partitioning mode.Compactornow discoversevent_hour=subpartitions per date from GCS and always preserves them when found — compacting each hour independently, output written back under its ownevent_hour=HH/folder. This isn't gated by config: it's structurally impossible to accidentally collapse a real hour partition, regardless of how a caller is configured.hourlyflag is a narrow assertion, not a mode switch: for event_sources actively expected to be hour-partitioned (e.g. live streaming sources), it asserts a date must have hour subpartitions and raises if none are found — signalling ingestion broke or a stale/wrong declaration, rather than silently falling back to flat compaction. It has no effect when hour subpartitions are actually present — those are always preserved either way.DailyCompactionUnit/HourlyCompactionUnit(a shared abstractCompactionUnitbase, siblings rather than one inheriting the other), moved to their owncompact_parquet/units.pymodule.DailyCompactionUnit.with_hour(hour)derives the corresponding hourly unit for the same date.Compactor.run()builds one flat list of units across the whole date range and loops over it once.gfw-opsversion to0.4.0.Test plan
tests/pipelines/compact_parquet/test_main.pycovers: hour discovery/sorting, per-unitpath()/__str__for both unit types,with_hour()derivation, the preserve-hours-regardless-of-flag behavior, thehourly=True-with-no-hours-found hard failure, and the existing swap/manifest/retry/copy-mode flows against the new unit-based APIgfw-opssuite passes (101 tests)ruff/mypyshow only pre-existing, unrelated debt (mypy actually dropped from 9 → 7 pre-existing findings as a side effect of this refactor) — no new violations introduced.