A data lakehouse combines the cheap, flexible storage of a data lake with the structure and speed of a data warehouse. On Google Cloud, a lakehouse can take the form of the BigQuery ecosystem itself, BigQuery’s own storage, reachable from outside through federated queries and BigLake. That works., but will get expensive at scale. I’ve watched teams hit that wall more than once. This also results in vendor lock-in, with proprietary storage format and proprietary compute engine. An open lakehouse breaks both locks. Use open table formats such as Iceberg or Hudi for storage, and whatever engine you want, whether Spark, Flink, Trino, or BigQuery itself, for reading, processing, and writing. Cost and lock-in. That’s what pushed the open lakehouse into existence.

Flink fits well into that picture, because it does streaming and batch natively in the same engine, with the same SQL dialect for both. That means one system can cover what would otherwise need two, a streaming pipeline and a batch pipeline running side by side. This article walks through writing to Iceberg from Flink SQL, streaming and batch both, through Google’s Lakehouse Iceberg REST catalog.
By the end you’ll have Flink running on a Compute Engine VM, the Lakehouse Iceberg REST catalog registered as a Flink SQL catalog, a streaming job reading mock events off Pub/Sub and writing an Iceberg table, that same table queried straight from BigQuery with no import step, and a batch job rebuilding a daily aggregate from it. The full source is a GitHub repo, so you can run this yourself.
Get the code
Every script and SQL file in this article lives in the companion repository:
https://github.com/sireeshapulipati/flink-iceberg-gcp-lakehouse
The README has detailed instructions and exact commands for every step, in order.
.env.example settings for every script
scripts/01_gcp_setup.sh bucket, catalog, Pub/Sub, service account, VM
scripts/02_install_flink.sh Java, Flink, and jars on the VM
scripts/03_run_streaming.sh submits the streaming job
scripts/04_run_batch.sh runs the batch job
scripts/cleanup.sh deletes the resources
sql/00_init.sql registers the Iceberg REST catalog
sql/10_streaming.sql Pub/Sub to Iceberg streaming job
sql/20_batch.sql daily aggregate rebuild in batch mode
sql/30_bigquery_verify.sql queries to run in BigQuery
publisher/publish_events.py mock clickstream events to Pub/Sub
The steps below explain what each part does and show the Flink SQL in full. Setup commands, the gcloud calls, the VM install, the publisher's shell bits, all live in the repo's scripts instead of getting repeated here.
Versions used
These are the versions used throughout this tutorial, and they need to line up with each other.
- Apache Flink: 1.20.x (LTS), Scala 2.12 build
- Java: 17
- Apache Iceberg: 1.11.0 (iceberg-flink-runtime-1.20, iceberg-gcp-bundle)
- Pub/Sub SQL connector: flink-sql-connector-gcp-pubsub 1.1.0-1.20
Iceberg 1.11.0 ships Flink runtimes for 1.20, 2.0, and 2.1, while the Pub/Sub connector only covers 1.20, 2.2, and 2.3. Flink 1.20 is the one version both sides support, which is the only reason this tutorial is pinned to it. The connector itself is worth flagging too: it’s an independent open-source project, not something Google or Apache maintains directly. It’s Apache 2.0 licensed and published on Maven Central under the group io.github.flink-gcp.
Flink for Streaming and Batch
Flink SQL runs both modes without switching engines, which is the whole reason it fits this kind of pipeline.
Streaming mode treats a source as unbounded. A statement like INSERT INTO keeps running continuously, processing new records as they arrive and committing results incrementally, a new Iceberg snapshot at every checkpoint, for example.
Batch mode treats the source as bounded instead. The same SELECT and INSERT INTO syntax applies, but the query runs once, over a fixed snapshot of the data, computes its result, and finishes. Nothing about the SQL dialect changes between the two, only the execution semantics.
Set up Flink on Google Cloud
Flink runs on Google Cloud a few different ways.
Compute Engine VM with standalone Flink. Tutorials, demos, and quick proofs of concept. I use this option here because it needs no cluster manager and gets you to a running job the fastest.
Dataproc with the Flink optional component. Fits teams already running Dataproc. Create the cluster with –optional-components=FLINK and confirm the Flink version bundled with the image matches the Iceberg runtime jar you download.
GKE with the Apache Flink Kubernetes Operator. Fits production deployments, session or application clusters, upgrades, and autoscaling.
This walkthrough uses the Compute Engine VM option, for the reason already given above, and the six steps below build the whole pipeline on top of it.
Step 1: Create the Google Cloud resources
Full commands for this step live in scripts/01_gcp_setup.sh in the repo, and it's safe to rerun if something fails partway through; it simply picks up where it left off. Broadly, the script turns on the APIs you'll need (BigLake, Pub/Sub, Compute Engine, BigQuery), creates a Cloud Storage bucket, and then creates the Lakehouse Iceberg REST catalog against that bucket in credential vending mode. Credential vending means the catalog hands out short-lived storage credentials at query time, so the Flink job itself never needs direct access to the bucket, only to the catalog that fronts it.
Once the catalog exists, there’s one manual step. Open its details page in the Lakehouse console, copy the auto-provisioned service account it shows you, and grant that account Storage Object User (roles/storage.objectUser) on the bucket. That account can take a minute or so to actually propagate through IAM, so if the grant fails the first time, just wait a bit and retry it.
The script also creates a clickstream-events Pub/Sub topic and subscription, a flink-sa service account for the VM to run as, and finally the flink-vm instance itself, running as that account. flink-sa gets four roles: roles/biglake.editor, roles/pubsub.subscriber, roles/pubsub.viewer, and roles/serviceusage.serviceUsageConsumer. The viewer and subscriber pair specifically covers what the Pub/Sub source needs on an existing subscription, pubsub.subscriptions.get and pubsub.subscriptions.consume. The same IAM propagation lag that affects the catalog's service account applies here too, which is why the script retries these grants automatically instead of failing on the first attempt.
Step 2: Install Flink on the VM
With the resources in place, the next step happens on the VM itself. SSH into the one Step 1 just created, gcloud compute ssh flink-vm –zone=us-central1-a, and stay in that session, since you'll come back to it again in Step 4. From there, run scripts/02_install_flink.sh (full commands are in the repo). It installs Java 17, downloads Flink 1.20.4, and drops three jars into lib/: iceberg-flink-runtime-1.20, iceberg-gcp-bundle, and flink-sql-connector-gcp-pubsub. It also pulls in a set of Hadoop-related jars that the Iceberg catalog needs even though nothing in this setup ever touches HDFS, and it adjusts a task-slot setting so the streaming job from Step 4 and the batch job from Step 6 can both run on this one VM at the same time.
All this does not need a key file anywhere, since the VM’s attached service account supplies application default credentials to every connector automatically. Flink’s web UI runs on port 8081, and you can reach it from your laptop with an SSH tunnel: gcloud compute ssh flink-vm –zone=us-central1-a — -L 8081:localhost:8081.
Step 3: Publish mock events to Pub/Sub
The VM is ready, but there’s nothing actually flowing through it yet. This step fixes that. Create publish_events.py on the VM.
import json
import random
import time
import uuid
from datetime import datetime, timezone
from google.cloud import pubsub_v1
PROJECT_ID = "my-project"
TOPIC_ID = "clickstream-events"
EVENT_TYPES = ["page_view", "add_to_cart", "checkout", "purchase"]
publisher = pubsub_v1.PublisherClient()
topic_path = publisher.topic_path(PROJECT_ID, TOPIC_ID)
def make_event():
now = datetime.now(timezone.utc)
return {
"user_id": f"user_{random.randint(1, 500)}",
"event_type": random.choice(EVENT_TYPES),
"event_ts": now.strftime("%Y-%m-%d %H:%M:%S.") + f"{now.microsecond // 1000:03d}",
"session_id": uuid.uuid4().hex[:8],
}
while True:
for _ in range(20):
publisher.publish(topic_path, json.dumps(make_event()).encode("utf-8"))
time.sleep(1)
Run it in a second SSH session rather than the one from Step 2, since that one needs to stay free for Flink itself. Set up a virtual environment, install google-cloud-pubsub, and run the script (full commands are in publisher/ in the repo). It publishes roughly 20 events a second, and the timestamps come out formatted as yyyy-MM-dd HH:mm:ss.SSS, which is the default format Flink's JSON reader expects.
Step 4: Stream from Pub/Sub into Iceberg
With events actually flowing, it’s time to build the Flink side that consumes them. Everything below, all the way through the streaming INSERT INTO, is exactly what's in sql/10_streaming.sql, and scripts/03_run_streaming.sh runs the whole thing in one shot rather than you typing each line in. In order, the script sets the checkpoint interval, creates a temporary table that reads from the Pub/Sub subscription, registers the Lakehouse Iceberg catalog, creates the namespace and the clickstream table inside it, and then starts a continuous INSERT INTO that streams from the Pub/Sub table into the Iceberg one. What follows walks through each of those statements so you know what it's actually doing.
Set checkpointing first, since two separate things depend on this one setting: the Pub/Sub source only acknowledges messages once a checkpoint completes, and the Iceberg sink only commits a new snapshot at that same moment. Set the interval once, and you’ve effectively set both.
SET 'execution.checkpointing.interval' = '60s';
Next, create the Pub/Sub source table in Flink’s default catalog. It’s a temporary table, so it only lives for the current session.
CREATE TEMPORARY TABLE default_catalog.default_database.clickstream_source (
user_id STRING,
event_type STRING,
event_ts TIMESTAMP(3),
session_id STRING,
message_id STRING METADATA FROM 'message-id' VIRTUAL
) WITH (
'connector' = 'pubsub',
'project' = 'my-project',
'subscription' = 'clickstream-sub',
'format' = 'json'
);
Now register the Lakehouse Iceberg REST catalog itself. These properties come straight from Google’s own Spark configuration for this same endpoint, just translated into Flink’s CREATE CATALOG syntax.
CREATE CATALOG biglake_catalog WITH (
'type' = 'iceberg',
'catalog-type' = 'rest',
'uri' = 'https://biglake.googleapis.com/iceberg/v1/restcatalog',
'warehouse' = 'bl://projects/my-project/catalogs/my-streaming-catalog',
'header.x-goog-user-project' = 'my-project',
'rest.auth.type' = 'org.apache.iceberg.gcp.auth.GoogleAuthManager',
'io-impl' = 'org.apache.iceberg.gcp.gcs.GCSFileIO',
'header.X-Iceberg-Access-Delegation' = 'vended-credentials'
);
warehouse takes the bl:// form for multiple-bucket catalogs and gs://bucket-name for single-bucket ones, and GoogleAuthManager needs Iceberg 1.10 or newer to work.
With the catalog registered, create the namespace and the table it will hold.
USE CATALOG biglake_catalog;
CREATE DATABASE IF NOT EXISTS events;
USE events;
CREATE TABLE clickstream (
user_id STRING,
event_type STRING,
event_ts TIMESTAMP(3),
session_id STRING,
message_id STRING,
event_date DATE
) PARTITIONED BY (event_date)
WITH ('format-version' = '2');
Iceberg tables in the Lakehouse runtime catalog use format version 2, which is why that property is set explicitly. And because Flink’s DDL partitions by plain columns rather than expressions, the table carries its own event_date column, derived from the event timestamp so there's something to partition on.
Everything so far has just been setup. This next statement is the one that actually starts the streaming write.
INSERT INTO clickstream
SELECT user_id, event_type, event_ts, session_id, message_id,
CAST(event_ts AS DATE) AS event_date
FROM default_catalog.default_database.clickstream_source;
Because the source is unbounded, this statement kicks off a continuous job rather than running once and finishing. You’ll see it show up as running in the Flink web UI, and a new snapshot appears after every checkpoint. The message_id column isn't just carried along for reference. It's there specifically so you can deduplicate downstream, since the Pub/Sub source only guarantees at-least-once delivery.
Step 5: Query the table from BigQuery
Data is landing in Iceberg now, so it’s worth checking that BigQuery can actually reach it. Tables in a Lakehouse runtime catalog use a four-part name there: project.catalog.namespace.table.
SELECT event_type, COUNT(*) AS event_count
FROM `my-project.my-streaming-catalog.events.clickstream`
GROUP BY event_type
ORDER BY event_count DESC;
BigQuery is reading the exact Iceberg metadata and Parquet files Flink wrote to Cloud Storage. There’s no import step, and no second copy of the data sitting inside BigQuery’s own storage.

