Ingest events into BigQuery
Analytics ingestion has a shape that predates the cloud: accept the event fast, queue it, and write it to the warehouse in batches. Nothing on the request path touches BigQuery, so a schema change or a slow warehouse cannot take the ingest endpoint down — the messages just queue up.
On GCP that is three resources and two hosts:
| Piece | What it is | Why |
|---|---|---|
PubSub.Topic + pull Subscription |
the queue | holds events until they are acked |
GCP.Function (Cloud Run Service) |
the front door | accepts and publishes, scales to zero |
GCP.Run.Job |
the drain | runs to completion, nothing billed between batches |
BigQuery.Dataset + Table |
the warehouse | what analysts query |
This guide walks through examples/gcp-event-pipeline.
The queue and the warehouse
Section titled “The queue and the warehouse”Resources that depend on another resource are built inside an
Effect.gen, because the topic has to exist before its name can be
referenced:
export const Events = GCP.PubSub.Topic("Events", {});
export const Analytics = GCP.BigQuery.Dataset("Analytics", { location: "US-CENTRAL1", forceDestroy: true,});
export const Inbox = Effect.gen(function* () { const topic = yield* Events; return yield* GCP.PubSub.Subscription("Inbox", { topic: topic.name, ackDeadlineSeconds: 60, });});
export const EventsTable = Effect.gen(function* () { const dataset = yield* Analytics; return yield* GCP.BigQuery.Table("EventsTable", { datasetId: dataset.datasetId, tableId: "events", schema: [ { name: "id", type: "STRING", mode: "REQUIRED" }, { name: "type", type: "STRING", mode: "REQUIRED" }, { name: "occurredAt", type: "TIMESTAMP", mode: "REQUIRED" }, { name: "payload", type: "JSON" }, ], });});Yielding Inbox from two different hosts is free: resources are keyed
by logical id, so both get the same subscription rather than two.
payload is a single JSON column instead of a column per attribute,
so a producer can add a field without a schema migration.
The front door
Section titled “The front door”export default class Ingest extends GCP.Function<Ingest>()( "Ingest", { main: import.meta.url, location: "us-central1", invokerIamDisabled: true, }, Effect.gen(function* () { const topic = yield* Events; const table = yield* EventsTable; const drain = yield* Drain;
// pubsub.publisher on the topic only. const publisher = yield* GCP.PubSub.WriteTopic(topic); // bigquery.dataViewer on the table, plus bigquery.jobUser on the // project — BigQuery only grants running query jobs there. const warehouse = yield* GCP.BigQuery.ReadTable(table); // Binding a Job to a Service grants run.jobsExecutorWithOverrides on // the job — this is how one host triggers another. const runDrain = yield* GCP.Run.RunJob(drain);
// An accessor: the table id is bound at deploy time and read back // inside the handler. const tableId = yield* table.tableId;The two data bindings are access-level clients: the service may publish to the topic and read the table, and nothing else. It cannot pull from the subscription or insert rows — that is the drain’s job.
GCP.Run.RunJob(drain) is the interesting one: binding a Job to a
Service grants roles/run.jobsExecutorWithOverrides on that one Job
to the service’s runtime account and returns a callable. That is how
one host triggers another, with no job name to interpolate and no IAM
policy to write.
POST /events
Section titled “POST /events” const event: EventRow = { id: crypto.randomUUID(), type: body.type, occurredAt: new Date().toISOString(), payload: JSON.stringify(body.payload ?? {}), };
yield* publisher .publish({ data: JSON.stringify(event), // Attributes are queryable without decoding the body, // which lets a filtered subscription fan out by type. attributes: { type: event.type }, }) .pipe(Effect.orDie);
return yield* HttpServerResponse.json( { id: event.id }, { status: 202 }, );publish takes a string or bytes and handles Pub/Sub’s base64
encoding; it returns the server-assigned message id.
202 rather than 201: the event is accepted, not yet stored.
POST /drain
Section titled “POST /drain” if (request.method === "POST" && url.pathname === "/drain") { const operation = yield* runDrain().pipe(Effect.orDie); return yield* HttpServerResponse.json( { execution: operation.name ?? null }, { status: 202 }, ); }Starting a Cloud Run Job returns as soon as the execution is created, not when it finishes. In production, Cloud Scheduler would call this on a cron; the route exists so you can force a batch by hand.
GET /events/count
Section titled “GET /events/count” const events = yield* tableId; // Unqualified table names resolve against the bound table's // dataset; `params` become named `@type` parameters. const rows = yield* warehouse .query( type ? `SELECT COUNT(*) AS n FROM \`${events}\` WHERE type = @type` : `SELECT COUNT(*) AS n FROM \`${events}\``, type ? { type } : undefined, ) .pipe(Effect.orDie);
return yield* HttpServerResponse.json({ count: Number(rows[0]?.n ?? 0), });ReadTable.query sets defaultDataset from the bound table, so the
bare table name resolves, and returns rows decoded to plain JavaScript
(INT64 → number). Named parameters keep a caller’s ?type= out of
the SQL text.
The drain
Section titled “The drain”GCP.Run.Job is the same platform with a run entry instead of
fetch. It starts, does the work, and exits.
export default class Drain extends GCP.Run.Job<Drain>()( "Drain", { main: import.meta.url, location: "us-central1", }, Effect.gen(function* () { const inbox = yield* Inbox; const table = yield* EventsTable;
// pubsub.subscriber on the subscription, and bigquery.dataEditor on // the table — nothing on the project. const subscription = yield* GCP.PubSub.ReadSubscription(inbox); const warehouse = yield* GCP.BigQuery.WriteTable(table);
return { // A pull may return fewer messages than are waiting, so drain // batch by batch until one comes back empty. run: Effect.gen(function* () { // A pull waits for messages; an empty subscription answers // nothing, so a quiet 10 seconds means the backlog is drained. const received = yield* subscription .pull({ maxMessages: BATCH }) .pipe(Effect.timeoutOption("10 seconds")); const messages = Option.getOrElse(received, () => []); if (messages.length === 0) { yield* Effect.log("drain: nothing left"); return 0; }
const rows = messages.map( (message) => JSON.parse(message.text) as EventRow, );
// insertIds make the streaming insert idempotent inside // BigQuery's dedup window, so a redelivered batch collapses. yield* warehouse.insert(rows, { insertIds: rows.map((row) => row.id) });
yield* subscription.acknowledge( messages.map((message) => message.ackId), );
yield* Effect.log(`drain: wrote ${rows.length} row(s)`); return messages.length; }).pipe( Effect.repeat({ until: (count) => count === 0 }), Effect.asVoid, Effect.orDie, ), }; }).pipe( Effect.provide([ GCP.PubSub.ReadSubscriptionHttp, GCP.BigQuery.WriteTableHttp, ]), ),) {}pull decodes each message for you (text, data, attributes,
ackId), and insert takes plain row objects. Because the table
depends on the dataset, binding the table is enough to make the job
wait for both before its first run.
Try it
Section titled “Try it”Both hosts are built from main, so Docker has to be running.
cd examples/gcp-event-pipelinebun alchemy deploycurl -X POST "$URL/events" -H 'content-type: application/json' \ -d '{"type":"signup","payload":{"plan":"pro"}}'# → { "id": "…" } (202)
curl "$URL/events/count?type=signup" # → { "count": 0 }curl -X POST "$URL/drain" # → { "execution": "…" } (202)# a few seconds latercurl "$URL/events/count?type=signup" # → { "count": 1 }test/integ.test.ts publishes three events under a unique type,
asserts BigQuery has not seen them, starts the drain, and polls the
count until the rows land.
Where to take it
Section titled “Where to take it”- Schedule the drain.
GCP.CloudScheduler.Jobcalling the samerun.jobs.runendpoint turns this into a batch pipeline with no manual step. - Fan out by type. A second
Subscriptionwith a filter on thetypeattribute routes a subset of events elsewhere without touching the producer. - Push instead of pull. Set
pushConfig.pushEndpointon the subscription to a Cloud Run URL if you want per-message delivery rather than batches.
- Serve an API on Cloud Run — the request-path counterpart.
- How bindings grant IAM.
- PubSub.Topic, PubSub.Subscription, BigQuery.Table, Run.Job reference.