diff --git a/README.md b/README.md index 0ce7ed1..e0d911e 100644 --- a/README.md +++ b/README.md @@ -19,7 +19,7 @@ See respective README files in sub-directories for details. - [LoR](lor/README.md): Using Dataflow, Cloud Run and Spanner to explore Lord of the Rings characters. - [Network Digital Twin](telco-and-csp/README.md): Advanced usage of graph to Spanner Graph to model, visualize, and query a complex telecommunications network. - [Transit Fraud Detector](TransitFraud/README.md): Advanced usage of graph capabilities to detect fraud. - +- [Orders Streaming Analytics](orders-streaming-analytics/README.md): Dual-path streaming with Spanner: Precomputed Lakehouse aggregates and Ad-Hoc analytics via Data Boost. ## Notebooks Some of these notebooks are hosted in external Google Cloud repositories. diff --git a/orders-streaming-analytics/.gitignore b/orders-streaming-analytics/.gitignore new file mode 100644 index 0000000..4a598e7 --- /dev/null +++ b/orders-streaming-analytics/.gitignore @@ -0,0 +1,4 @@ +src/java-tpch-stream-generator/target +output +checkpoints +.env diff --git a/orders-streaming-analytics/README.md b/orders-streaming-analytics/README.md new file mode 100644 index 0000000..43a8b61 --- /dev/null +++ b/orders-streaming-analytics/README.md @@ -0,0 +1,27 @@ +# Spanner DataBoost Streaming Demo + +This demo implements a dual-path streaming architecture ingesting orders (TCP-H dataset) where each branch addresses a distinct analytical goal. + +The left branch (Lakehouse Aggregates path) handles predictable, day-to-day reporting by precomputing cumulative aggregates, such as daily order totals, and saving them directly into Spanner aggregate tables only when relevant data changes. + +The right branch (Direct Spanner path) continuously maintains the raw orders table to support ad-hoc queries not covered by precomputations. By leveraging Spanner Data Boost alongside its columnar engine, these ad-hoc queries run efficiently without impacting the primary Spanner instance, thereby eliminating the need to over-provision it. + +![DataBoost Streaming Demo](SpannerStreamingLakehouse-Dataflow.drawio.png) + +### System requirements + +- gcloud cli +- uv https://docs.astral.sh/uv/#installation (for dependencies management) +- jq https://jqlang.org/ (for parsing gcloud json response reliably, used by the bash scripts during the notebook) + + +## Setup + +The entire demo architecture notes, GCP infrastructure provisioning, Spark job sources, producer execution, and dashboard/ad-hoc query code is in [`demo.ipynb`](demo.ipynb). + +```bash +uv sync +uv run jupyter lab +``` + +Open `demo.ipynb` and run the cells top to bottom. diff --git a/orders-streaming-analytics/SpannerStreamingLakehouse-Dataflow.drawio.png b/orders-streaming-analytics/SpannerStreamingLakehouse-Dataflow.drawio.png new file mode 100644 index 0000000..df494ae Binary files /dev/null and b/orders-streaming-analytics/SpannerStreamingLakehouse-Dataflow.drawio.png differ diff --git a/orders-streaming-analytics/demo.ipynb b/orders-streaming-analytics/demo.ipynb new file mode 100644 index 0000000..f8b63a5 --- /dev/null +++ b/orders-streaming-analytics/demo.ipynb @@ -0,0 +1,2959 @@ +{ + "cells": [ + { + "cell_type": "markdown", + "id": "2fc87198", + "metadata": {}, + "source": [ + "# Spanner Data Boost Streaming Demo\n", + "\n", + "![DataBoost Streaming Demo](SpannerStreamingLakehouse-Dataflow.drawio.png)\n", + "\n", + "This notebook is the single source of truth for the demo: it drives GCP\n", + "resource provisioning (via `gcloud`, executed as shell cells), builds and\n", + "submits the Spark streaming jobs, runs the producer, provisions the BigQuery\n", + "reservation and scheduled query, starts the continuous queries, and contains\n", + "the dashboard and ad-hoc query code used for analysis.\n", + "\n", + "**Two execution contexts:**\n", + "\n", + "- **Local cells** (marked `%%bash`) run on your workstation with `gcloud`\n", + " configured, provisioning infrastructure and submitting jobs.\n", + "- **Cluster-side cells** (Part 13) run inside a JupyterLab **Python 3** kernel\n", + " launched from the Managed Spark cluster's Web Interfaces, because they need\n", + " private-network access to Spanner/Iceberg that only exists inside the VPC.\n", + " Copy those cells into a notebook there (or open this same file in that\n", + " JupyterLab instance) and run them with a Python 3 kernel.\n", + "\n", + "## Architecture\n", + "\n", + "```text\n", + "Custom VPC: spark-databoost-demo-vpc\n", + "Subnet: 10.10.0.0/24\n", + " |\n", + " +-- Private Google Access\n", + " |\n", + " +-- Cloud Router + Cloud NAT\n", + " | Outbound access for Maven packages and Docker images\n", + " |\n", + " +-- IAP\n", + " | SSH access to the Kafka VM without a public IP\n", + " |\n", + " +-- Private Kafka VM\n", + " | +-- Container-Optimized OS\n", + " | +-- apache/kafka-native:4.1.2\n", + " | +-- tpch-generator:local (built from source src/java-tpch-stream-generator)\n", + " |\n", + " +-- Private single-node Managed Spark (Dataproc) cluster\n", + " +-- Jupyter optional component\n", + " +-- Kafka, Spanner, and Iceberg dependencies\n", + " +-- Kafka-to-Lakehouse streaming job\n", + " +-- Kafka-to-Spanner streaming job\n", + "\n", + "Cloud Storage:\n", + " +-- Job bucket (staging and checkpoints)\n", + " +-- Warehouse bucket (Iceberg data and metadata)\n", + "\n", + "Lakehouse aggregate path:\n", + " +-- BigQuery scheduled query every 5 minutes\n", + " | +-- Identify aggregate keys affected in the last 5 minutes\n", + " | +-- Recompute cumulative COUNT and SUM snapshots for those keys\n", + " | +-- Append snapshots to three native BigQuery staging tables\n", + " +-- Three BigQuery continuous queries\n", + " +-- Consume staging-table appends with APPENDS()\n", + " +-- Calculate averages with SAFE_DIVIDE\n", + " +-- Upsert snapshots into three Spanner aggregate tables\n", + "\n", + "Spanner:\n", + " +-- Raw orders populated directly from Kafka by Spark\n", + " +-- Three aggregate tables populated through the lakehouse path\n", + "```\n" + ] + }, + { + "cell_type": "markdown", + "id": "b3567d29", + "metadata": {}, + "source": [ + "## 1. Prerequisites and variables\n", + "\n", + "Before starting:\n", + "\n", + "- Run `gcloud auth login`\n", + "- Choose a region that supports all services used by the demo; `ZONE` must belong to `REGION`.\n", + "\n", + "The user running setup needs these roles:\n", + "\n", + "| Category | Role name and ID | Purpose |\n", + "| :--- | :--- | :--- |\n", + "| **IAM and policy** | **Project IAM Admin**
`roles/resourcemanager.projectIamAdmin` | Grant project-level IAM roles to the Spark service account |\n", + "| | **Service Account Admin**
`roles/iam.serviceAccountAdmin` | Create, manage, and delete the service account |\n", + "| | **Service Account User**
`roles/iam.serviceAccountUser` | Attach and impersonate the service account for demo jobs |\n", + "| **Network and compute** | **Compute Network Admin**
`roles/compute.networkAdmin` | Create and delete the VPC, subnet, firewall rules, router, and NAT |\n", + "| | **Compute Instance Admin**
`roles/compute.instanceAdmin.v1` | Create, manage, and delete the Kafka VM |\n", + "| | **IAP-secured Tunnel User**
`roles/iap.tunnelResourceAccessor` | Use SSH and SCP through IAP to reach the private VM |\n", + "| **Data and analytics** | **Cloud Spanner Admin**
`roles/spanner.admin` | Create the Spanner instance, database, schema, and execute SQL |\n", + "| | **Dataproc Admin**
`roles/dataproc.admin` | Create the cluster, submit jobs, access JupyterLab, and delete the cluster |\n", + "| | **BigLake Admin**
`roles/biglake.admin` | Create and delete the Iceberg REST catalog |\n", + "| | **BigQuery Admin**
`roles/bigquery.admin` | Create datasets, tables, reservations, continuous queries, and scheduled queries |\n", + "| **Storage and services** | **Storage Admin**
`roles/storage.admin` | Create and delete buckets and grant bucket IAM |\n", + "| | **Service Usage Admin**
`roles/serviceusage.serviceUsageAdmin` | Enable required APIs |\n", + "\n", + "Edit the values in the next cell. Running it writes a local `.env` file used by later shell, SCP, and SSH cells.\n", + "\n", + "**Everything has a default that works EXCEPT the PROJECT_ID that mandatory must be edited with your project name.**\n", + "\n", + "> **Warning:** `%%writefile` overwrites an existing `.env`." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "51446391", + "metadata": {}, + "outputs": [], + "source": [ + "%%writefile .env\n", + "# !! MANDATORY !! Edit PROJECT_ID\n", + "export PROJECT_ID=\"your-project-id\"\n", + "export REGION=\"us-central1\"\n", + "export ZONE=\"us-central1-a\"\n", + "\n", + "# Network\n", + "export NETWORK=\"spark-databoost-demo-vpc\"\n", + "export SUBNET=\"spark-databoost-demo-subnet\"\n", + "export SUBNET_CIDR=\"10.10.0.0/24\"\n", + "export ROUTER=\"${NETWORK}-router\"\n", + "export NAT=\"${NETWORK}-nat\"\n", + "\n", + "# Demo resources\n", + "export DEMO_ID=\"spark-databoost-demo\"\n", + "export KAFKA_VM=\"${DEMO_ID}-kafka\"\n", + "export CLUSTER=\"${DEMO_ID}-cluster\"\n", + "export SA_NAME=\"${DEMO_ID}-sa\"\n", + "export SPARK_SA=\"${SA_NAME}@${PROJECT_ID}.iam.gserviceaccount.com\"\n", + "\n", + "# Storage\n", + "export JOB_BUCKET=\"${PROJECT_ID}-${DEMO_ID}-jobs\"\n", + "export WAREHOUSE_BUCKET=\"${PROJECT_ID}-${DEMO_ID}-warehouse\"\n", + "\n", + "# Lakehouse\n", + "export LAKEHOUSE_CATALOG=\"${WAREHOUSE_BUCKET}\"\n", + "export LAKEHOUSE_NAMESPACE=\"demo\"\n", + "\n", + "# Spanner\n", + "export SPANNER_INSTANCE=\"${DEMO_ID}-instance\"\n", + "export SPANNER_DATABASE=\"demo\"\n", + "\n", + "# Kafka\n", + "export KAFKA_TOPIC=\"orders\"\n", + "export KAFKA_BROKER=\"${KAFKA_VM}.${ZONE}.c.${PROJECT_ID}.internal:9092\"\n", + "\n", + "# BigQuery\n", + "export BQ_DATASET=\"spark_databoost_demo_lakehouse_ds\"\n", + "export BQ_RESERVATION=\"${DEMO_ID}-bq-reservation\"" + ] + }, + { + "cell_type": "markdown", + "id": "98832c43", + "metadata": {}, + "source": [ + "Load the variables and configure `gcloud`:" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "b7c2f72e", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud config set project \"$PROJECT_ID\"\n", + "gcloud config set dataproc/region \"$REGION\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "f364299c", + "metadata": {}, + "source": [ + "## 2. Enable APIs" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "b925ee93", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "\n", + "gcloud services enable \\\n", + " compute.googleapis.com \\\n", + " dataproc.googleapis.com \\\n", + " iam.googleapis.com \\\n", + " biglake.googleapis.com \\\n", + " bigquery.googleapis.com \\\n", + " bigqueryreservation.googleapis.com \\\n", + " bigquerydatatransfer.googleapis.com \\\n", + " storage.googleapis.com \\\n", + " spanner.googleapis.com \\\n", + " serviceusage.googleapis.com\n" + ] + }, + { + "cell_type": "markdown", + "id": "6c79441f", + "metadata": {}, + "source": [ + "## 3. Create private network, subnet, NAT, and firewall rules\n", + "\n", + "### 3.1 Create the VPC and subnet\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "b69c1d71", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud compute networks describe \"$NETWORK\" &>/dev/null; then\n", + " echo \"Network $NETWORK already exists, skipping.\"\n", + "else\n", + " gcloud compute networks create \"$NETWORK\" \\\n", + " --subnet-mode=custom\n", + "fi\n", + "\n", + "if gcloud compute networks subnets describe \"$SUBNET\" --region=\"$REGION\" &>/dev/null; then\n", + " echo \"Subnet $SUBNET already exists, skipping.\"\n", + "else\n", + " gcloud compute networks subnets create \"$SUBNET\" \\\n", + " --network=\"$NETWORK\" \\\n", + " --region=\"$REGION\" \\\n", + " --range=\"$SUBNET_CIDR\" \\\n", + " --enable-private-ip-google-access\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "762357a7", + "metadata": {}, + "source": [ + "### 3.2 Allow internal subnet communication" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "d7019b52", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud compute firewall-rules describe \"allow-${NETWORK}-internal\" &>/dev/null; then\n", + " echo \"Firewall rule allow-${NETWORK}-internal already exists, skipping.\"\n", + "else\n", + " gcloud compute firewall-rules create \"allow-${NETWORK}-internal\" \\\n", + " --network=\"$NETWORK\" \\\n", + " --direction=INGRESS \\\n", + " --action=ALLOW \\\n", + " --rules=tcp,udp,icmp \\\n", + " --source-ranges=\"$SUBNET_CIDR\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "c021fa80", + "metadata": {}, + "source": [ + "### 3.3 Allow IAP SSH access" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "21622322", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud compute firewall-rules describe \"allow-${NETWORK}-ssh-iap\" &>/dev/null; then\n", + " echo \"Firewall rule allow-${NETWORK}-ssh-iap already exists, skipping.\"\n", + "else\n", + " gcloud compute firewall-rules create \"allow-${NETWORK}-ssh-iap\" \\\n", + " --network=\"$NETWORK\" \\\n", + " --direction=INGRESS \\\n", + " --action=ALLOW \\\n", + " --rules=tcp:22 \\\n", + " --source-ranges=\"35.235.240.0/20\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "0691c970", + "metadata": {}, + "source": [ + "### 3.4 Create Cloud Router and Cloud NAT" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "32f7924c", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud compute routers describe \"$ROUTER\" --region=\"$REGION\" &>/dev/null; then\n", + " echo \"Router $ROUTER already exists, skipping.\"\n", + "else\n", + " gcloud compute routers create \"$ROUTER\" \\\n", + " --network=\"$NETWORK\" \\\n", + " --region=\"$REGION\"\n", + "fi\n", + "\n", + "if gcloud compute routers nats describe \"$NAT\" --router=\"$ROUTER\" --region=\"$REGION\" &>/dev/null; then\n", + " echo \"NAT $NAT already exists, skipping.\"\n", + "else\n", + " gcloud compute routers nats create \"$NAT\" \\\n", + " --router=\"$ROUTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --auto-allocate-nat-external-ips \\\n", + " --nat-custom-subnet-ip-ranges=\"$SUBNET\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "8b44dd45", + "metadata": {}, + "source": [ + "## 4. Create Service Account and IAM grants\n", + "\n", + "Create the demo service account:\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "7d26e56c", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud iam service-accounts describe \"$SPARK_SA\" &>/dev/null; then\n", + " echo \"Service account $SPARK_SA already exists, skipping.\"\n", + "else\n", + " gcloud iam service-accounts create \"$SA_NAME\" \\\n", + " --display-name=\"Spark Data Boost demo\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "b8c61eca", + "metadata": {}, + "source": [ + "Grant the required project-level roles" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "aef5fe58", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "for ROLE in \\\n", + " roles/dataproc.worker \\\n", + " roles/spanner.databaseUser \\\n", + " roles/spanner.viewer \\\n", + " roles/spanner.databaseReaderWithDataBoost \\\n", + " roles/biglake.editor \\\n", + " roles/bigquery.jobUser \\\n", + " roles/bigquery.dataEditor \\\n", + " roles/serviceusage.serviceUsageConsumer \\\n", + " roles/logging.logWriter; do\n", + " echo \"Granting $ROLE to $SPARK_SA...\"\n", + " gcloud projects add-iam-policy-binding \"$PROJECT_ID\" \\\n", + " --member=\"serviceAccount:${SPARK_SA}\" \\\n", + " --role=\"$ROLE\" \\\n", + " --condition=None \\\n", + " --quiet > /dev/null\n", + "done\n", + "echo \"Done.\"" + ] + }, + { + "cell_type": "markdown", + "id": "c0dc764d", + "metadata": {}, + "source": [ + "Allow the current user to attach the service account:" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "5fd3a874", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "export CURRENT_USER=\"$(gcloud config get-value account)\"\n", + "\n", + "gcloud iam service-accounts add-iam-policy-binding \"$SPARK_SA\" \\\n", + " --member=\"user:${CURRENT_USER}\" \\\n", + " --role=\"roles/iam.serviceAccountUser\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "02004880-a3cf-45d3-9522-bdf7570f9ea5", + "metadata": {}, + "source": [ + "Grant the BigQuery Data Transfer Service agent permission to impersonate `$SPARK_SA`." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "ae348bb5-10c6-48df-b806-9efc759bc5a0", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "PROJECT_NUMBER=$(gcloud projects describe \"$PROJECT_ID\" --format='value(projectNumber)')\n", + "DTS_SA=\"service-${PROJECT_NUMBER}@gcp-sa-bigquerydatatransfer.iam.gserviceaccount.com\"\n", + "\n", + "echo \"Granting Token Creator on $SPARK_SA to $DTS_SA...\"\n", + "gcloud iam service-accounts add-iam-policy-binding \"$SPARK_SA\" \\\n", + " --member=\"serviceAccount:${DTS_SA}\" \\\n", + " --role=\"roles/iam.serviceAccountTokenCreator\" \\\n", + " --condition=None \\\n", + " --quiet >/dev/null" + ] + }, + { + "cell_type": "markdown", + "id": "52608563", + "metadata": {}, + "source": [ + "## 5. Create job/checkpoint and warehouse buckets\n", + "\n", + "### 5.1 Create the job bucket\n", + "\n", + "The job bucket stores Dataproc staging files, PySpark applications, and streaming checkpoints.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "93ab18b2", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud storage buckets describe \"gs://${JOB_BUCKET}\" &>/dev/null; then\n", + " echo \"Bucket gs://${JOB_BUCKET} already exists, skipping creation.\"\n", + "else\n", + " gcloud storage buckets create \"gs://${JOB_BUCKET}\" \\\n", + " --location=\"$REGION\" \\\n", + " --uniform-bucket-level-access\n", + "fi\n", + "\n", + "gcloud storage buckets add-iam-policy-binding \"gs://${JOB_BUCKET}\" \\\n", + " --member=\"serviceAccount:${SPARK_SA}\" \\\n", + " --role=\"roles/storage.objectAdmin\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "22bbd5a9", + "metadata": {}, + "source": [ + "### 5.2 Create the warehouse bucket\n", + "\n", + "The warehouse bucket stores Iceberg table data and metadata.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "a0622cff", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud storage buckets describe \"gs://${WAREHOUSE_BUCKET}\" &>/dev/null; then\n", + " echo \"Bucket gs://${WAREHOUSE_BUCKET} already exists, skipping.\"\n", + "else\n", + " gcloud storage buckets create \"gs://${WAREHOUSE_BUCKET}\" \\\n", + " --location=\"$REGION\" \\\n", + " --uniform-bucket-level-access\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "5ba800ef", + "metadata": {}, + "source": [ + "## 6. Create the Lakehouse Iceberg REST catalog\n", + "\n", + "In the Google Cloud console:\n", + "\n", + "1. Open **BigQuery / Lakehouse**.\n", + "2. Select **Create catalog**.\n", + "3. Select **Iceberg REST catalog**.\n", + "4. Select **Single bucket catalog**.\n", + "5. Select `gs://${WAREHOUSE_BUCKET}` as the Cloud Storage bucket.\n", + "6. Select **Credential vending** as the authentication method.\n", + "7. Create the catalog.\n", + "8. Select **Set bucket permissions**.\n", + "\n", + "The final action gives the auto-provisioned Lakehouse catalog service account access to the warehouse bucket. Spark receives short-lived, scoped storage credentials through the REST catalog.\n" + ] + }, + { + "cell_type": "markdown", + "id": "6087d348", + "metadata": {}, + "source": [ + "## 7. Create Kafka VM and start Kafka\n", + "\n", + "### 7.1 Create the private Kafka VM\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "876f9287", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud compute instances describe \"$KAFKA_VM\" --zone=\"$ZONE\" &>/dev/null; then\n", + " echo \"Instance $KAFKA_VM already exists, skipping.\"\n", + "else\n", + " gcloud compute instances create \"$KAFKA_VM\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --machine-type=\"e2-small\" \\\n", + " --network=\"$NETWORK\" \\\n", + " --subnet=\"$SUBNET\" \\\n", + " --no-address \\\n", + " --image-family=\"cos-stable\" \\\n", + " --image-project=\"cos-cloud\" \\\n", + " --boot-disk-size=\"10GB\" \\\n", + " --boot-disk-type=\"pd-standard\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "1f9d6643", + "metadata": {}, + "source": [ + "### 7.2 Copy the source files to the VM" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "45934f57", + "metadata": { + "scrolled": true + }, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud compute scp \\\n", + " --recurse \\\n", + " src \\\n", + " .env \\\n", + " \"${KAFKA_VM}:~/\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --tunnel-through-iap\n" + ] + }, + { + "cell_type": "markdown", + "id": "d4bec1cc", + "metadata": {}, + "source": [ + "### 7.3 Start Kafka and build the producer image\n", + "\n", + "This step starts the Kafka broker and builds the data generator Docker image from source. \n", + "Expect this cell to take 3 to 5 minutes to complete." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "bc44976f", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud compute ssh \"$KAFKA_VM\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --tunnel-through-iap \\\n", + " --command='\n", + " set -euo pipefail\n", + "\n", + " source ~/.env\n", + "\n", + " : \"${KAFKA_BROKER:?KAFKA_BROKER is not defined in ~/.env}\"\n", + " : \"${KAFKA_TOPIC:?KAFKA_TOPIC is not defined in ~/.env}\"\n", + "\n", + " sudo iptables -C INPUT -p tcp --dport 9092 -j ACCEPT \\\n", + " 2>/dev/null || \\\n", + " sudo iptables -A INPUT -p tcp --dport 9092 -j ACCEPT\n", + "\n", + " docker rm -f kafka 2>/dev/null || true\n", + "\n", + " docker run -d \\\n", + " --name kafka \\\n", + " --network host \\\n", + " -e KAFKA_NODE_ID=1 \\\n", + " -e KAFKA_PROCESS_ROLES=broker,controller \\\n", + " -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \\\n", + " -e KAFKA_LISTENERS=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093 \\\n", + " -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://${KAFKA_BROKER} \\\n", + " -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \\\n", + " -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \\\n", + " -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \\\n", + " -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \\\n", + " -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \\\n", + " -e KAFKA_DEFAULT_REPLICATION_FACTOR=1 \\\n", + " apache/kafka-native:4.1.2\n", + "\n", + " echo \"Waiting for Kafka and creating topic: $KAFKA_TOPIC\"\n", + "\n", + " for attempt in $(seq 1 10); do\n", + " echo \"Attempt $attempt of 10 in 30 seconds...\"\n", + " sleep 30\n", + "\n", + " if docker run --rm \\\n", + " --network host \\\n", + " apache/kafka:4.1.2 \\\n", + " /opt/kafka/bin/kafka-topics.sh \\\n", + " --bootstrap-server localhost:9092 \\\n", + " --create \\\n", + " --if-not-exists \\\n", + " --topic \"$KAFKA_TOPIC\" \\\n", + " --partitions 1 \\\n", + " --replication-factor 1; then\n", + " echo \"Topic is ready: $KAFKA_TOPIC\"\n", + " break\n", + " fi\n", + " if [ \"$attempt\" -eq 10 ]; then\n", + " echo \"Failed to create topic after 10 attempts\" >&2\n", + " docker logs kafka >&2\n", + " exit 1\n", + " fi\n", + " done\n", + "\n", + " cd ~/src/java-tpch-stream-generator\n", + " docker build --tag tpch-generator:local .\n", + " '\n" + ] + }, + { + "cell_type": "markdown", + "id": "88d40a93", + "metadata": {}, + "source": [ + "## 8. Create Spanner instance and schema\n", + "\n", + "### 8.1 Create the Spanner instance\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "66d0e166", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud spanner instances describe \"$SPANNER_INSTANCE\" &>/dev/null; then\n", + " echo \"Spanner instance $SPANNER_INSTANCE already exists, skipping.\"\n", + "else\n", + " gcloud spanner instances create \"$SPANNER_INSTANCE\" \\\n", + " --config=\"regional-${REGION}\" \\\n", + " --description=\"Spanner Data Boost demo\" \\\n", + " --edition=\"ENTERPRISE\" \\\n", + " --processing-units=100\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "82f88f84", + "metadata": {}, + "source": [ + "### 8.2 Create the database and tables" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9e35121b", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "if gcloud spanner databases describe \"$SPANNER_DATABASE\" --instance=\"$SPANNER_INSTANCE\" &>/dev/null; then\n", + " echo \"Spanner database $SPANNER_DATABASE already exists.\"\n", + "else\n", + " gcloud spanner databases create \"$SPANNER_DATABASE\" \\\n", + " --instance=\"$SPANNER_INSTANCE\"\n", + "fi\n", + "\n", + "gcloud spanner databases ddl update \"$SPANNER_DATABASE\" \\\n", + " --instance=\"$SPANNER_INSTANCE\" \\\n", + " --ddl=\"CREATE TABLE IF NOT EXISTS orders (\n", + " order_key INT64 NOT NULL,\n", + " customer_key INT64,\n", + " order_status STRING(1) NOT NULL,\n", + " total_price NUMERIC,\n", + " order_date DATE NOT NULL,\n", + " order_priority STRING(15),\n", + " clerk STRING(15) NOT NULL,\n", + " ship_priority INT64,\n", + " comment STRING(79),\n", + " kafka_timestamp TIMESTAMP,\n", + " ingested_time TIMESTAMP\n", + " ) PRIMARY KEY (order_key);\n", + "\n", + " CREATE TABLE IF NOT EXISTS order_daily_aggregate (\n", + " order_date DATE NOT NULL,\n", + " order_count INT64 NOT NULL,\n", + " order_total NUMERIC NOT NULL,\n", + " order_total_average NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP\n", + " ) PRIMARY KEY (order_date);\n", + "\n", + " CREATE TABLE IF NOT EXISTS clerk_status_daily_aggregate (\n", + " clerk STRING(15) NOT NULL,\n", + " order_status STRING(1) NOT NULL,\n", + " order_date DATE NOT NULL,\n", + " orders_processed INT64 NOT NULL,\n", + " total_revenue NUMERIC NOT NULL,\n", + " avg_order_value NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP\n", + " ) PRIMARY KEY (clerk, order_status, order_date);\n", + "\n", + " CREATE TABLE IF NOT EXISTS order_status_daily_aggregate (\n", + " order_status STRING(1) NOT NULL,\n", + " order_date DATE NOT NULL,\n", + " orders_processed INT64 NOT NULL,\n", + " total_revenue NUMERIC NOT NULL,\n", + " avg_order_value NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP\n", + " ) PRIMARY KEY (order_status, order_date);\"" + ] + }, + { + "cell_type": "markdown", + "id": "c9c0ea2b", + "metadata": {}, + "source": [ + "### 8.3 Enable columnar storage for Orders table" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9f4952ff", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud spanner databases ddl update $SPANNER_DATABASE \\\n", + " --instance=$SPANNER_INSTANCE \\\n", + " --ddl=\"ALTER TABLE orders SET OPTIONS (columnar_policy = 'enabled');\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "64494a4d", + "metadata": {}, + "source": [ + "## 9. Create a single-node Managed Spark cluster\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "26e25066", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "MANAGED_SPARK_IMAGE=\"2.3-debian12\"\n", + "\n", + "if gcloud dataproc clusters describe \"$CLUSTER\" --region=\"$REGION\" &>/dev/null; then\n", + " echo \"Dataproc cluster $CLUSTER already exists, skipping.\"\n", + "else\n", + " gcloud dataproc clusters create \"$CLUSTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --subnet=\"$SUBNET\" \\\n", + " --no-address \\\n", + " --single-node \\\n", + " --image-version=\"$MANAGED_SPARK_IMAGE\" \\\n", + " --master-machine-type=n2d-standard-8 \\\n", + " --master-boot-disk-type=pd-ssd \\\n", + " --master-boot-disk-size=\"100GB\" \\\n", + " --bucket=\"$JOB_BUCKET\" \\\n", + " --temp-bucket=\"$JOB_BUCKET\" \\\n", + " --service-account=\"$SPARK_SA\" \\\n", + " --optional-components=JUPYTER \\\n", + " --enable-component-gateway \\\n", + " --scopes=\"cloud-platform\"\n", + "fi\n" + ] + }, + { + "cell_type": "markdown", + "id": "64660f1f", + "metadata": {}, + "source": [ + "Expect the cluster to be created in a few minutes." + ] + }, + { + "cell_type": "markdown", + "id": "29448c90", + "metadata": {}, + "source": [ + "## 10. Upload and submit both streaming jobs\n", + "\n", + "**Streaming flow:** The 2 Spark jobs process data in micro-batches governed by:\n", + "\n", + "- **processing-time:** The fixed clock interval between micro-batch triggers. When each interval arrives, Spark immediately polls Kafka and processes whatever records are currently available at that exact moment.\n", + "\n", + "- **max-offsets-per-trigger:** A hard upper limit (ceiling) on the number of Kafka records pulled in a single micro-batch. If Kafka has a massive backlog, Spark caps the batch size at this number to prevent memory overload.\n", + "\n", + "### 10.1 Write the PySpark application sources\n", + "\n", + "The two streaming job sources live in these cells so the code and its\n", + "submission command stay together. Running the cell writes the file to\n", + "`src/` for upload.\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "8f007d59", + "metadata": {}, + "outputs": [], + "source": [ + "%%writefile src/kafka_to_lakehouse.py\n", + "#!/usr/bin/env python3\n", + "\n", + "# Requires Jars in Spark env.\n", + "# - spark-iceberg connector\n", + "# - spark-sql-kafka connector\n", + "import argparse\n", + "from pyspark.sql import SparkSession, DataFrame\n", + "from pyspark.sql.functions import col, split, to_date, current_timestamp\n", + "from pyspark.sql.types import DecimalType, LongType, IntegerType\n", + "\n", + "\n", + "def parse_args():\n", + " \"\"\"Parses command-line arguments using argparse.\"\"\"\n", + " parser = argparse.ArgumentParser(\n", + " description=\"Stream data from Kafka into an Apache Iceberg lakehouse table.\"\n", + " )\n", + "\n", + " parser.add_argument(\"--bootstrap\", required=True, help=\"Kafka bootstrap servers connection string\")\n", + " parser.add_argument(\"--kafka-topic\", required=True, help=\"Kafka topic name to subscribe to\")\n", + " parser.add_argument(\"--lakehouse-catalog\", required=True, help=\"Target Iceberg catalog name\")\n", + " parser.add_argument(\"--lakehouse-namespace\", required=True, help=\"Target Iceberg namespace/database\")\n", + " parser.add_argument(\"--lakehouse-table\", required=True, help=\"Target Iceberg table name\")\n", + " parser.add_argument(\"--spark-stream-checkpoint\", required=True, help=\"Remote or local for streaming checkpoints\")\n", + " parser.add_argument(\n", + " \"--processing-time\",\n", + " default=\"30 seconds\",\n", + " help=\"Processing time interval for stream trigger (e.g., '10 seconds', '1 minute', '0 seconds'). Default: '30 seconds'\"\n", + " )\n", + " parser.add_argument(\n", + " \"--max-offsets-per-trigger\",\n", + " default=\"100\",\n", + " help=\"Maximum number of Kafka offsets processed per trigger cycle. Default: 100\"\n", + " )\n", + "\n", + " return parser.parse_args()\n", + "\n", + "\n", + "def read_orders(spark: SparkSession, bootstrap_servers: str, topic: str, max_offsets: str) -> DataFrame:\n", + " reader = (\n", + " spark.readStream.format(\"kafka\")\n", + " .option(\"kafka.bootstrap.servers\", bootstrap_servers)\n", + " .option(\"subscribe\", topic)\n", + " .option(\"startingOffsets\", \"earliest\")\n", + " .option(\"maxOffsetsPerTrigger\", max_offsets)\n", + " .option(\"failOnDataLoss\", \"true\")\n", + " )\n", + "\n", + " # Extract raw value and Kafka record timestamp\n", + " payload = reader.load().select(\n", + " col(\"value\").cast(\"string\").alias(\"value\"),\n", + " col(\"timestamp\").alias(\"kafka_timestamp\")\n", + " )\n", + " fields = split(col(\"value\"), r\"\\|\")\n", + " return payload.select(\n", + " fields.getItem(0).cast(LongType()).alias(\"order_key\"),\n", + " fields.getItem(1).cast(LongType()).alias(\"customer_key\"),\n", + " fields.getItem(2).alias(\"order_status\"),\n", + " fields.getItem(3).cast(DecimalType(15, 2)).alias(\"total_price\"),\n", + " to_date(fields.getItem(4), \"yyyy-MM-dd\").alias(\"order_date\"),\n", + " fields.getItem(5).alias(\"order_priority\"),\n", + " fields.getItem(6).alias(\"clerk\"),\n", + " fields.getItem(7).cast(IntegerType()).alias(\"ship_priority\"),\n", + " fields.getItem(8).alias(\"comment\"),\n", + " col(\"kafka_timestamp\"),\n", + " current_timestamp().alias(\"ingested_time\")\n", + " )\n", + "\n", + "\n", + "def main():\n", + " # Parse arguments\n", + " args = parse_args()\n", + "\n", + " spark = (\n", + " SparkSession.builder\n", + " .appName(\"kafka-to-lakehouse-orders\")\n", + " .config(\"spark.sql.iceberg.check-nullability\", \"false\")\n", + " .getOrCreate()\n", + " )\n", + " spark.sparkContext.setLogLevel(\"WARN\")\n", + "\n", + " # Setup Lakehouse\n", + " namespace_name = (\n", + " f\"`{args.lakehouse_catalog}`.\"\n", + " f\"`{args.lakehouse_namespace}`\"\n", + " )\n", + "\n", + " table_name = (\n", + " f\"`{args.lakehouse_catalog}`.\"\n", + " f\"`{args.lakehouse_namespace}`.\"\n", + " f\"`{args.lakehouse_table}`\"\n", + " )\n", + "\n", + " spark.sql(\n", + " f\"CREATE NAMESPACE IF NOT EXISTS {namespace_name}\"\n", + " )\n", + "\n", + " spark.sql(\n", + " f\"\"\"\n", + " CREATE TABLE IF NOT EXISTS {table_name} (\n", + " order_key BIGINT,\n", + " customer_key BIGINT,\n", + " order_status STRING NOT NULL,\n", + " total_price DECIMAL(15,2),\n", + " order_date DATE NOT NULL,\n", + " order_priority STRING,\n", + " clerk STRING NOT NULL,\n", + " ship_priority INT,\n", + " comment STRING,\n", + " kafka_timestamp TIMESTAMP,\n", + " ingested_time TIMESTAMP\n", + " )\n", + " USING iceberg\n", + " PARTITIONED BY (hours(ingested_time))\n", + " \"\"\"\n", + " )\n", + "\n", + " # read orders from Kafka\n", + " orders = read_orders(\n", + " spark,\n", + " args.bootstrap,\n", + " args.kafka_topic,\n", + " args.max_offsets_per_trigger\n", + " )\n", + "\n", + " # stream order to Lakehouse catalog\n", + " query = (orders.writeStream.format(\"iceberg\").outputMode(\"append\")\n", + " .option(\"checkpointLocation\", args.spark_stream_checkpoint)\n", + " .trigger(processingTime=args.processing_time)\n", + " .queryName(\"kafka-to-lakehouse-orders\").toTable(table_name))\n", + " query.awaitTermination()\n", + " spark.stop()\n", + "\n", + "\n", + "if __name__ == \"__main__\":\n", + " main()\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "25d93259", + "metadata": {}, + "outputs": [], + "source": [ + "%%writefile src/kafka_to_spanner.py\n", + "#!/usr/bin/env python3\n", + "\n", + "# Requires Jars in Spark env.\n", + "# - spark-cloud-spanner connector\n", + "# - spark-sql-kafka connector\n", + "import argparse\n", + "from pyspark.sql import DataFrame, SparkSession\n", + "from pyspark.sql.functions import col, split, to_date, current_timestamp\n", + "from pyspark.sql.types import DecimalType, LongType, IntegerType\n", + "\n", + "\n", + "def parse_args():\n", + " \"\"\"Parses command-line arguments using argparse.\"\"\"\n", + " parser = argparse.ArgumentParser(\n", + " description=\"Stream data from Kafka into a Google Cloud Spanner table.\"\n", + " )\n", + "\n", + " parser.add_argument(\"--bootstrap\", required=True, help=\"Kafka bootstrap servers connection string\")\n", + " parser.add_argument(\"--kafka-topic\", required=True, help=\"Kafka topic name to subscribe to\")\n", + " parser.add_argument(\"--spanner-project\", required=True, help=\"Google Cloud Project ID\")\n", + " parser.add_argument(\"--spanner-instance\", required=True, help=\"Cloud Spanner instance ID\")\n", + " parser.add_argument(\"--spanner-database\", required=True, help=\"Cloud Spanner database name\")\n", + " parser.add_argument(\"--spanner-table\", required=True, help=\"Cloud Spanner table name\")\n", + " parser.add_argument(\"--spark-stream-checkpoint\", required=True,\n", + " help=\"Remote or local path for streaming checkpoints\")\n", + " parser.add_argument(\n", + " \"--processing-time\",\n", + " default=\"5 seconds\",\n", + " help=\"Processing time interval for stream trigger (e.g., '10 seconds', '1 minute', '0 seconds'). Default: '5 seconds'\"\n", + " )\n", + " parser.add_argument(\n", + " \"--max-offsets-per-trigger\",\n", + " default=\"100\",\n", + " help=\"Maximum number of Kafka offsets processed per trigger cycle. Default: 100\"\n", + " )\n", + "\n", + " return parser.parse_args()\n", + "\n", + "\n", + "def read_orders(spark: SparkSession, bootstrap_servers: str, topic: str, max_offsets_per_trigger: str) -> DataFrame:\n", + " reader = (\n", + " spark.readStream.format(\"kafka\")\n", + " .option(\"kafka.bootstrap.servers\", bootstrap_servers)\n", + " .option(\"subscribe\", topic)\n", + " .option(\"startingOffsets\", \"earliest\")\n", + " .option(\"maxOffsetsPerTrigger\", max_offsets_per_trigger)\n", + " .option(\"failOnDataLoss\", \"true\")\n", + " )\n", + "\n", + " # Extract raw value and Kafka record timestamp\n", + " payload = reader.load().select(\n", + " col(\"value\").cast(\"string\").alias(\"value\"),\n", + " col(\"timestamp\").alias(\"kafka_timestamp\")\n", + " )\n", + " fields = split(col(\"value\"), r\"\\|\")\n", + " return payload.select(\n", + " fields.getItem(0).cast(LongType()).alias(\"order_key\"),\n", + " fields.getItem(1).cast(LongType()).alias(\"customer_key\"),\n", + " fields.getItem(2).alias(\"order_status\"),\n", + " fields.getItem(3).cast(DecimalType(15, 2)).alias(\"total_price\"),\n", + " to_date(fields.getItem(4), \"yyyy-MM-dd\").alias(\"order_date\"),\n", + " fields.getItem(5).alias(\"order_priority\"),\n", + " fields.getItem(6).alias(\"clerk\"),\n", + " fields.getItem(7).cast(IntegerType()).alias(\"ship_priority\"),\n", + " fields.getItem(8).alias(\"comment\"),\n", + " col(\"kafka_timestamp\"),\n", + " current_timestamp().alias(\"ingested_time\")\n", + " )\n", + "\n", + "\n", + "def write_batch(batch: DataFrame, project: str, instance: str, database: str, table: str):\n", + " if batch.isEmpty():\n", + " return\n", + " (batch.write.format(\"cloud-spanner\")\n", + " .option(\"projectId\", project).option(\"instanceId\", instance)\n", + " .option(\"databaseId\", database).option(\"table\", table)\n", + " .option(\"mutationType\", \"insert_or_update\")\n", + " .option(\"assumeIdempotentRows\", \"true\").mode(\"append\").save())\n", + "\n", + "\n", + "def main():\n", + " # Parse arguments\n", + " args = parse_args()\n", + "\n", + " spark = SparkSession.builder.appName(\"kafka-to-spanner-orders\").getOrCreate()\n", + " spark.sparkContext.setLogLevel(\"WARN\")\n", + "\n", + " # read orders from Kafka\n", + " orders = read_orders(\n", + " spark,\n", + " args.bootstrap,\n", + " args.kafka_topic,\n", + " args.max_offsets_per_trigger\n", + " )\n", + "\n", + " # Stream orders to Spanner using mapped underscore properties\n", + " query = (orders.writeStream.foreachBatch(\n", + " lambda b, _: write_batch(\n", + " b,\n", + " args.spanner_project,\n", + " args.spanner_instance,\n", + " args.spanner_database,\n", + " args.spanner_table\n", + " )\n", + " )\n", + " .option(\"checkpointLocation\", args.spark_stream_checkpoint)\n", + " .trigger(processingTime=args.processing_time)\n", + " .queryName(\"kafka-to-spanner-orders\").start())\n", + " query.awaitTermination()\n", + " spark.stop()\n", + "\n", + "\n", + "if __name__ == \"__main__\":\n", + " main()\n" + ] + }, + { + "cell_type": "markdown", + "id": "7db3e5b8", + "metadata": {}, + "source": [ + "### 10.2 Upload the PySpark applications" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "9298ece6", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud storage cp \\\n", + " ./src/kafka_to_lakehouse.py \\\n", + " ./src/kafka_to_spanner.py \\\n", + " \"gs://${JOB_BUCKET}/jobs/\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "22bd9223", + "metadata": {}, + "source": [ + "### 10.3 Submit the Kafka-to-Lakehouse job" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "e1d4e6e7", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "KAFKA_PKG=\"org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3\"\n", + "\n", + "LAKEHOUSE_PROPS=\"^#^\\\n", + "spark.jars.ivy=/tmp/ivy-lakehouse#\\\n", + "spark.jars.packages=${KAFKA_PKG}#\\\n", + "spark.driver.memory=2g#\\\n", + "spark.sql.defaultCatalog=${LAKEHOUSE_CATALOG}#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}=org.apache.iceberg.spark.SparkCatalog#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.type=rest#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.uri=https://biglake.googleapis.com/iceberg/v1/restcatalog#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.warehouse=gs://${WAREHOUSE_BUCKET}#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.io-impl=org.apache.iceberg.gcp.gcs.GCSFileIO#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.header.x-goog-user-project=${PROJECT_ID}#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.rest.auth.type=google#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.header.X-Iceberg-Access-Delegation=vended-credentials#\\\n", + "spark.sql.catalog.${LAKEHOUSE_CATALOG}.rest-metrics-reporting-enabled=false#\\\n", + "spark.sql.extensions=org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions\"\n", + "\n", + "LAKEHOUSE_JOB_ID=\"${DEMO_ID}-lakehouse-$(date +%Y%m%d%H%M%S)\"\n", + "\n", + "ICEBERG_VERSION=\"1.10.0\"\n", + "\n", + "gcloud dataproc jobs submit pyspark \\\n", + " \"gs://${JOB_BUCKET}/jobs/kafka_to_lakehouse.py\" \\\n", + " --id=\"$LAKEHOUSE_JOB_ID\" \\\n", + " --cluster=\"$CLUSTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --properties=\"$LAKEHOUSE_PROPS\" \\\n", + " --jars=https://storage-download.googleapis.com/maven-central/maven2/org/apache/iceberg/iceberg-spark-runtime-3.5_2.12/${ICEBERG_VERSION}/iceberg-spark-runtime-3.5_2.12-${ICEBERG_VERSION}.jar,https://storage-download.googleapis.com/maven-central/maven2/org/apache/iceberg/iceberg-gcp-bundle/${ICEBERG_VERSION}/iceberg-gcp-bundle-${ICEBERG_VERSION}.jar \\\n", + " --async \\\n", + " -- \\\n", + " --bootstrap=\"$KAFKA_BROKER\" \\\n", + " --kafka-topic=\"$KAFKA_TOPIC\" \\\n", + " --lakehouse-catalog=\"$LAKEHOUSE_CATALOG\" \\\n", + " --lakehouse-namespace=\"$LAKEHOUSE_NAMESPACE\" \\\n", + " --lakehouse-table=\"orders\" \\\n", + " --spark-stream-checkpoint=\"gs://${JOB_BUCKET}/checkpoints/lakehouse\" \\\n", + " --processing-time=\"10 seconds\" \\\n", + " --max-offsets-per-trigger=\"100\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "70709f9f", + "metadata": {}, + "source": [ + "### 10.4 Submit the Kafka-to-Spanner job" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "6d41315c", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "KAFKA_PKG=\"org.apache.spark:spark-sql-kafka-0-10_2.12:3.5.3\"\n", + "SPANNER_PKG=\"com.google.cloud.spark.spanner:spark-3.5-spanner:1.5.0\"\n", + "\n", + "SPANNER_PROPS=\"^#^\\\n", + "spark.jars.ivy=/tmp/ivy-spanner#\\\n", + "spark.jars.packages=${KAFKA_PKG},${SPANNER_PKG}#\\\n", + "spark.master=local[2]#\\\n", + "spark.driver.memory=2g\"\n", + "\n", + "SPANNER_JOB_ID=\"${DEMO_ID}-spanner-$(date +%Y%m%d%H%M%S)\"\n", + "\n", + "gcloud dataproc jobs submit pyspark \\\n", + " \"gs://${JOB_BUCKET}/jobs/kafka_to_spanner.py\" \\\n", + " --id=\"$SPANNER_JOB_ID\" \\\n", + " --cluster=\"$CLUSTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --properties=\"$SPANNER_PROPS\" \\\n", + " --async \\\n", + " -- \\\n", + " --bootstrap=\"$KAFKA_BROKER\" \\\n", + " --kafka-topic=\"$KAFKA_TOPIC\" \\\n", + " --spanner-project=\"$PROJECT_ID\" \\\n", + " --spanner-instance=\"$SPANNER_INSTANCE\" \\\n", + " --spanner-database=\"$SPANNER_DATABASE\" \\\n", + " --spanner-table=\"orders\" \\\n", + " --spark-stream-checkpoint=\"gs://${JOB_BUCKET}/checkpoints/spanner\" \\\n", + " --processing-time=\"5 seconds\" \\\n", + " --max-offsets-per-trigger=\"100\"\n" + ] + }, + { + "cell_type": "markdown", + "id": "1686a6773c7803c7", + "metadata": {}, + "source": [ + "### 10.5 View the jobs in Google Cloud console\n", + "\n", + "In console, go to Managed Spark > Cluster > Jobs to check if the jobs are running.\n", + "The jobs would take 1-2 minutes to initialize and be ready to process data." + ] + }, + { + "cell_type": "markdown", + "id": "d29f0935", + "metadata": {}, + "source": [ + "## 11. Compute aggregates from Lakehouse to Spanner using BigQuery Scheduled & Continuous Queries\n", + "\n", + "This section implements the automated end-to-end aggregation and streaming pipeline into Spanner:\n", + "\n", + "- **5-minute scheduled query (batch precomputation):** scans the Iceberg Lakehouse orders table for affected keys, computes historical `COUNT` and `SUM` snapshots, and appends them to three native BigQuery staging tables (`stg_order_daily`, `stg_clerk_status_daily`, and `stg_order_status_daily`).\n", + "- **Three BigQuery Continuous Queries (reverse ETL):** consume changes from the staging tables through `APPENDS()`, calculate averages with `SAFE_DIVIDE`, and export directly into the corresponding Spanner aggregate tables.\n", + "\n", + "The native staging layer is required because BigQuery Continuous Queries cannot read external Iceberg tables directly. https://docs.cloud.google.com/bigquery/docs/continuous-queries-introduction#limitations" + ] + }, + { + "cell_type": "markdown", + "id": "bq_11_5_md", + "metadata": {}, + "source": [ + "### 11.1 Create BigQuery Enterprise reservation and assignment\n", + "\n", + "Continuous Queries require an Enterprise or Enterprise Plus reservation with a `CONTINUOUS` assignment. This demo creates a dedicated 50-slot reservation for the three jobs." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "bq_11_5_code", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "echo \"=== Setting up BigQuery Enterprise Reservation for Continuous Queries ===\"\n", + "\n", + "# 1. Create or resize the Enterprise reservation\n", + "if ! bq show --reservation --location=\"$REGION\" --project_id=\"$PROJECT_ID\" \"$BQ_RESERVATION\" &>/dev/null; then\n", + " echo \"Creating Enterprise reservation: $BQ_RESERVATION...\"\n", + " bq mk --reservation \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --slots=50 \\\n", + " --edition=ENTERPRISE \\\n", + " \"$BQ_RESERVATION\"\n", + "else\n", + " echo \"Updating reservation $BQ_RESERVATION to 50 slots...\"\n", + " bq update --reservation \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --slots=50 \\\n", + " \"$BQ_RESERVATION\"\n", + "fi\n", + "\n", + "# 2. Create the CONTINUOUS reservation assignment\n", + "RAW_ASSIGNMENTS=$(\n", + " bq ls \\\n", + " --reservation_assignment \\\n", + " --location=\"$REGION\" \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --format=json 2>/dev/null || true\n", + ")\n", + "\n", + "EXISTING_ASSIGNMENT=\"\"\n", + "# Only parse with jq if bq returned a JSON array\n", + "if [[ \"$RAW_ASSIGNMENTS\" =~ ^[[:space:]]*\\[ ]]; then\n", + " EXISTING_ASSIGNMENT=$(\n", + " echo \"$RAW_ASSIGNMENTS\" | jq -r --arg assignee \"projects/${PROJECT_ID}\" '\n", + " .[]\n", + " | select(.jobType == \"CONTINUOUS\" and .assignee == $assignee)\n", + " | .name // empty\n", + " '\n", + " )\n", + "fi\n", + "\n", + "if [ -z \"$EXISTING_ASSIGNMENT\" ]; then\n", + " echo \"Creating reservation assignment for CONTINUOUS jobs...\"\n", + " bq mk \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --reservation_assignment \\\n", + " --reservation_id=\"$BQ_RESERVATION\" \\\n", + " --job_type=CONTINUOUS \\\n", + " --assignee_type=PROJECT \\\n", + " --assignee_id=\"$PROJECT_ID\"\n", + "elif [[ \"$EXISTING_ASSIGNMENT\" == *\"/reservations/${BQ_RESERVATION}/assignments/\"* ]]; then\n", + " echo \"CONTINUOUS reservation assignment already active: $EXISTING_ASSIGNMENT\"\n", + "else\n", + " echo \"Project $PROJECT_ID already has a CONTINUOUS assignment to another reservation:\" >&2\n", + " echo \"$EXISTING_ASSIGNMENT\" >&2\n", + " exit 1\n", + "fi\n", + "\n", + "echo \"Enterprise reservation configured for Continuous Queries.\"" + ] + }, + { + "cell_type": "markdown", + "id": "25e56734", + "metadata": {}, + "source": [ + "### 11.2 Create native BigQuery staging tables\n", + "\n", + "Create the BigQuery dataset and three native staging tables (`stg_order_daily`, `stg_clerk_status_daily`, and `stg_order_status_daily`) that hold precomputed counts and totals." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "1297dcf6", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "# 1. Create BigQuery dataset if not exists\n", + "if ! bq show --dataset \"${PROJECT_ID}:${BQ_DATASET}\" &>/dev/null; then\n", + " echo \"Creating BigQuery dataset: ${BQ_DATASET}...\"\n", + " bq mk --dataset --location=\"$REGION\" \"${PROJECT_ID}:${BQ_DATASET}\"\n", + "else\n", + " echo \"BigQuery dataset ${BQ_DATASET} already exists.\"\n", + "fi\n", + "\n", + "# 2. Create the 3 native staging tables\n", + "echo \"Creating staging tables in dataset: ${BQ_DATASET}...\"\n", + "bq query --project_id=\"$PROJECT_ID\" --location=\"$REGION\" --use_legacy_sql=false \"\n", + "-- 1. Daily Staging Table\n", + "CREATE TABLE IF NOT EXISTS \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_daily\\` (\n", + " order_date DATE NOT NULL,\n", + " order_count INT64 NOT NULL,\n", + " order_total NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP NOT NULL\n", + ")\n", + "PARTITION BY DATE(last_updated)\n", + "OPTIONS (\n", + " partition_expiration_days = 2\n", + ");\n", + "\n", + "-- 2. Clerk & Status Staging Table\n", + "CREATE TABLE IF NOT EXISTS \\`${PROJECT_ID}.${BQ_DATASET}.stg_clerk_status_daily\\` (\n", + " clerk STRING NOT NULL,\n", + " order_status STRING NOT NULL,\n", + " order_date DATE NOT NULL,\n", + " orders_processed INT64 NOT NULL,\n", + " total_revenue NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP NOT NULL\n", + ")\n", + "PARTITION BY DATE(last_updated)\n", + "CLUSTER BY clerk, order_status\n", + "OPTIONS (\n", + " partition_expiration_days = 2\n", + ");\n", + "\n", + "-- 3. Order Status Staging Table\n", + "CREATE TABLE IF NOT EXISTS \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_status_daily\\` (\n", + " order_status STRING NOT NULL,\n", + " order_date DATE NOT NULL,\n", + " orders_processed INT64 NOT NULL,\n", + " total_revenue NUMERIC NOT NULL,\n", + " last_updated TIMESTAMP NOT NULL\n", + ")\n", + "PARTITION BY DATE(last_updated)\n", + "CLUSTER BY order_status\n", + "OPTIONS (\n", + " partition_expiration_days = 2\n", + ");\n", + "\"\n", + "echo \"Staging tables ready.\"" + ] + }, + { + "cell_type": "markdown", + "id": "03ae5945-3352-4309-9a9a-72d2df6263a4", + "metadata": {}, + "source": [ + "### 11.3 Create the precompute scheduled query\n", + "\n", + "The query uses a 5-minute lookback to identify recently affected keys, then recomputes complete `COUNT` and `SUM` snapshots for those keys." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "f602866f-2f45-4857-b85d-7a5c3f21c094", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "mkdir -p src/bq_scheduled\n", + "\n", + "cat < src/bq_scheduled/precompute_aggregates_to_staging.sql\n", + "-- =========================================================================\n", + "-- Step 1: Identify dates affected in the last 5-minute schedule interval\n", + "-- and pull their historical data for complete cumulative snapshots.\n", + "-- =========================================================================\n", + "CREATE TEMP TABLE active_orders AS\n", + "WITH recent_dates AS (\n", + " SELECT order_date\n", + " FROM \\`${PROJECT_ID}.${LAKEHOUSE_CATALOG}.${LAKEHOUSE_NAMESPACE}.orders\\`\n", + " WHERE ingested_time > TIMESTAMP_SUB(@run_time, INTERVAL 5 MINUTE)\n", + " AND ingested_time <= @run_time\n", + ")\n", + "SELECT\n", + " order_date,\n", + " clerk,\n", + " order_status,\n", + " total_price,\n", + " (ingested_time > TIMESTAMP_SUB(@run_time, INTERVAL 5 MINUTE)\n", + " AND ingested_time <= @run_time) AS is_recent\n", + "FROM \\`${PROJECT_ID}.${LAKEHOUSE_CATALOG}.${LAKEHOUSE_NAMESPACE}.orders\\`\n", + "WHERE order_date IN (SELECT order_date FROM recent_dates);\n", + "\n", + "-- =========================================================================\n", + "-- Step 2: Compute all 3 aggregation levels simultaneously.\n", + "--\n", + "-- Use GROUPING SETS for:\n", + "-- 1. Efficiency: Scans and aggregates 'active_orders' once.\n", + "-- 2. Safe Routing: GROUPING() reliably identifies which aggregation level each row belongs to.\n", + "-- =========================================================================\n", + "CREATE TEMP TABLE aggregate_snapshots AS\n", + "SELECT\n", + " order_date,\n", + " clerk,\n", + " order_status,\n", + " COUNT(*) AS orders_count,\n", + " CAST(COALESCE(SUM(total_price), 0) AS NUMERIC) AS total_value,\n", + " GROUPING(clerk) AS is_clerk_aggregated,\n", + " GROUPING(order_status) AS is_status_aggregated,\n", + " @run_time AS last_updated\n", + "FROM active_orders\n", + "GROUP BY GROUPING SETS (\n", + " (order_date), -- Granularity 1: Daily total\n", + " (order_status, order_date), -- Granularity 2: Order status by date\n", + " (clerk, order_status, order_date) -- Granularity 3: Clerk & status by date\n", + ")\n", + "HAVING LOGICAL_OR(is_recent);\n", + "\n", + "-- =========================================================================\n", + "-- Step 3: Route into your 3 existing staging tables using is_***_aggregated \n", + "-- conditions.\n", + "-- =========================================================================\n", + "-- 1. Daily Aggregates\n", + "INSERT INTO \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_daily\\` (\n", + " order_date,\n", + " order_count,\n", + " order_total,\n", + " last_updated\n", + ")\n", + "SELECT\n", + " order_date,\n", + " orders_count,\n", + " total_value,\n", + " last_updated\n", + "FROM aggregate_snapshots\n", + "WHERE is_clerk_aggregated = 1 AND is_status_aggregated = 1;\n", + "\n", + "-- 2. Clerk & Status Daily Aggregates\n", + "INSERT INTO \\`${PROJECT_ID}.${BQ_DATASET}.stg_clerk_status_daily\\` (\n", + " clerk,\n", + " order_status,\n", + " order_date,\n", + " orders_processed,\n", + " total_revenue,\n", + " last_updated\n", + ")\n", + "SELECT\n", + " clerk,\n", + " order_status,\n", + " order_date,\n", + " orders_count,\n", + " total_value,\n", + " last_updated\n", + "FROM aggregate_snapshots\n", + "WHERE is_clerk_aggregated = 0 AND is_status_aggregated = 0;\n", + "\n", + "-- 3. Order Status Daily Aggregates\n", + "INSERT INTO \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_status_daily\\` (\n", + " order_status,\n", + " order_date,\n", + " orders_processed,\n", + " total_revenue,\n", + " last_updated\n", + ")\n", + "SELECT\n", + " order_status,\n", + " order_date,\n", + " orders_count,\n", + " total_value,\n", + " last_updated\n", + "FROM aggregate_snapshots\n", + "WHERE is_clerk_aggregated = 1 AND is_status_aggregated = 0;\n", + "EOF" + ] + }, + { + "cell_type": "markdown", + "id": "04f37535-a84b-4bc9-9348-167c1bf20447", + "metadata": {}, + "source": [ + "### 11.4 Register the scheduled query\n", + "\n", + "Register the SQL file as a BigQuery scheduled query running every 5 minutes." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "e0346c70-0833-49a7-bea3-973c8d78f738", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "DISPLAY_NAME=\"precompute_aggregates_to_staging\"\n", + "PARAMS_JSON=$(jq -n --rawfile q \"src/bq_scheduled/${DISPLAY_NAME}.sql\" '{\"query\": $q}')\n", + "\n", + "CONFIG_IDS=$(\n", + " bq ls \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --transfer_config \\\n", + " --transfer_location=\"$REGION\" \\\n", + " --format=json 2>/dev/null |\n", + " jq -r --arg display_name \"$DISPLAY_NAME\" '\n", + " .[]\n", + " | select(.displayName == $display_name)\n", + " | .name // empty\n", + " '\n", + ")\n", + "while IFS= read -r CONFIG_ID; do\n", + " [ -z \"$CONFIG_ID\" ] && continue\n", + " echo \"Removing existing scheduled query config: $CONFIG_ID\"\n", + " bq rm -f --transfer_config \"$CONFIG_ID\"\n", + "done <<< \"$CONFIG_IDS\"\n", + "\n", + "echo \"Creating scheduled query: $DISPLAY_NAME...\"\n", + "bq mk \\\n", + " --transfer_config \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --data_source=scheduled_query \\\n", + " --display_name=\"$DISPLAY_NAME\" \\\n", + " --location=\"$REGION\" \\\n", + " --schedule=\"every 5 minutes\" \\\n", + " --service_account_name=\"$SPARK_SA\" \\\n", + " --params=\"$PARAMS_JSON\"" + ] + }, + { + "cell_type": "markdown", + "id": "588fed20-4291-4617-a8b8-b184c1b4e306", + "metadata": {}, + "source": [ + "### 11.5 Create the Continuous Query SQL files for Spanner reverse ETL\n", + "\n", + "Create three Continuous Query definitions. Each query reads appends to its staging table, calculates the scalar average, and performs the reverse ETL upsert to its corresponding Spanner table.\n", + "\n", + "Each query replays one day of staging history when it starts. Combined with the two-day staging retention, this lets restarted jobs reprocess recent snapshots instead of creating a gap. Spanner's `change_timestamp_column` resolves replayed snapshots by `last_updated`." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "7b9bee01-47c2-4a7f-9791-d1b507e087cd", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "mkdir -p src/bq_continuous\n", + "\n", + "# ===========================================================================\n", + "# 1. Continuous Query -> Spanner: order_daily_aggregate\n", + "# ===========================================================================\n", + "cat < src/bq_continuous/cq_order_daily.sql\n", + "EXPORT DATA OPTIONS (\n", + " uri = 'https://spanner.googleapis.com/projects/${PROJECT_ID}/instances/${SPANNER_INSTANCE}/databases/${SPANNER_DATABASE}',\n", + " format = 'CLOUD_SPANNER',\n", + " spanner_options = '{\"table\": \"order_daily_aggregate\", \"change_timestamp_column\": \"last_updated\"}'\n", + ") AS\n", + "SELECT\n", + " order_date,\n", + " order_count,\n", + " order_total,\n", + " CAST(SAFE_DIVIDE(order_total, order_count) AS NUMERIC) AS order_total_average,\n", + " last_updated\n", + "FROM APPENDS(\n", + " TABLE \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_daily\\`,\n", + " CURRENT_TIMESTAMP() - INTERVAL 1 DAY\n", + ");\n", + "EOF\n", + "\n", + "# ===========================================================================\n", + "# 2. Continuous Query -> Spanner: clerk_status_daily_aggregate\n", + "# ===========================================================================\n", + "cat < src/bq_continuous/cq_clerk_status_daily.sql\n", + "EXPORT DATA OPTIONS (\n", + " uri = 'https://spanner.googleapis.com/projects/${PROJECT_ID}/instances/${SPANNER_INSTANCE}/databases/${SPANNER_DATABASE}',\n", + " format = 'CLOUD_SPANNER',\n", + " spanner_options = '{\"table\": \"clerk_status_daily_aggregate\", \"change_timestamp_column\": \"last_updated\"}'\n", + ") AS\n", + "SELECT\n", + " clerk,\n", + " order_status,\n", + " order_date,\n", + " orders_processed,\n", + " total_revenue,\n", + " CAST(SAFE_DIVIDE(total_revenue, orders_processed) AS NUMERIC) AS avg_order_value,\n", + " last_updated\n", + "FROM APPENDS(\n", + " TABLE \\`${PROJECT_ID}.${BQ_DATASET}.stg_clerk_status_daily\\`,\n", + " CURRENT_TIMESTAMP() - INTERVAL 1 DAY\n", + ");\n", + "EOF\n", + "\n", + "# ===========================================================================\n", + "# 3. Continuous Query -> Spanner: order_status_daily_aggregate\n", + "# ===========================================================================\n", + "cat < src/bq_continuous/cq_order_status_daily.sql\n", + "EXPORT DATA OPTIONS (\n", + " uri = 'https://spanner.googleapis.com/projects/${PROJECT_ID}/instances/${SPANNER_INSTANCE}/databases/${SPANNER_DATABASE}',\n", + " format = 'CLOUD_SPANNER',\n", + " spanner_options = '{\"table\": \"order_status_daily_aggregate\", \"change_timestamp_column\": \"last_updated\"}'\n", + ") AS\n", + "SELECT\n", + " order_status,\n", + " order_date,\n", + " orders_processed,\n", + " total_revenue,\n", + " CAST(SAFE_DIVIDE(total_revenue, orders_processed) AS NUMERIC) AS avg_order_value,\n", + " last_updated\n", + "FROM APPENDS(\n", + " TABLE \\`${PROJECT_ID}.${BQ_DATASET}.stg_order_status_daily\\`,\n", + " CURRENT_TIMESTAMP() - INTERVAL 1 DAY\n", + ");\n", + "EOF\n", + "\n", + "echo \"Continuous query SQL definitions created in src/bq_continuous/.\"" + ] + }, + { + "cell_type": "markdown", + "id": "7f21d425-4aca-42a2-8b7f-b7c9c4193439", + "metadata": {}, + "source": [ + "### 11.6 Start the Continuous Queries\n", + "\n", + "Submit the three Continuous Queries in the background with `--continuous=true` and `--sync=false`. The launch cell first cancels running jobs created by an earlier execution, preventing duplicate consumers when the cell is rerun.\n", + "\n", + "Continuous Queries started with a service account have a maximum runtime of 150 days. Rerun this cell before or after that limit, and after any extended interruption. The one-day replay makes a restart safe as long as the required staging partitions have not expired." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "52c14a67-3f91-4901-9b64-df9e909998d1", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "RUNNING_CQ_JOBS=$(\n", + " bq query \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --use_legacy_sql=false \\\n", + " --format=csv \\\n", + " --quiet \\\n", + " \"\n", + " SELECT job_id\n", + " FROM \\`region-${REGION}\\`.INFORMATION_SCHEMA.JOBS_BY_PROJECT\n", + " WHERE continuous IS TRUE\n", + " AND state = 'RUNNING'\n", + " AND (\n", + " STARTS_WITH(job_id, 'cq_order_daily_')\n", + " OR STARTS_WITH(job_id, 'cq_clerk_status_daily_')\n", + " OR STARTS_WITH(job_id, 'cq_order_status_daily_')\n", + " )\n", + " \" |\n", + " tail -n +2\n", + ")\n", + "\n", + "while IFS= read -r JOB_ID; do\n", + " [ -z \"$JOB_ID\" ] && continue\n", + " echo \"Cancelling existing Continuous Query: $JOB_ID...\"\n", + " bq --project_id=\"$PROJECT_ID\" --location=\"$REGION\" cancel \"$JOB_ID\"\n", + "done <<< \"$RUNNING_CQ_JOBS\"\n", + "\n", + "# Generate a timestamp suffix to guarantee uniqueness across runs\n", + "RUN_TS=$(date +%Y%m%d_%H%M%S)\n", + "\n", + "echo \"=== Launching BigQuery Continuous Queries with Service Account ($SPARK_SA) ===\"\n", + "\n", + "# 1. Order Daily Aggregate\n", + "echo \"Starting Continuous Query: order_daily_aggregate...\"\n", + "bq query \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --job_id=\"cq_order_daily_${RUN_TS}\" \\\n", + " --label=\"demo:sparkdataboost\" \\\n", + " --use_legacy_sql=false \\\n", + " --continuous=true \\\n", + " --sync=false \\\n", + " --connection_property=service_account=\"$SPARK_SA\" \\\n", + " < src/bq_continuous/cq_order_daily.sql\n", + "\n", + "# 2. Clerk & Status Daily Aggregate\n", + "echo \"Starting Continuous Query: clerk_status_daily_aggregate...\"\n", + "bq query \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --job_id=\"cq_clerk_status_daily_${RUN_TS}\" \\\n", + " --label=\"demo:sparkdataboost\" \\\n", + " --use_legacy_sql=false \\\n", + " --continuous=true \\\n", + " --sync=false \\\n", + " --connection_property=service_account=\"$SPARK_SA\" \\\n", + " < src/bq_continuous/cq_clerk_status_daily.sql\n", + "\n", + "# 3. Order Status Daily Aggregate\n", + "echo \"Starting Continuous Query: order_status_daily_aggregate...\"\n", + "bq query \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --location=\"$REGION\" \\\n", + " --job_id=\"cq_order_status_daily_${RUN_TS}\" \\\n", + " --label=\"demo:sparkdataboost\" \\\n", + " --use_legacy_sql=false \\\n", + " --continuous=true \\\n", + " --sync=false \\\n", + " --connection_property=service_account=\"$SPARK_SA\" \\\n", + " < src/bq_continuous/cq_order_status_daily.sql\n", + "\n", + "echo \"All 3 Continuous Queries successfully started in the background under $SPARK_SA.\"" + ] + }, + { + "cell_type": "markdown", + "id": "9538ab60", + "metadata": {}, + "source": [ + "## 12. Produce and verify demo data\n", + "\n", + "### 12.1 Run the producer\n", + "\n", + "The data flow is now ready to process orders. We can start the Kafka producer." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "4a955250", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud compute ssh \"$KAFKA_VM\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --tunnel-through-iap \\\n", + " --command='\n", + " set -euo pipefail\n", + "\n", + " source ~/.env\n", + "\n", + " docker run --rm \\\n", + " --name tpch-generator \\\n", + " --network host \\\n", + " tpch-generator:local \\\n", + " --local \\\n", + " --bootstrap-servers=\"${KAFKA_BROKER}\" \\\n", + " --max-messages=1000 \\\n", + " --target-throughput=100 \\\n", + " --topic=\"${KAFKA_TOPIC}\" \\\n", + " --date-from=\"1 month ago\" \\\n", + " --date-to=\"today\"\n", + " '\n" + ] + }, + { + "cell_type": "markdown", + "id": "01003dd5", + "metadata": {}, + "source": [ + "With the settings above, the producer will send 1000 messages at a rate of 100 messages per second, which will take about 10 seconds to complete.\n", + "\n", + "In GCP console, you can check that Lakehouse and Spanner tables are being populated.\n" + ] + }, + { + "cell_type": "markdown", + "id": "bq_12_2_md", + "metadata": {}, + "source": [ + "### 12.2 Verify Spanner aggregate table population\n", + "\n", + "Note: The scheduled query runs every five minutes. Allow at least one complete schedule interval before expecting aggregate rows.\n", + "\n", + "Check Spanner after the scheduled query has produced staging snapshots and the Continuous Queries have exported them. If the result is empty, try again shortly. Data propagation to Spanner may take a moment." + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "bq_12_2_code", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "echo \"=== Spanner Aggregate Tables Record Counts ===\"\n", + "for TABLE in order_daily_aggregate clerk_status_daily_aggregate order_status_daily_aggregate; do\n", + " echo \"Table: $TABLE\"\n", + " gcloud spanner databases execute-sql \"$SPANNER_DATABASE\" \\\n", + " --instance=\"$SPANNER_INSTANCE\" \\\n", + " --sql=\"SELECT COUNT(*) AS orders_count, MAX(last_updated) AS latest_sync FROM $TABLE;\"\n", + "done\n" + ] + }, + { + "cell_type": "markdown", + "id": "d57bf893", + "metadata": {}, + "source": [ + "## 13. Visualize the data in JupyterLab Notebook\n", + "\n", + "In your Managed Spark cluster, navigate to Web Interfaces.\n", + "\n", + "![spark-webinterfaces.png](imgs/spark-webinterfaces.png)\n", + "\n", + "Open a **JupyterLab** notebook.\n", + "\n", + "From the Launcher, create a new **Python 3** notebook. (Note: Select Python 3, not PySpark.)\n", + "\n", + "> **Run the next two code cells there, not on your local kernel.** They need\n", + "> the cluster's private-network access to Spanner and rely on packages\n", + "> (`pyspark`, `matplotlib`, `seaborn`) preinstalled on the cluster's Jupyter\n", + "> image.\n", + "\n", + "### Part 1: Dashboard querying pre-aggregated tables\n", + "\n", + "Paste this cell into the cluster's Python 3 notebook to query the\n", + "pre-computed aggregate tables from the Lakehouse path.\n", + "\n", + "**Edit `PROJECT_ID = \"your-project-id\"` with your GCP project.**\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "89557a37", + "metadata": {}, + "outputs": [], + "source": [ + "# =========================================================================\n", + "# SPANNER ORDERS DASHBOARD\n", + "# Rolling one-month reporting period\n", + "#\n", + "# Dashboard panels:\n", + "# 1. Daily total revenue\n", + "# 2A. Daily revenue by order status\n", + "# 2B. Daily revenue percentage by order status\n", + "# 3. Top 10 clerks by order count across all statuses\n", + "# 4. Top 10 clerks by total revenue across all statuses\n", + "# =========================================================================\n", + "\n", + "import matplotlib.pyplot as plt\n", + "import seaborn as sns\n", + "import pandas as pd\n", + "\n", + "from pyspark.sql import SparkSession\n", + "from pyspark.sql import functions as F\n", + "\n", + "PROJECT_ID = \"your-project-id\"\n", + "SPANNER_INSTANCE = \"spark-databoost-demo-instance\"\n", + "SPANNER_DATABASE = \"demo\"\n", + "\n", + "if PROJECT_ID == \"your-project-id\":\n", + " raise ValueError(\"Set PROJECT_ID to the project used in section 1\")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 1. INITIALIZE SPARK SESSION WITH THE SPANNER CONNECTOR\n", + "# -------------------------------------------------------------------------\n", + "\n", + "spark = (\n", + " SparkSession.builder\n", + " .appName(\"Spanner Aggregates Dashboard\")\n", + " .config(\n", + " \"spark.jars\",\n", + " \"gs://spark-lib/spanner/spark-3.5-spanner-1.5.0.jar\"\n", + " )\n", + " .getOrCreate()\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 2. DEFINE THE ROLLING ONE-MONTH REPORTING PERIOD\n", + "# -------------------------------------------------------------------------\n", + "\n", + "# Inclusive date range:\n", + "# start_date = one calendar month before today\n", + "# end_date = today\n", + "\n", + "start_date = F.add_months(F.current_date(), -1)\n", + "end_date = F.current_date()\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 3. READ AND FILTER THE PRE-AGGREGATED SPANNER TABLES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "df_day_agg = (\n", + " spark.read\n", + " .format(\"cloud-spanner\")\n", + " .option(\"projectId\", PROJECT_ID)\n", + " .option(\"instanceId\", SPANNER_INSTANCE)\n", + " .option(\"databaseId\", SPANNER_DATABASE)\n", + " .option(\"table\", \"order_daily_aggregate\")\n", + " .load()\n", + " .filter(\n", + " F.to_date(F.col(\"order_date\")).between(start_date, end_date)\n", + " )\n", + ")\n", + "\n", + "\n", + "df_status_agg = (\n", + " spark.read\n", + " .format(\"cloud-spanner\")\n", + " .option(\"projectId\", PROJECT_ID)\n", + " .option(\"instanceId\", SPANNER_INSTANCE)\n", + " .option(\"databaseId\", SPANNER_DATABASE)\n", + " .option(\"table\", \"order_status_daily_aggregate\")\n", + " .load()\n", + " .filter(\n", + " F.to_date(F.col(\"order_date\")).between(start_date, end_date)\n", + " )\n", + ")\n", + "\n", + "\n", + "df_clerk_status_agg = (\n", + " spark.read\n", + " .format(\"cloud-spanner\")\n", + " .option(\"projectId\", PROJECT_ID)\n", + " .option(\"instanceId\", SPANNER_INSTANCE)\n", + " .option(\"databaseId\", SPANNER_DATABASE)\n", + " .option(\"table\", \"clerk_status_daily_aggregate\")\n", + " .load()\n", + " .filter(\n", + " F.to_date(F.col(\"order_date\")).between(start_date, end_date)\n", + " )\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 4. CONVERT THE SPARK DATAFRAMES TO PANDAS DATAFRAMES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "pdf_day = df_day_agg.orderBy(\"order_date\").toPandas()\n", + "pdf_status = df_status_agg.orderBy(\"order_date\").toPandas()\n", + "pdf_clerk = df_clerk_status_agg.orderBy(\"order_date\").toPandas()\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 5. PREPARE THE PANDAS DATAFRAMES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "for pdf in [pdf_day, pdf_status, pdf_clerk]:\n", + " pdf[\"order_date\"] = pd.to_datetime(\n", + " pdf[\"order_date\"],\n", + " errors=\"coerce\"\n", + " )\n", + "\n", + "\n", + "# Convert Spanner NUMERIC and INT64 values into Pandas numeric types.\n", + "\n", + "pdf_day[\"order_total\"] = pd.to_numeric(\n", + " pdf_day[\"order_total\"],\n", + " errors=\"coerce\"\n", + ").fillna(0.0)\n", + "\n", + "pdf_status[\"total_revenue\"] = pd.to_numeric(\n", + " pdf_status[\"total_revenue\"],\n", + " errors=\"coerce\"\n", + ").fillna(0.0)\n", + "\n", + "pdf_clerk[\"total_revenue\"] = pd.to_numeric(\n", + " pdf_clerk[\"total_revenue\"],\n", + " errors=\"coerce\"\n", + ").fillna(0.0)\n", + "\n", + "pdf_clerk[\"orders_processed\"] = pd.to_numeric(\n", + " pdf_clerk[\"orders_processed\"],\n", + " errors=\"coerce\"\n", + ").fillna(0)\n", + "\n", + "\n", + "# Make order-status labels more descriptive.\n", + "\n", + "status_labels = {\n", + " \"F\": \"F - Fulfilled\",\n", + " \"O\": \"O - Open\",\n", + " \"P\": \"P - Processing\"\n", + "}\n", + "\n", + "pdf_status[\"order_status_label\"] = (\n", + " pdf_status[\"order_status\"]\n", + " .map(status_labels)\n", + " .fillna(pdf_status[\"order_status\"])\n", + ")\n", + "\n", + "pdf_clerk[\"order_status_label\"] = (\n", + " pdf_clerk[\"order_status\"]\n", + " .map(status_labels)\n", + " .fillna(pdf_clerk[\"order_status\"])\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 6. CREATE THE ORDER-STATUS PIVOT TABLES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "pivot_status_abs = (\n", + " pdf_status\n", + " .pivot_table(\n", + " index=\"order_date\",\n", + " columns=\"order_status_label\",\n", + " values=\"total_revenue\",\n", + " aggfunc=\"sum\",\n", + " fill_value=0\n", + " )\n", + " .sort_index()\n", + ")\n", + "\n", + "pivot_status_pct = (\n", + " pivot_status_abs\n", + " .div(\n", + " pivot_status_abs.sum(axis=1).replace(0, pd.NA),\n", + " axis=0\n", + " )\n", + " .mul(100)\n", + " .fillna(0)\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 7. GENERATE THE ACTUAL REPORTING-PERIOD LABEL\n", + "# -------------------------------------------------------------------------\n", + "\n", + "all_dates = pd.concat(\n", + " [\n", + " pdf_day[\"order_date\"],\n", + " pdf_status[\"order_date\"],\n", + " pdf_clerk[\"order_date\"]\n", + " ],\n", + " ignore_index=True\n", + ").dropna()\n", + "\n", + "if not all_dates.empty:\n", + " period_start_label = all_dates.min().strftime(\"%b %d, %Y\")\n", + " period_end_label = all_dates.max().strftime(\"%b %d, %Y\")\n", + " period_label = f\"{period_start_label} to {period_end_label}\"\n", + "else:\n", + " period_label = \"Rolling One-Month Period\"\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 8. CALCULATE THE TOP 10 CLERKS BY ORDER COUNT\n", + "# -------------------------------------------------------------------------\n", + "\n", + "# The aggregate table has one row per clerk, status, and date.\n", + "# Summing only by clerk combines all statuses and dates in the period.\n", + "\n", + "top_clerks_by_count = (\n", + " pdf_clerk\n", + " .groupby(\"clerk\", as_index=False)\n", + " .agg(order_count=(\"orders_processed\", \"sum\"))\n", + " .nlargest(10, \"order_count\")\n", + " .sort_values(\"order_count\", ascending=False)\n", + " .reset_index(drop=True)\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 9. CALCULATE THE TOP 10 CLERKS BY TOTAL REVENUE\n", + "# -------------------------------------------------------------------------\n", + "\n", + "# This ranking is independent of the order-count ranking.\n", + "\n", + "top_clerks_by_revenue = (\n", + " pdf_clerk\n", + " .groupby(\"clerk\", as_index=False)\n", + " .agg(total_revenue=(\"total_revenue\", \"sum\"))\n", + " .nlargest(10, \"total_revenue\")\n", + " .sort_values(\"total_revenue\", ascending=False)\n", + " .reset_index(drop=True)\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 10. BUILD THE FIVE-PANEL DASHBOARD\n", + "# -------------------------------------------------------------------------\n", + "\n", + "sns.set_theme(style=\"whitegrid\")\n", + "\n", + "fig, axes = plt.subplots(\n", + " nrows=5,\n", + " ncols=1,\n", + " figsize=(14, 29),\n", + " gridspec_kw={\n", + " \"height_ratios\": [1, 1.25, 1.25, 1.6, 1.6]\n", + " }\n", + ")\n", + "\n", + "fig.suptitle(\n", + " (\n", + " \"Spanner Orders Dashboard\\n\"\n", + " f\"Rolling One-Month Period: {period_label}\"\n", + " ),\n", + " fontsize=17,\n", + " fontweight=\"bold\",\n", + " y=0.995\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# GRAPH 1: TOTAL DAILY REVENUE TREND\n", + "# -------------------------------------------------------------------------\n", + "\n", + "sns.lineplot(\n", + " ax=axes[0],\n", + " data=pdf_day,\n", + " x=\"order_date\",\n", + " y=\"order_total\",\n", + " marker=\"o\",\n", + " color=\"navy\",\n", + " linewidth=2.5\n", + ")\n", + "\n", + "axes[0].set_title(\n", + " f\"1. Daily Total Revenue | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\"\n", + ")\n", + "axes[0].set_xlabel(\"Order Date\")\n", + "axes[0].set_ylabel(\"Revenue ($)\")\n", + "axes[0].tick_params(axis=\"x\", rotation=30)\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# GRAPH 2A: DAILY REVENUE BREAKDOWN BY ORDER STATUS\n", + "# -------------------------------------------------------------------------\n", + "\n", + "pivot_status_abs.plot(\n", + " kind=\"bar\",\n", + " stacked=True,\n", + " ax=axes[1],\n", + " colormap=\"Set2\",\n", + " edgecolor=\"black\",\n", + " linewidth=0.8,\n", + " alpha=0.85\n", + ")\n", + "\n", + "axes[1].set_title(\n", + " f\"2A. Daily Revenue by Order Status | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\"\n", + ")\n", + "axes[1].set_ylabel(\"Total Revenue ($)\")\n", + "axes[1].set_xlabel(\"Order Date\")\n", + "axes[1].legend(title=\"Order Status\", loc=\"upper left\")\n", + "axes[1].set_xticklabels(\n", + " [date.strftime(\"%b %d\") for date in pivot_status_abs.index],\n", + " rotation=30,\n", + " ha=\"right\"\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# GRAPH 2B: DAILY REVENUE DISTRIBUTION BY ORDER STATUS\n", + "# -------------------------------------------------------------------------\n", + "\n", + "pivot_status_pct.plot(\n", + " kind=\"bar\",\n", + " stacked=True,\n", + " ax=axes[2],\n", + " colormap=\"Set2\",\n", + " edgecolor=\"black\",\n", + " linewidth=0.8,\n", + " alpha=0.85\n", + ")\n", + "\n", + "axes[2].set_title(\n", + " f\"2B. Daily Revenue Distribution by Order Status | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\"\n", + ")\n", + "axes[2].set_ylabel(\"Percentage (%)\")\n", + "axes[2].set_xlabel(\"Order Date\")\n", + "axes[2].set_ylim(0, 100)\n", + "axes[2].legend(title=\"Order Status\", loc=\"upper left\")\n", + "axes[2].set_xticklabels(\n", + " [date.strftime(\"%b %d\") for date in pivot_status_pct.index],\n", + " rotation=30,\n", + " ha=\"right\"\n", + ")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# GRAPH 3: TOP 10 CLERKS BY ORDER COUNT, ALL STATUSES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "if (\n", + " not top_clerks_by_count.empty\n", + " and top_clerks_by_count[\"order_count\"].sum() > 0\n", + "):\n", + " count_pie_labels = (\n", + " top_clerks_by_count[\"clerk\"]\n", + " .astype(str)\n", + " .str.replace(r\"^Clerk#0*\", \"Clerk#\", regex=True)\n", + " )\n", + "\n", + " top_10_total_orders = top_clerks_by_count[\"order_count\"].sum()\n", + "\n", + " def format_order_pie_label(percent):\n", + " count = int(round(percent * top_10_total_orders / 100))\n", + " return f\"{percent:.1f}%\\n{count:,} orders\"\n", + "\n", + " count_pie_colors = sns.color_palette(\n", + " \"Set2\",\n", + " n_colors=len(top_clerks_by_count)\n", + " )\n", + "\n", + " count_wedges, count_label_texts, count_value_texts = axes[3].pie(\n", + " top_clerks_by_count[\"order_count\"],\n", + " labels=count_pie_labels,\n", + " autopct=format_order_pie_label,\n", + " startangle=90,\n", + " counterclock=False,\n", + " colors=count_pie_colors,\n", + " pctdistance=0.72,\n", + " labeldistance=1.08,\n", + " wedgeprops={\n", + " \"edgecolor\": \"white\",\n", + " \"linewidth\": 1.5\n", + " },\n", + " textprops={\n", + " \"fontsize\": 9\n", + " }\n", + " )\n", + "\n", + " for value_text in count_value_texts:\n", + " value_text.set_fontsize(8)\n", + " value_text.set_fontweight(\"bold\")\n", + "\n", + " axes[3].set_title(\n", + " f\"3. Top 10 Clerks by Order Count, All Statuses | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\",\n", + " pad=18\n", + " )\n", + " axes[3].axis(\"equal\")\n", + " axes[3].legend(\n", + " count_wedges,\n", + " [\n", + " f\"{clerk}: {int(count):,} orders\"\n", + " for clerk, count in zip(\n", + " count_pie_labels,\n", + " top_clerks_by_count[\"order_count\"]\n", + " )\n", + " ],\n", + " title=(\n", + " \"Clerk and Order Count\\n\"\n", + " f\"Top 10 Total: {int(top_10_total_orders):,}\"\n", + " ),\n", + " loc=\"center left\",\n", + " bbox_to_anchor=(1.02, 0.5),\n", + " fontsize=9\n", + " )\n", + "else:\n", + " axes[3].text(\n", + " 0.5,\n", + " 0.5,\n", + " \"No clerk order data available\\nfor the selected reporting period\",\n", + " horizontalalignment=\"center\",\n", + " verticalalignment=\"center\",\n", + " transform=axes[3].transAxes,\n", + " fontsize=12\n", + " )\n", + " axes[3].set_title(\n", + " f\"3. Top 10 Clerks by Order Count, All Statuses | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\"\n", + " )\n", + " axes[3].axis(\"off\")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# GRAPH 4: TOP 10 CLERKS BY TOTAL REVENUE, ALL STATUSES\n", + "# -------------------------------------------------------------------------\n", + "\n", + "if (\n", + " not top_clerks_by_revenue.empty\n", + " and top_clerks_by_revenue[\"total_revenue\"].sum() > 0\n", + "):\n", + " revenue_pie_labels = (\n", + " top_clerks_by_revenue[\"clerk\"]\n", + " .astype(str)\n", + " .str.replace(r\"^Clerk#0*\", \"Clerk#\", regex=True)\n", + " )\n", + "\n", + " top_10_total_revenue = top_clerks_by_revenue[\"total_revenue\"].sum()\n", + "\n", + " def format_revenue_pie_label(percent):\n", + " return f\"{percent:.1f}%\"\n", + "\n", + " revenue_pie_colors = sns.color_palette(\n", + " \"Set3\",\n", + " n_colors=len(top_clerks_by_revenue)\n", + " )\n", + "\n", + " revenue_wedges, revenue_label_texts, revenue_value_texts = axes[4].pie(\n", + " top_clerks_by_revenue[\"total_revenue\"],\n", + " labels=revenue_pie_labels,\n", + " autopct=format_revenue_pie_label,\n", + " startangle=90,\n", + " counterclock=False,\n", + " colors=revenue_pie_colors,\n", + " pctdistance=0.72,\n", + " labeldistance=1.08,\n", + " wedgeprops={\n", + " \"edgecolor\": \"white\",\n", + " \"linewidth\": 1.5\n", + " },\n", + " textprops={\n", + " \"fontsize\": 9\n", + " }\n", + " )\n", + "\n", + " for value_text in revenue_value_texts:\n", + " value_text.set_fontsize(9)\n", + " value_text.set_fontweight(\"bold\")\n", + "\n", + " axes[4].set_title(\n", + " f\"4. Top 10 Clerks by Total Order Revenue, All Statuses | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\",\n", + " pad=18\n", + " )\n", + " axes[4].axis(\"equal\")\n", + " axes[4].legend(\n", + " revenue_wedges,\n", + " [\n", + " f\"{clerk}: ${revenue:,.2f}\"\n", + " for clerk, revenue in zip(\n", + " revenue_pie_labels,\n", + " top_clerks_by_revenue[\"total_revenue\"]\n", + " )\n", + " ],\n", + " title=(\n", + " \"Clerk and Total Revenue\\n\"\n", + " f\"Top 10 Total: ${top_10_total_revenue:,.2f}\"\n", + " ),\n", + " loc=\"center left\",\n", + " bbox_to_anchor=(1.02, 0.5),\n", + " fontsize=9\n", + " )\n", + "else:\n", + " axes[4].text(\n", + " 0.5,\n", + " 0.5,\n", + " \"No clerk revenue data available\\nfor the selected reporting period\",\n", + " horizontalalignment=\"center\",\n", + " verticalalignment=\"center\",\n", + " transform=axes[4].transAxes,\n", + " fontsize=12\n", + " )\n", + " axes[4].set_title(\n", + " f\"4. Top 10 Clerks by Total Order Revenue, All Statuses | {period_label}\",\n", + " fontsize=13,\n", + " fontweight=\"bold\"\n", + " )\n", + " axes[4].axis(\"off\")\n", + "\n", + "\n", + "# -------------------------------------------------------------------------\n", + "# 11. FINALIZE AND DISPLAY THE DASHBOARD\n", + "# -------------------------------------------------------------------------\n", + "\n", + "# Reserve space on the right for both pie-chart legends.\n", + "plt.tight_layout(rect=[0, 0, 0.82, 0.985])\n", + "plt.show()\n" + ] + }, + { + "cell_type": "markdown", + "id": "e63d61a2", + "metadata": {}, + "source": [ + "### Part 2: Dashboard with ad-hoc queries using Data Boost\n", + "\n", + "First, stop the kernel of the first notebook to free-up resources in the Spark cluster.\n", + "\n", + "\n", + "\n", + "Using Spanner Data Boost, we can run ad-hoc queries on the raw `orders`\n", + "table. Paste this cell into a second **Python 3** notebook on the cluster.\n", + "\n", + "**Edit `PROJECT_ID = \"your-project-id\"` with your GCP project.**\n", + "\n", + "Note:\n", + "- Data Boost is enabled via `.option(\"enableDataBoost\", \"true\")`.\n", + "- To force using the Columnar scan we are using Spanner `@{scan_method=columnar}` hint in the SQL query sent to Spanner.\n", + "- Sending a direct query `.option(\"query\", query)` is introduced by spark-3.5-spanner-1.5.0 (version >= 1.5.0).\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "6ad4e999", + "metadata": {}, + "outputs": [], + "source": [ + "from datetime import date, timedelta\n", + "from IPython.display import display, Markdown\n", + "from pyspark.sql import SparkSession\n", + "from pyspark.sql import functions as F\n", + "\n", + "PROJECT_ID = \"your-project-id\"\n", + "SPANNER_INSTANCE = \"spark-databoost-demo-instance\"\n", + "SPANNER_DATABASE = \"demo\"\n", + "\n", + "if PROJECT_ID == \"your-project-id\":\n", + " raise ValueError(\"Set PROJECT_ID to the project used in section 1\")\n", + "\n", + "spark = (\n", + " SparkSession.builder\n", + " .appName(\"Spanner Aggregates Data Boost Ad-hoc Queries - Direct Query\")\n", + " .config(\n", + " \"spark.jars\",\n", + " \"gs://spark-lib/spanner/spark-3.5-spanner-1.5.0.jar\"\n", + " )\n", + " .getOrCreate()\n", + ")\n", + "\n", + "cutoff_date = (date.today() - timedelta(days=3)).isoformat()\n", + "\n", + "# Force columnar scan via the @{scan_method=columnar} hint\n", + "query = f\"\"\"\n", + "@{{scan_method=columnar}}\n", + "SELECT \n", + " order_key,\n", + " customer_key,\n", + " order_status,\n", + " total_price,\n", + " order_date,\n", + " order_priority,\n", + " clerk\n", + "FROM orders\n", + "WHERE order_status != 'F'\n", + " AND order_date < DATE '{cutoff_date}'\n", + "\"\"\"\n", + "\n", + "# Read using direct query and Data Boost\n", + "unfulfilled_orders_raw = (\n", + " spark.read\n", + " .format(\"cloud-spanner\")\n", + " .option(\"projectId\", PROJECT_ID)\n", + " .option(\"instanceId\", SPANNER_INSTANCE)\n", + " .option(\"databaseId\", SPANNER_DATABASE)\n", + " .option(\"query\", query)\n", + " .option(\"enableDataBoost\", \"true\")\n", + " .load()\n", + ")\n", + "\n", + "# Order in Spark\n", + "unfulfilled_orders = unfulfilled_orders_raw.orderBy(F.col(\"order_date\").asc())\n", + "\n", + "display(Markdown(\"### Spark Physical Plan\"))\n", + "unfulfilled_orders.explain(\"formatted\")\n", + "\n", + "display(Markdown(\"### ⚠️ Unfulfilled Orders (> 3 Days Old)\"))\n", + "unfulfilled_orders.show(truncate=False)" + ] + }, + { + "cell_type": "markdown", + "id": "440c73a4", + "metadata": {}, + "source": [ + "## 14. Cleanup" + ] + }, + { + "cell_type": "markdown", + "id": "2ec349b7", + "metadata": {}, + "source": [ + "### 14.1 Clean-up Spark resources" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "eb0dd62c", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "# 1. Stop Spark Dataproc streaming jobs\n", + "echo \"=== Stopping Spark Dataproc Jobs ===\"\n", + "for JOB_ID in $(\n", + " gcloud dataproc jobs list \\\n", + " --region=\"$REGION\" \\\n", + " --cluster=\"$CLUSTER\" \\\n", + " --filter='status.state=RUNNING OR status.state=PENDING OR status.state=SETUP_DONE' \\\n", + " --format='value(reference.jobId)' 2>/dev/null\n", + "); do\n", + " echo \"Killing Dataproc job: $JOB_ID...\"\n", + " gcloud dataproc jobs kill \"$JOB_ID\" --region=\"$REGION\" --quiet || true\n", + "done\n", + "\n", + "# 2. Delete the Spark cluster\n", + "gcloud dataproc clusters delete \"$CLUSTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --quiet || true\n" + ] + }, + { + "cell_type": "markdown", + "id": "14_bq_clean_md", + "metadata": {}, + "source": [ + "### 14.2 Clean-up BigQuery resources" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "14_bq_clean_code", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "echo \"=== Cleaning up BigQuery Resources ===\"\n", + "\n", + "# 1. Cancel running BigQuery Continuous Queries\n", + "echo \"=== Stopping demo Continuous Queries ===\"\n", + "\n", + "for JOB_ID in $(\n", + " bq query --project_id=\"$PROJECT_ID\" --location=\"$REGION\" --use_legacy_sql=false --format=csv --quiet \"\n", + " SELECT job_id\n", + " FROM \\`region-${REGION}\\`.INFORMATION_SCHEMA.JOBS_BY_PROJECT\n", + " WHERE continuous IS TRUE\n", + " AND state = 'RUNNING'\n", + " AND (\n", + " STARTS_WITH(job_id, 'cq_order_daily_')\n", + " OR STARTS_WITH(job_id, 'cq_clerk_status_daily_')\n", + " OR STARTS_WITH(job_id, 'cq_order_status_daily_')\n", + " )\n", + " \" 2>/dev/null | tail -n +2\n", + "); do\n", + " [[ -z \"$JOB_ID\" ]] && continue\n", + " echo \"Cancelling Continuous Query: $JOB_ID...\"\n", + " bq --project_id=\"$PROJECT_ID\" --location=\"$REGION\" cancel \"$JOB_ID\" || true\n", + "done\n", + "\n", + "# 2. Delete the scheduled query\n", + "DISPLAY_NAME=\"precompute_aggregates_to_staging\"\n", + "while IFS= read -r CONFIG_ID; do\n", + " [[ -z \"$CONFIG_ID\" ]] && continue\n", + " echo \"Deleting scheduled query: $DISPLAY_NAME ($CONFIG_ID)...\"\n", + " bq rm -f --transfer_config \"$CONFIG_ID\" || true\n", + "done < <(\n", + " bq ls \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --transfer_config \\\n", + " --transfer_location=\"$REGION\" \\\n", + " --format=json 2>/dev/null |\n", + " jq -r --arg display_name \"$DISPLAY_NAME\" '\n", + " .[]\n", + " | select(.displayName == $display_name)\n", + " | .name // empty\n", + " '\n", + ")\n", + "\n", + "# 3. Delete intermediate staging dataset and tables\n", + "echo \"Deleting BigQuery staging dataset: ${BQ_DATASET}...\"\n", + "bq rm -r -f -d \"${PROJECT_ID}:${BQ_DATASET}\" || true\n", + "\n", + "# 4. Delete CONTINUOUS reservation assignments\n", + "echo \"=== Deleting BigQuery Reservation Assignments ===\"\n", + "ASSIGNMENT_NAMES=$(\n", + " bq ls \\\n", + " --reservation_assignment \\\n", + " --location=\"$REGION\" \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " --format=prettyjson 2>/dev/null |\n", + " jq -r \\\n", + " --arg reservation \"$BQ_RESERVATION\" \\\n", + " '\n", + " .[]\n", + " | select(\n", + " .name\n", + " | contains(\"/reservations/\" + $reservation + \"/assignments/\")\n", + " )\n", + " | .name\n", + " '\n", + ")\n", + "\n", + "if [ -z \"$ASSIGNMENT_NAMES\" ]; then\n", + " echo \"No assignments found for reservation: $BQ_RESERVATION\"\n", + "else\n", + " while IFS= read -r ASSIGNMENT_NAME; do\n", + " [[ -z \"$ASSIGNMENT_NAME\" ]] && continue\n", + " ASSIGNMENT_ID=\"${ASSIGNMENT_NAME##*/}\"\n", + " RESERVATION_PATH=\"${ASSIGNMENT_NAME%/assignments/*}\"\n", + " RESERVATION_ID=\"${RESERVATION_PATH##*/}\"\n", + " ASSIGNMENT_REF=\"${RESERVATION_ID}.${ASSIGNMENT_ID}\"\n", + "\n", + " echo \"Deleting assignment: $ASSIGNMENT_REF\"\n", + " bq rm \\\n", + " --force \\\n", + " --reservation_assignment \\\n", + " --location=\"$REGION\" \\\n", + " --project_id=\"$PROJECT_ID\" \\\n", + " \"$ASSIGNMENT_REF\"\n", + " done <<< \"$ASSIGNMENT_NAMES\"\n", + "fi\n", + "\n", + "# 5. Delete the reservation\n", + "if bq show --reservation --location=\"$REGION\" --project_id=\"$PROJECT_ID\" \"$BQ_RESERVATION\" &>/dev/null; then\n", + " echo \"Deleting reservation: $BQ_RESERVATION...\"\n", + " bq rm -f --reservation --location=\"$REGION\" --project_id=\"$PROJECT_ID\" \"$BQ_RESERVATION\" || true\n", + "fi\n", + "\n", + "echo \"All BigQuery resources cleaned up.\"" + ] + }, + { + "cell_type": "markdown", + "id": "72651b10", + "metadata": {}, + "source": [ + "### 14.3 Delete compute and database resources" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "7b6eb0d6", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud spanner instances delete \"$SPANNER_INSTANCE\" \\\n", + " --quiet || true\n", + "\n", + "gcloud compute instances delete \"$KAFKA_VM\" \\\n", + " --zone=\"$ZONE\" \\\n", + " --quiet || true\n" + ] + }, + { + "cell_type": "markdown", + "id": "09724021", + "metadata": {}, + "source": [ + "### 14.4 Delete the Lakehouse catalog\n", + "\n", + "In the Google Cloud console:\n", + "\n", + "1. Open **BigQuery / Lakehouse**.\n", + "2. Select the Iceberg REST catalog.\n", + "3. Delete the catalog.\n" + ] + }, + { + "cell_type": "markdown", + "id": "d0dbe748", + "metadata": {}, + "source": [ + "### 14.5 Delete the storage buckets" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "e4080267", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud storage rm \\\n", + " --recursive \\\n", + " \"gs://${JOB_BUCKET}\" || true\n", + "\n", + "gcloud storage rm \\\n", + " --recursive \\\n", + " \"gs://${WAREHOUSE_BUCKET}\" || true\n" + ] + }, + { + "cell_type": "markdown", + "id": "b590aaee", + "metadata": {}, + "source": [ + "### 14.6 Delete the network resources" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "be94815a", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud compute routers nats delete \"$NAT\" \\\n", + " --router=\"$ROUTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --quiet || true\n", + "\n", + "gcloud compute routers delete \"$ROUTER\" \\\n", + " --region=\"$REGION\" \\\n", + " --quiet || true\n", + "\n", + "gcloud compute firewall-rules delete \\\n", + " \"allow-${NETWORK}-internal\" \\\n", + " \"allow-${NETWORK}-ssh-iap\" \\\n", + " --quiet || true\n", + "\n", + "gcloud compute networks subnets delete \"$SUBNET\" \\\n", + " --region=\"$REGION\" \\\n", + " --quiet || true\n", + "gcloud compute networks delete \"$NETWORK\" \\\n", + " --quiet || true\n" + ] + }, + { + "cell_type": "markdown", + "id": "c30a5a0e", + "metadata": {}, + "source": [ + "### 14.7 Delete the Service Account" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "2923a8f7", + "metadata": {}, + "outputs": [], + "source": [ + "%%bash\n", + "source .env\n", + "set -euo pipefail\n", + "\n", + "gcloud iam service-accounts delete \"$SPARK_SA\" \\\n", + " --quiet || true\n" + ] + }, + { + "cell_type": "code", + "execution_count": null, + "id": "df8043ee", + "metadata": {}, + "outputs": [], + "source": [ + "# Copyright 2025 Google LLC\n", + "#\n", + "# Licensed under the Apache License, Version 2.0 (the \"License\");\n", + "# you may not use this file except in compliance with the License.\n", + "# You may obtain a copy of the License at\n", + "#\n", + "# https://www.apache.org/licenses/LICENSE-2.0\n", + "#\n", + "# Unless required by applicable law or agreed to in writing, software\n", + "# distributed under the License is distributed on an \"AS IS\" BASIS,\n", + "# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.\n", + "# See the License for the specific language governing permissions and\n", + "# limitations under the License." + ] + } + ], + "metadata": { + "kernelspec": { + "display_name": "Python 3 (ipykernel)", + "language": "python", + "name": "python3" + }, + "language_info": { + "codemirror_mode": { + "name": "ipython", + "version": 3 + }, + "file_extension": ".py", + "mimetype": "text/x-python", + "name": "python", + "nbconvert_exporter": "python", + "pygments_lexer": "ipython3", + "version": "3.14.7" + } + }, + "nbformat": 4, + "nbformat_minor": 5 +} diff --git a/orders-streaming-analytics/imgs/shutdown-kernel.png b/orders-streaming-analytics/imgs/shutdown-kernel.png new file mode 100644 index 0000000..3dc4816 Binary files /dev/null and b/orders-streaming-analytics/imgs/shutdown-kernel.png differ diff --git a/orders-streaming-analytics/imgs/spark-webinterfaces.png b/orders-streaming-analytics/imgs/spark-webinterfaces.png new file mode 100644 index 0000000..9fa1b2f Binary files /dev/null and b/orders-streaming-analytics/imgs/spark-webinterfaces.png differ diff --git a/orders-streaming-analytics/pyproject.toml b/orders-streaming-analytics/pyproject.toml new file mode 100644 index 0000000..6a55380 --- /dev/null +++ b/orders-streaming-analytics/pyproject.toml @@ -0,0 +1,9 @@ +[project] +name = "gcp-spanner-boost-streaming-demo" +version = "0.1.0" +description = "Spanner Data Boost streaming demo, driven end-to-end from demo.ipynb" +readme = "README.md" +requires-python = ">=3.12" +dependencies = [ + "jupyterlab>=4.0", +] diff --git a/orders-streaming-analytics/src/java-tpch-stream-generator/.dockerignore b/orders-streaming-analytics/src/java-tpch-stream-generator/.dockerignore new file mode 100644 index 0000000..b3a2895 --- /dev/null +++ b/orders-streaming-analytics/src/java-tpch-stream-generator/.dockerignore @@ -0,0 +1,8 @@ +.git +.gitignore +.idea +.vscode +target +*.iml +.DS_Store +README.md \ No newline at end of file diff --git a/orders-streaming-analytics/src/java-tpch-stream-generator/Dockerfile b/orders-streaming-analytics/src/java-tpch-stream-generator/Dockerfile new file mode 100644 index 0000000..048fa8e --- /dev/null +++ b/orders-streaming-analytics/src/java-tpch-stream-generator/Dockerfile @@ -0,0 +1,51 @@ +# syntax=docker/dockerfile:1 + +# Copyright 2025 Google LLC +# +# Licensed under the Apache License, Version 2.0 (the "License"); +# you may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +FROM maven:3.9-eclipse-temurin-25 AS builder + +WORKDIR /workspace + +COPY pom.xml . +RUN mvn --batch-mode dependency:go-offline + +COPY src ./src + +RUN mvn --batch-mode clean package -DskipTests + + +FROM eclipse-temurin:25-jre + +WORKDIR /app + +RUN groupadd \ + --system \ + --gid 10001 \ + appgroup \ + && useradd \ + --system \ + --uid 10001 \ + --gid appgroup \ + --no-create-home \ + appuser + +COPY --from=builder \ + --chown=appuser:appgroup \ + /workspace/target/tpch-generator.jar \ + /app/tpch-generator.jar + +USER appuser + +ENTRYPOINT ["java", "-jar", "/app/tpch-generator.jar"] diff --git a/orders-streaming-analytics/src/java-tpch-stream-generator/README.md b/orders-streaming-analytics/src/java-tpch-stream-generator/README.md new file mode 100644 index 0000000..299f831 --- /dev/null +++ b/orders-streaming-analytics/src/java-tpch-stream-generator/README.md @@ -0,0 +1,82 @@ +# Java TPCH Stream Generator + +- A Kafka producer that publishes TPC-H Orders in the official pipe-delimited format. +- Supports flexible message constraints: + - **Max messages:** Stop producing and exit after a specific message count. + - **Target throughput:** Throttling in messages/second. + - **Date filtering & relative windows:** Limit records to a date range (`--date-from`, `--date-to`) using ISO dates (`YYYY-MM-DD`) or relative expressions (e.g., `"1 month ago"`, `"today"`). + - **Automatic date rebasing:** Dynamically scales historical TPC-H dates (1992–1998) into modern date windows automatically. + +```bash +Usage: tpch-generator.jar [-hV] [--local] --bootstrap-servers= + [--date-from=] [--date-to=] + --max-messages= --target-throughput= + --topic= +TPC-H Orders Kafka Producer + --bootstrap-servers= + Kafka bootstrap servers + --date-from= + Start date range (e.g., '1995-01-01', '1 month ago', + '7 days ago') + --date-to= + End date range (e.g., '1995-12-31', 'today'). Defaults + to today if date-from is specified. + -h, --help Show this help message and exit. + --local Use local Kafka without authentication + --max-messages= + Maximum number of messages to publish + --target-throughput= + Target throughput in messages per second. Use 0 for + unlimited. + --topic= Kafka topic + -V, --version Print version information and exit. + +``` + +--- + +### Date Filtering & Rebasing + +The `--date-from` and `--date-to` options accept both **ISO dates** (`YYYY-MM-DD`) and **relative expressions**: + +* **Relative formats:** `"1 day ago"`, `"2 weeks ago"`, `"1 month ago"`, `"5 years ago"`, `"today"`, `"yesterday"`. +* **Default behavior:** If `--date-from` is provided without `--date-to`, `--date-to` defaults to `today`. + +#### How Date Processing Works Automatically: + +1. **Native Filtering:** If your date range falls strictly within the TPC-H benchmark bounds (**`1992-01-01` to `1998-08-02**`), messages are filtered without altering the dataset timestamps. +2. **Automatic Rebasing:** If your requested date range extends outside the 1992–1998 window (such as relative ranges like `--date-from="1 month ago"`), the producer automatically projects and scales the historical TPC-H timestamps directly into your target window. + +--- + +### Usage Docker + +#### Relative Modern Date Window (Auto-Rebased) + +```bash +docker build -t tpch-generator:local . + +docker run --rm tpch-generator:local \ + --local \ + --bootstrap-servers=localhost:9092 \ + --max-messages=1000 \ + --target-throughput=100 \ + --topic=tpch-orders \ + --date-from="1 month ago" \ + --date-to="today" + +``` + +#### Historical Date Range (Native TPC-H Filtering) + +```bash +docker run --rm tpch-generator:local \ + --local \ + --bootstrap-servers=localhost:9092 \ + --max-messages=1000 \ + --target-throughput=100 \ + --topic=tpch-orders \ + --date-from="1995-01-01" \ + --date-to="1995-06-30" + +``` diff --git a/orders-streaming-analytics/src/java-tpch-stream-generator/pom.xml b/orders-streaming-analytics/src/java-tpch-stream-generator/pom.xml new file mode 100644 index 0000000..975265b --- /dev/null +++ b/orders-streaming-analytics/src/java-tpch-stream-generator/pom.xml @@ -0,0 +1,93 @@ + + + + + 4.0.0 + + com.google.tpch + java-tpch-stream-generator + 1.0-SNAPSHOT + + java-tpch-stream-generator + + + UTF-8 + 25 + + + + + io.trino.tpch + tpch + 1.4 + + + org.apache.kafka + kafka-clients + 4.3.1 + + + info.picocli + picocli + 4.7.7 + + + + com.google.cloud.hosted.kafka + managed-kafka-auth-login-handler + 1.0.5 + + + + + tpch-generator + + + + org.apache.maven.plugins + maven-compiler-plugin + 3.14.1 + + + org.apache.maven.plugins + maven-shade-plugin + 3.6.2 + + + package + + shade + + + false + + + + com.google.tpch.Main + + + + + + + + + diff --git a/orders-streaming-analytics/src/java-tpch-stream-generator/src/main/java/com/google/tpch/Main.java b/orders-streaming-analytics/src/java-tpch-stream-generator/src/main/java/com/google/tpch/Main.java new file mode 100644 index 0000000..bffa849 --- /dev/null +++ b/orders-streaming-analytics/src/java-tpch-stream-generator/src/main/java/com/google/tpch/Main.java @@ -0,0 +1,485 @@ +/* + * Copyright 2025 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package com.google.tpch; + +import io.trino.tpch.Order; +import io.trino.tpch.TpchTable; + +import java.time.LocalDate; +import java.time.format.DateTimeParseException; +import java.util.Iterator; +import java.util.Properties; +import java.util.concurrent.Callable; +import java.util.concurrent.atomic.AtomicReference; +import java.util.regex.Matcher; +import java.util.regex.Pattern; + +import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; +import org.apache.kafka.clients.producer.ProducerRecord; +import org.apache.kafka.common.serialization.StringSerializer; +import picocli.CommandLine; +import picocli.CommandLine.Command; +import picocli.CommandLine.Option; +import picocli.CommandLine.ParameterException; + +@Command( + name = "app", + mixinStandardHelpOptions = true, + description = "TPC-H Orders Kafka Producer" +) +public final class Main implements Callable { + + private static final LocalDate TPCH_MIN_DATE = LocalDate.of(1992, 1, 1); + private static final LocalDate TPCH_MAX_DATE = LocalDate.of(1998, 8, 2); + private static final long TPCH_MIN_EPOCH = TPCH_MIN_DATE.toEpochDay(); + private static final long TPCH_SPAN_DAYS = TPCH_MAX_DATE.toEpochDay() - TPCH_MIN_EPOCH; + + @Option( + names = "--local", + description = "Use local Kafka without authentication" + ) + private boolean localMode; + + @Option( + names = "--max-messages", + required = true, + description = "Maximum number of messages to publish" + ) + private long maxMessages; + + @Option( + names = "--target-throughput", + required = true, + description = "Target throughput in messages per second. Use 0 for unlimited." + ) + private int targetThroughput; + + @Option( + names = "--bootstrap-servers", + required = true, + description = "Kafka bootstrap servers" + ) + private String bootstrapServers; + + @Option( + names = "--topic", + required = true, + description = "Kafka topic" + ) + private String topic; + + @Option( + names = "--date-from", + description = "Start date range (e.g., '1995-01-01', '1 month ago', '7 days ago')" + ) + private String dateFromRaw; + + @Option( + names = "--date-to", + description = "End date range (e.g., '1995-12-31', 'today'). Defaults to today if date-from is specified." + ) + private String dateToRaw; + + public static void main(String[] args) { + int exitCode = new CommandLine(new Main()).execute(args); + System.exit(exitCode); + } + + @Override + public Integer call() { + validateArguments(); + + LocalDate dateFrom = parseDateOption(dateFromRaw, "--date-from"); + LocalDate dateTo = parseDateOption(dateToRaw, "--date-to"); + + // Fill defaults if only one bound is provided + if (dateFrom != null && dateTo == null) { + dateTo = LocalDate.now(); + } else if (dateFrom == null && dateTo != null) { + dateFrom = TPCH_MIN_DATE; + } + + if (dateFrom != null && dateTo != null && dateFrom.isAfter(dateTo)) { + throw new ParameterException( + new CommandLine(this), + "--date-from (" + dateFrom + ") cannot be after --date-to (" + dateTo + ")" + ); + } + + // Auto-detect rebasing: enable if either date falls outside standard TPC-H 1992-1998 bounds + boolean rebaseDates = dateFrom != null && dateTo != null && ( + dateFrom.isBefore(TPCH_MIN_DATE) || dateFrom.isAfter(TPCH_MAX_DATE) || + dateTo.isBefore(TPCH_MIN_DATE) || dateTo.isAfter(TPCH_MAX_DATE) + ); + + run(new Configuration( + localMode, + maxMessages, + targetThroughput, + bootstrapServers, + topic, + dateFrom, + dateTo, + rebaseDates + )); + + return 0; + } + + private LocalDate parseDateOption(String raw, String optionName) { + if (raw == null || raw.isBlank()) { + return null; + } + try { + return parseDateExpression(raw); + } catch (IllegalArgumentException e) { + throw new ParameterException( + new CommandLine(this), + "Invalid value for " + optionName + ": " + e.getMessage(), + e + ); + } + } + + private static LocalDate parseDateExpression(String input) { + String trimmed = input.trim().toLowerCase(); + + if ("today".equals(trimmed) || "now".equals(trimmed)) { + return LocalDate.now(); + } + if ("yesterday".equals(trimmed)) { + return LocalDate.now().minusDays(1); + } + + Pattern pattern = Pattern.compile("^(\\d+)\\s+(day|week|month|year)s?\\s+ago$"); + Matcher matcher = pattern.matcher(trimmed); + + if (matcher.matches()) { + long amount = Long.parseLong(matcher.group(1)); + String unit = matcher.group(2); + + return switch (unit) { + case "day" -> LocalDate.now().minusDays(amount); + case "week" -> LocalDate.now().minusWeeks(amount); + case "month" -> LocalDate.now().minusMonths(amount); + case "year" -> LocalDate.now().minusYears(amount); + default -> throw new IllegalArgumentException("Unsupported unit: " + unit); + }; + } + + try { + return LocalDate.parse(trimmed); + } catch (DateTimeParseException e) { + throw new IllegalArgumentException( + "Could not parse date '" + input + "'. Expected format YYYY-MM-DD or relative like '1 month ago'." + ); + } + } + + private void validateArguments() { + CommandLine commandLine = new CommandLine(this); + + if (maxMessages <= 0) { + throw new ParameterException( + commandLine, + "--max-messages must be greater than zero." + ); + } + + if (targetThroughput < 0) { + throw new ParameterException( + commandLine, + "--target-throughput cannot be negative." + ); + } + + if (bootstrapServers.isBlank()) { + throw new ParameterException( + commandLine, + "--bootstrap-servers cannot be blank." + ); + } + + if (topic.isBlank()) { + throw new ParameterException( + commandLine, + "--topic cannot be blank." + ); + } + } + + private static void run(Configuration configuration) { + Properties properties = createProducerProperties(configuration); + + System.out.printf( + "Configuration:%n" + + " Kafka mode: %s%n" + + " Bootstrap server: %s%n" + + " Topic: %s%n" + + " Max messages: %d%n" + + " Throughput: %s%n" + + " Date Range: %s to %s%n" + + " Date Strategy: %s%n", + configuration.localMode() + ? "local, unauthenticated" + : "GCP Managed Kafka", + configuration.bootstrapServers(), + configuration.topic(), + configuration.maxMessages(), + configuration.targetThroughput() == 0 + ? "unlimited" + : configuration.targetThroughput() + " msg/sec", + configuration.dateFrom() != null ? configuration.dateFrom() : "unconstrained", + configuration.dateTo() != null ? configuration.dateTo() : "unconstrained", + configuration.rebaseDates() ? "Rebased automatically" : "Filtered natively" + ); + + AtomicReference sendFailure = new AtomicReference<>(); + + try (KafkaProducer producer = new KafkaProducer<>(properties)) { + + double scaleFactor = 0.1; + + Iterator orderIterator = + TpchTable.ORDERS + .createGenerator(scaleFactor, 1, 1) + .iterator(); + + long messagesSent = 0; + long startTimeNs = System.nanoTime(); + + while (messagesSent < configuration.maxMessages() && orderIterator.hasNext()) { + + Exception previousFailure = sendFailure.get(); + if (previousFailure != null) { + throw new IllegalStateException( + "A Kafka send operation failed.", + previousFailure + ); + } + + Order order = orderIterator.next(); + LocalDate originalOrderDate = LocalDate.ofEpochDay(order.orderDate()); + String recordPayload = order.toLine(); + + if (configuration.dateFrom() != null && configuration.dateTo() != null) { + if (configuration.rebaseDates()) { + // Project TPC-H epoch progress linearly into [dateFrom, dateTo] + long targetSpanDays = configuration.dateTo().toEpochDay() - configuration.dateFrom().toEpochDay(); + double progress = (double) (order.orderDate() - TPCH_MIN_EPOCH) / TPCH_SPAN_DAYS; + long rebasedEpochDay = configuration.dateFrom().toEpochDay() + (long) (progress * targetSpanDays); + + LocalDate rebasedDate = LocalDate.ofEpochDay(rebasedEpochDay); + recordPayload = replaceOrderDateInPayload(recordPayload, rebasedDate); + } else { + // Native filtering for dates within 1992-1998 + if (originalOrderDate.isBefore(configuration.dateFrom()) || + originalOrderDate.isAfter(configuration.dateTo())) { + continue; + } + } + } + + // Bit-reverse the sequential orderKey to prevent write hotspots in distributed targets + long distributedOrderKey = bitReverse(order.orderKey()); + + // Replace O_ORDERKEY (field 0) in the pipe-delimited payload + recordPayload = replaceOrderKeyInPayload(recordPayload, distributedOrderKey); + + + ProducerRecord record = + new ProducerRecord<>( + configuration.topic(), + String.valueOf(distributedOrderKey), + recordPayload + ); + + producer.send( + record, + (metadata, exception) -> { + if (exception != null) { + sendFailure.compareAndSet( + null, + exception + ); + } + }); + + messagesSent++; + + throttle( + messagesSent, + configuration.targetThroughput(), + startTimeNs + ); + } + + producer.flush(); + + Exception finalFailure = sendFailure.get(); + if (finalFailure != null) { + throw new IllegalStateException( + "At least one Kafka record could not be published.", + finalFailure + ); + } + + long elapsedNs = System.nanoTime() - startTimeNs; + double elapsedSeconds = elapsedNs / 1_000_000_000.0; + double actualThroughput = elapsedSeconds == 0.0 + ? messagesSent + : messagesSent / elapsedSeconds; + + System.out.printf( + "Completed successfully. Published %d records in %.2f seconds (%.2f msg/sec).%n", + messagesSent, + elapsedSeconds, + actualThroughput + ); + } + } + + /** + * Bit-reverses a 64-bit integer while preserving it as a strictly positive number (> 0). + * This mimics Cloud Spanner's bit_reversed_positive sequence behavior. + */ + private static long bitReverse(long value) { + return Long.reverse(value) >>> 1; + } + + private static String replaceOrderKeyInPayload(String rawPayload, long newOrderKey) { + String[] fields = rawPayload.split("\\|", -1); + if (fields.length > 0) { + fields[0] = Long.toString(newOrderKey); // Field index 0 is O_ORDERKEY + } + return String.join("|", fields); + } + + private static String replaceOrderDateInPayload(String rawPayload, LocalDate newDate) { + String[] fields = rawPayload.split("\\|", -1); + if (fields.length > 4) { + fields[4] = newDate.toString(); // Field index 4 is O_ORDERDATE + } + return String.join("|", fields); + } + + private static Properties createProducerProperties(Configuration configuration) { + Properties properties = new Properties(); + + properties.put( + ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, + configuration.bootstrapServers() + ); + + properties.put( + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, + StringSerializer.class.getName() + ); + + properties.put( + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, + StringSerializer.class.getName() + ); + + properties.put( + ProducerConfig.LINGER_MS_CONFIG, + "50" + ); + + properties.put( + ProducerConfig.ACKS_CONFIG, + "all" + ); + + properties.put( + ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, + "true" + ); + + properties.put( + ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, + "120000" + ); + + properties.put( + ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, + "30000" + ); + + if (!configuration.localMode()) { + configureGoogleManagedKafkaAuthentication(properties); + } + + return properties; + } + + private static void configureGoogleManagedKafkaAuthentication(Properties properties) { + properties.put("security.protocol", "SASL_SSL"); + properties.put("sasl.mechanism", "OAUTHBEARER"); + properties.put( + "sasl.login.callback.handler.class", + "com.google.cloud.hosted.kafka.auth.GcpLoginCallbackHandler" + ); + properties.put( + "sasl.jaas.config", + "org.apache.kafka.common.security.oauthbearer." + + "OAuthBearerLoginModule required;" + ); + } + + private static void throttle( + long messagesSent, + int targetThroughput, + long startTimeNs) { + if (targetThroughput == 0) { + return; + } + + long expectedElapsedNs = (messagesSent * 1_000_000L) / targetThroughput * 1_000L; + long actualElapsedNs = System.nanoTime() - startTimeNs; + long delayNs = expectedElapsedNs - actualElapsedNs; + + if (delayNs <= 0) { + return; + } + + long delayMillis = delayNs / 1_000_000L; + int additionalNanos = (int) (delayNs % 1_000_000L); + + try { + Thread.sleep(delayMillis, additionalNanos); + } catch (InterruptedException exception) { + Thread.currentThread().interrupt(); + + throw new IllegalStateException( + "Producer was interrupted while throttling.", + exception + ); + } + } + + private record Configuration( + boolean localMode, + long maxMessages, + int targetThroughput, + String bootstrapServers, + String topic, + LocalDate dateFrom, + LocalDate dateTo, + boolean rebaseDates) { + } +}