Step 6: Rebuild a daily aggregate in batch mode
This step adds the batch side, without touching what’s already running. There’s no need to cancel Step 4’s streaming job first. This runs right alongside it, which is exactly what the second task slot from Step 2 was for.

Everything below is what’s in sql/20_batch.sql, and scripts/04_run_batch.sh runs it in one shot. The script re-registers the Lakehouse Iceberg catalog first, since a new SQL client session starts with nothing registered (CREATE CATALOG IF NOT EXISTS makes that safe to repeat even if Step 4's session is technically still open somewhere), switches execution to batch mode, creates the clickstream_daily table if it doesn't already exist, and runs an INSERT OVERWRITE that aggregates the streaming table by day and event type.
CREATE CATALOG IF NOT EXISTS biglake_catalog WITH (
'type' = 'iceberg',
'catalog-type' = 'rest',
'uri' = 'https://biglake.googleapis.com/iceberg/v1/restcatalog',
'warehouse' = 'bl://projects/my-project/catalogs/my-streaming-catalog',
'header.x-goog-user-project' = 'my-project',
'rest.auth.type' = 'org.apache.iceberg.gcp.auth.GoogleAuthManager',
'io-impl' = 'org.apache.iceberg.gcp.gcs.GCSFileIO',
'header.X-Iceberg-Access-Delegation' = 'vended-credentials'
);
USE CATALOG biglake_catalog;
USE events;
Now switch the session to batch mode. In this mode, the query reads one snapshot of the table, and INSERT OVERWRITE replaces only the partitions that query actually produces, which is exactly what makes it safe to rerun as many times as you need to.
SET 'execution.runtime-mode' = 'batch';
SET 'table.dml-sync' = 'true';
CREATE TABLE IF NOT EXISTS clickstream_daily (
event_date DATE,
event_type STRING,
event_count BIGINT
) PARTITIONED BY (event_date)
WITH ('format-version' = '2');
INSERT OVERWRITE clickstream_daily
SELECT event_date, event_type, COUNT(*) AS event_count
FROM clickstream
GROUP BY event_date, event_type;

Unlike the streaming job, this one blocks. The client sits there until the job actually finishes, then hands control back to you. Query my-project.my-streaming-catalog.events.clickstream_daily from BigQuery afterward and you'll see the result sitting there already.

In a real production setup, you'd trigger this job on a schedule instead of running it by hand, through something like Cloud Composer or Cloud Scheduler.
Production notes
That covers the pipeline working end to end, streaming and batch both, but a demo running cleanly once isn’t the same thing as it holding up in production. A few things are worth knowing before you take any of this further.
Small files build up faster than you’d expect here. Every checkpoint writes new data files, so at a 60-second interval, a single partition can end up with a lot of files by the end of just one day. It’s worth scheduling compaction for this (Iceberg’s rewrite data files action, run from a Spark job, works well), and choosing a checkpoint interval that actually matches how fresh you need the data to be, rather than whatever felt reasonable when you were first setting things up.
The Pub/Sub source only guarantees at-least-once delivery, not exactly-once, so duplicates are possible. Deduplicating on message_id handles this, either in your downstream queries directly or through a separate batch cleanup job.
There’s a hard limit worth knowing about too: the Lakehouse runtime catalog caps the Iceberg metadata.json file at 1MB, and it only supports Iceberg format version 2, not version 1.
The Java classpath issue from Step 2 is worth repeating here as well, since it’s easy to hit if you’re working from a different setup than this one: the Iceberg REST catalog needs those Hadoop-related jars on the classpath, or CREATE CATALOG simply fails, even though nothing in this pipeline ever touches HDFS.
You can cleanup at the end using: scripts/cleanup.sh. It tears down the VM, subscription, topic, service account, catalog, bucket, all of it. Killing the VM kills the Flink cluster too, so the streaming job and the publisher just stop, nothing to cancel by hand first.
Flink SQL and Apache Iceberg on Google Cloud: Streaming and Batch Into One Open Lakehouse was originally published in Google Cloud – Community on Medium, where people are continuing the conversation by highlighting and responding to this story.
Source Credit: https://medium.com/google-cloud/flink-sql-and-apache-iceberg-on-google-cloud-streaming-and-batch-into-one-open-lakehouse-5521296f627c?source=rss—-e52cf94d98af—4
