Repository navigation
feat: parquet generation function - #1845
Draft
davidgamez wants to merge 5 commits into
Draft
davidgamez wants to merge 5 commits into
davidgamez wants to merge 5 commits into
Conversation
Contributor
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Concurrency, enqueue, purge, publication, and SQL-safety defects can leave builds duplicated, inaccessible, or incorrectly advertised.
Review effort: Balanced
Findings: 7
Open (10)
Prevent mark_triggered from regressing active or completed builds · New Escape archive filenames before interpolating DuckDB identifiers · New Commit claims before progress updates and rollback callback errors · New Fail builds when making uploaded objects public fails · New Scope expiry cleanup to the currently published artifact version · New Prevent concurrent cleanup from deleting newly published objects · New Reapply volume mount when the function source changes · New Propagate task creation failures and use fresh attempt names · New Rollback the session before recording build failures · New Remove undocumented tables field from the example response · New
What changed in this PR
Adds on-demand GTFS-to-Parquet generation, status APIs, retention cleanup, deployment infrastructure, and local development tooling.
Changes:
- Implements Parquet conversion, publishing, progress tracking, and tests.
- Adds Operations API endpoints and scheduled retention cleanup.
- Configures Cloud Functions, task queues, storage CORS, and local tooling.
| File | Description |
|---|---|
.gitignore |
Ignores local Parquet output. |
README.md |
Documents local feed browsing. |
docs/OperationsAPI.yaml |
Defines Parquet API endpoints and models. |
scripts/parquet-generate-local.sh |
Bootstraps local generation. |
infra/batch/main.tf |
Enables HEAD CORS requests. |
infra/functions-python/main.tf |
Deploys builder, queue, purge, and volume. |
infra/functions-python/vars.tf |
Adds Parquet infrastructure variables. |
functions-python/helpers/ephemeral_workdir.py |
Generalizes work-directory ownership. |
functions-python/helpers/task_execution/task_execution_tracker.py |
Adds exclusive task claims and leases. |
functions-python/helpers/tests/test_task_execution_tracker.py |
Tests task claim behavior. |
functions-python/helpers/utils.py |
Enqueues Parquet build tasks. |
functions-python/operations_api/.openapi-generator/FILES |
Registers generated Parquet API files. |
functions-python/operations_api/README.md |
Documents local Parquet development. |
functions-python/operations_api/src/main.py |
Mounts the Parquet router. |
functions-python/operations_api/src/feeds_operations/impl/parquet_api_impl.py |
Implements status and generation endpoints. |
functions-python/operations_api/tests/feeds_operations/impl/test_parquet_api_impl.py |
Tests endpoint state handling. |
functions-python/parquet_builder/README.md |
Documents builder behavior and usage. |
functions-python/parquet_builder/function_config.json |
Configures the Cloud Function. |
functions-python/parquet_builder/requirements.txt |
Adds runtime dependencies. |
functions-python/parquet_builder/requirements_dev.txt |
Adds test dependencies. |
functions-python/parquet_builder/src/converter.py |
Converts GTFS tables to Parquet. |
functions-python/parquet_builder/src/main.py |
Orchestrates cloud builds and publishing. |
functions-python/parquet_builder/src/progress.py |
Throttles progress updates. |
functions-python/parquet_builder/src/scripts/generate_local.py |
Implements local generation and serving. |
functions-python/parquet_builder/tests/test_converter.py |
Tests conversion contracts. |
functions-python/parquet_builder/tests/test_generate_local.py |
Tests local generation and range serving. |
functions-python/parquet_builder/tests/test_main.py |
Tests build orchestration. |
functions-python/parquet_builder/tests/test_progress.py |
Tests progress throttling. |
functions-python/pmtiles_builder/src/main.py |
Uses the shared work-directory helper. |
functions-python/pmtiles_builder/tests/test_workdir.py |
Updates work-directory tests. |
functions-python/tasks_executor/src/main.py |
Registers the purge task. |
functions-python/tasks_executor/src/tasks/parquet/__init__.py |
Adds the Parquet task package. |
functions-python/tasks_executor/src/tasks/parquet/purge_expired_parquet.py |
Implements retention cleanup. |
functions-python/tasks_executor/tests/tasks/parquet/test_purge_expired_parquet.py |
Tests retention cleanup. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
Comment on lines
+156
to
+159
| tracker.mark_triggered( | ||
| dataset.stable_id, | ||
| metadata={"phase": "start", "done": 0, "total": 0, "detail": ""}, | ||
| ) |
| return False | ||
| try: | ||
| con.execute(f""" | ||
| CREATE VIEW "{table}" AS |
Comment on lines
+219
to
+220
| tracker.heartbeat(dataset_stable_id, metadata=state) | ||
| db_session.commit() |
Comment on lines
+532
to
+537
| try: | ||
| blob.make_public() | ||
| except Exception as error: | ||
| # Uniform bucket-level access would make this unnecessary and impossible at the | ||
| # same time; the objects are public by bucket policy in that case. | ||
| logger.warning("Could not make %s public: %s", blob.name, error) |
| completed_at, | ||
| COALESCE((metadata ->> 'retention_days')::int, :default_days) AS retention_days | ||
| FROM task_execution_log | ||
| WHERE task_name = :task_name |
| purged = [] | ||
| for dataset_stable_id in datasets: | ||
| try: | ||
| deleted_objects += _delete_prefix(bucket, dataset_stable_id) |
Comment on lines
+1696
to
+1701
| triggers_replace = { | ||
| function_name = google_cloudfunctions2_function.parquet_builder.name | ||
| region = var.gcp_region | ||
| project = var.project_id | ||
| size = var.parquet_builder_in_memory_size | ||
| } |
| project_id=project_id, | ||
| gcp_region=gcp_region, | ||
| queue_name=queue_name, | ||
| task_name=None if force else task_name, |
Comment on lines
+275
to
+281
| except Exception as error: | ||
| logger.exception("Parquet build failed for %s", dataset_stable_id) | ||
| # Releases the claim as well as recording why, so the dataset can be retried | ||
| # without waiting out the lease. | ||
| tracker.mark_failed(dataset_stable_id, error_message=str(error)) | ||
| db_session.commit() | ||
| raise |
| ```json | ||
| { "status": "absent", "feed_stable_id": "mdb-1210", "dataset_stable_id": "mdb-1210-202402121801" } | ||
| { "status": "preparing", "phase": "convert", "done": 12, "total": 32, "detail": "stop_times" } | ||
| { "status": "ready", "base_url": "https://.../parquet", "tables": [{"name": "stops"}] } |
This branch has not been deployed
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.



Summary:
Adds on-demand Parquet rendering of GTFS datasets so a browser can query a feed in place over HTTP range requests, without downloading or unpacking it.
parquet_builderCloud Function: converts a dataset to one Parquet file per GTFS table via DuckDB plus amanifest.json, published publicly togs://<datasets-bucket>/<feed>/<dataset>/parquet/. All columns are text; empty CSV fields become NULL./v1/operations/gtfs_feeds/{id}/parquetand/v1/operations/gtfs_datasets/{id}/parquet. GET is read-only and returns 200 for every state includingabsent; 404 means the feed or dataset does not exist. POST is idempotent.TaskExecutionTracker.try_acquire(new conditional-upsert claim ontask_execution_log) is the guarantee, not the Cloud Task name — delivery is at-least-once.manifest.jsonlast, stale objects pruned after. Pinned byTestPublishIsNotObservablyPartial.customTimeofnow + retention_days(default 30, 1..60), and one bucket lifecycle rule deletes ondays_since_custom_time = 0. The API reports an expired row asabsentso the next request rebuilds. No scheduled job.["GET", "HEAD"]— the reader probes each table with HEAD, and without it the browser blocks every request.scripts/parquet-generate-local.shruns the same conversion with no GCP;--serveadds a range-capable, CORS-enabled server.api/src/shared/common/gcp_utils.pygains an opt-inraise_on_erroroncreate_http_task_with_name, default off, used only by the Parquet enqueue. No change to the public Feeds API or the DB schema.Expected behavior:
POST returns 202. Polling GET walks
absent→preparing(phase,done,total,detail; bytes duringdownload, counts otherwise) →readywith abase_url. Failures reportfailedwith a readablemessage, and the function still returns 200 — nothing here fails in a way a retry would fix.On its expiry date the dataset reports
absentand rebuilds on next request. Files are removed by the bucket rule shortly after, on GCS's schedule.{base_url}/manifest.jsonis the only place the table list, counts and sizes live.Testing tips:
Fastest path, no GCP:
scripts/parquet-generate-local.sh mdb-1210 --serve # http://localhost:8090Against the operations web app:
Unit tests:
scripts/api-tests.sh --folder functions-python/parquet_builder.Expiry: POST with
retention_days: 1, thenWorth exercising:
forceagainst a running build, an emptyfrequencies.txt, a corrupt archive, and a browse immediately after a first-ever build. Full endpoint-level walkthrough is infunctions-python/parquet_builder/README.md.Deployment note: the CORS fix and the lifecycle rule are both in
infra/batch, sodatasets-batch-deployermust run for an environment or the viewer stays broken there.Open item:
days_since_custom_time = 0is undocumented as a minimum and wants confirming against a real bucket. Fallback iscustomTime = build_timewith the rule at60; only reclamation timing changes.Please make sure these boxes are checked before submitting your pull request - thanks!
./scripts/api-tests.shto make sure you didn't break anything