Skip to content

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.

Resources that depend on another resource are built inside an Effect.gen, because the topic has to exist before its name can be referenced:

src/resources.ts
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.

src/Ingest.ts
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.

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.

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.

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.

GCP.Run.Job is the same platform with a run entry instead of fetch. It starts, does the work, and exits.

src/Drain.ts
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.

Both hosts are built from main, so Docker has to be running.

Terminal window
cd examples/gcp-event-pipeline
bun alchemy deploy
Terminal window
curl -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 later
curl "$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.

  • Schedule the drain. GCP.CloudScheduler.Job calling the same run.jobs.run endpoint turns this into a batch pipeline with no manual step.
  • Fan out by type. A second Subscription with a filter on the type attribute routes a subset of events elsewhere without touching the producer.
  • Push instead of pull. Set pushConfig.pushEndpoint on the subscription to a Cloud Run URL if you want per-message delivery rather than batches.