Skip to content

Iceberg tables

Cloudflare’s catalog (formerly R2 Data Catalog) turns an R2 bucket into an Apache Iceberg warehouse. Iceberg stores a table as data files plus metadata, so any engine that speaks the Iceberg REST protocol can read and write it: Spark, DuckDB, PyIceberg, Snowflake, Trino, and others.

The catalog is part of Basin, Cloudflare’s analytics stack on R2. Tables get into it in two ways:

  • A Pipelines catalog sink creates a table and writes events into it.
  • Cloudflare.Basin.Table creates and manages a table directly, for tables your own code or engines write.

This page covers enabling the catalog, managing tables, and connecting query engines.

alchemy.run.ts
const lake = yield* Cloudflare.R2.Bucket("Lake");
const catalog = yield* Cloudflare.Basin.Catalog("Catalog", {
bucket: lake,
});

Basin.Catalog is the same resource as Cloudflare.R2.DataCatalog. It outputs:

  • catalogUri: the Iceberg REST endpoint, for example https://catalog.cloudflarestorage.com/<account>/<bucket>.
  • name: the warehouse name, <account>_<bucket>.

Iceberg tables accumulate small files and old snapshots. The catalog can compact files and expire snapshots for every table in the warehouse:

const catalog = yield* Cloudflare.Basin.Catalog("Catalog", {
bucket: lake,
compaction: { state: "enabled", targetSizeMb: "128" },
snapshotExpiration: { state: "enabled", maxSnapshotAge: "7d", minSnapshotsToKeep: 10 },
token: yield* Config.Redacted("CATALOG_MAINTENANCE_TOKEN"),
});

The maintenance jobs run with token, a Cloudflare API token with access to R2 and the catalog. They stay pending until a token is set. Only the fields you set are enforced; the rest keep Cloudflare’s defaults.

Cloudflare.Basin.Table creates a namespace and table through the catalog’s Iceberg REST API, with columns from an Effect Schema:

import * as Schema from "effect/Schema";
const Order = Schema.Struct({
orderId: Schema.String,
amountCents: Schema.Number,
at: Schema.Date,
note: Schema.optional(Schema.String),
});
const orders = yield* Cloudflare.Basin.Table("Orders", {
catalog,
namespace: "sales",
schema: Order,
partitionBy: ["day(at)"],
properties: { "write.format.default": "parquet" },
});

The table name is generated from the stack, stage, and logical ID unless you set name. The table outputs identifier (for example sales.orders), tableUuid, location, and metadataLocation.

partitionBy takes Iceberg transforms: a column name, day(at), month(at), bucket[16](orderId), and so on. properties are Iceberg table properties, updated in place.

The table authenticates with your deploying Cloudflare credentials. Pass token to use a different one.

Adding an optional column evolves the table in place. Existing data stays, and the new column reads as empty for old rows:

const Order = Schema.Struct({
orderId: Schema.String,
amountCents: Schema.Number,
at: Schema.Date,
note: Schema.optional(Schema.String),
coupon: Schema.optional(Schema.String),
});

A change that would lose data fails the plan before anything is applied: dropping or narrowing a column, or adding a required one. Changing namespace, name, or partitionBy also fails the plan, because each would mean dropping the table.

maintenance overrides the catalog’s settings for one table:

const orders = yield* Cloudflare.Basin.Table("Orders", {
catalog,
namespace: "sales",
schema: Order,
maintenance: {
compaction: { targetSizeMb: "256" },
snapshotExpiration: { olderThan: "3 days", retainLast: 5 },
},
});

Per-table snapshot expiration only works after snapshot expiration is enabled on the catalog.

Destroying a Table resource leaves the table and its data in place by default, so removing it from your program never deletes data:

Cloudflare.Basin.Table("Orders", { catalog, namespace: "sales", schema: Order, delete: "retain" }); // default
Cloudflare.Basin.Table("Orders", { catalog, namespace: "sales", schema: Order, delete: "drop" }); // drop the table, keep its files
Cloudflare.Basin.Table("Orders", { catalog, namespace: "sales", schema: Order, delete: "purge" }); // drop the table and its files

Any Iceberg REST client connects with the catalog’s URI, the warehouse name, and a Cloudflare API token with access to R2 and the catalog. Export the first two from your stack:

return {
catalogUri: catalog.catalogUri,
warehouse: catalog.name,
};

Then point the client at them. With PyIceberg:

from pyiceberg.catalog.rest import RestCatalog
catalog = RestCatalog(
name="lake",
uri=CATALOG_URI,
warehouse=WAREHOUSE,
token=CLOUDFLARE_API_TOKEN,
)
orders = catalog.load_table(("sales", "orders"))
print(orders.scan().to_pandas())

From Effect code, @distilled.cloud/iceberg is a typed client for the same API:

import * as Iceberg from "@distilled.cloud/iceberg";
const table = yield* Iceberg.loadTable({ namespace: "sales", table: "orders" }).pipe(
Effect.provide(Iceberg.IcebergProtocol),
Effect.provide(Iceberg.fromCatalogConfig({ uri: catalogUri, warehouse, token })),
);

Related:

  • Pipelines: ingest events into a table with a catalog sink.
  • R2: the bucket under the warehouse.

Reference: