t3-code-android-nightly/.repos/alchemy-effect/examples/gcp-event-pipeline
Julius Marminge 6f9cea00ae
chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170)
Co-authored-by: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
2026-10-05 13:22:30 -07:00
..
src chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00
test chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00
alchemy.run.ts chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00
package.json chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00
README.md chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00
tsconfig.json chore(refs): sync Effect and Alchemy references to 4.0.1 and beta.80 (#16170) 2026-10-05 13:22:30 -07:00

Event pipeline: Pub/Sub to BigQuery

An analytics pipeline. A public Cloud Run service accepts events and publishes them to Pub/Sub. A Cloud Run Job drains the backlog into BigQuery in batches. Nothing touches BigQuery on the request path, so a slow warehouse cannot take ingestion down.

curl -X POST "$URL/events" -H 'content-type: application/json' -d '{"type":"signup","payload":{"plan":"pro"}}'
# 202 {"id":"5f0c…"}
curl -X POST "$URL/drain"
# 202 {"execution":"projects/…/executions/…"}
curl "$URL/events/count?type=signup"
# 200 {"count":1}

Architecture

  • Events — Pub/Sub topic every event is published to.
  • Inbox — pull subscription on Events. It holds messages until they are acked, so the drain can run whenever a batch is due.
  • Analytics / EventsTable — BigQuery dataset and events table (id, type, occurredAt, payload as JSON).
  • Ingest — public GCP.Function (a Cloud Run service):
    • POST /events — publish { type, payload } to the topic.
    • POST /drain — start a Drain execution now.
    • GET /events/count?type= — count rows in BigQuery.
  • Drain — GCP.Run.Job. It pulls batches of up to 100 messages, inserts them into the table, and acks them, until a pull comes back empty.

The drain inserts first and acks second. A crash between the two redelivers the batch; BigQuery insertIds collapse the duplicates inside its dedup window.

Bindings and IAM

Host Binding IAM granted to the host's service account
Ingest GCP.PubSub.WriteTopic(Events) + WriteTopicHttp roles/pubsub.publisher on the topic
Ingest GCP.BigQuery.ReadTable(EventsTable) + ReadTableHttp roles/bigquery.dataViewer on the table; roles/bigquery.jobUser on the project (query jobs can only be granted there)
Ingest GCP.Run.RunJob(Drain) + RunJobHttp roles/run.jobsExecutorWithOverrides on the Drain job
Drain GCP.PubSub.ReadSubscription(Inbox) + ReadSubscriptionHttp roles/pubsub.subscriber on the subscription
Drain GCP.BigQuery.WriteTable(EventsTable) + WriteTableHttp roles/bigquery.dataEditor on the table

Deploy

Credentials come from your alchemy profile: run alchemy profile once and pick GCP (Service account JSON for a key file, or Stored for an access token or key kept in ~/.alchemy/credentials, plus a default region), then deploy with --profile <name>.

From the repository root:

pnpm install
cd examples/gcp-event-pipeline
pnpm deploy --profile <name>

Docker must be running, because both hosts are built locally from main. To drain on a schedule instead of on demand, point a Cloud Scheduler job at Drain.

Live test

ALCHEMY_PROFILE=<name> bun test --timeout 1200000

The test deploys the stack and publishes events, checking that nothing reaches BigQuery before a drain. It then starts Drain through POST /drain and polls GET /events/count until the rows land. At the end it destroys the stack.

A Cloud Run Job execution can take a couple of minutes to start.

Destroy

pnpm destroy --profile <name>

This deletes the service, the job, their IAM grants, the subscription, the topic, and the dataset with its table.