Skip to content

Project Aegis: Enterprise streaming analytics and autonomous agentic operations platform

Integrating Managed Open-Source Technologies with Google Cloud Infrastructure Demonstrating Managed Apache Kafka, Apache Spark (Dataproc Serverless with C++ Velox acceleration), Cloud Bigtable, BigQuery, and GEAP (Gemini Enterprise Agent Platform).

Overview and project goal

Project Aegis is an enterprise-grade reference architecture and interactive demonstration platform built for Google Cloud Customer Engineers (CEs) and Solution Architects. It demonstrates how modern enterprises can seamlessly combine Managed Open-Source Software (MOSS)—such as Apache Kafka and Apache Spark—with Google Cloud Services—including Cloud Bigtable, BigQuery, GEAP (Gemini Enterprise Agent Platform), and Cloud Monitoring and Logging.

The platform ingests high-volume industrial IoT (IIoT) telemetry from 15 simulated assets, performs sub-second windowed analytics, detects thermal and compute anomalies, routes alerts through Model Armor security guardrails, and executes autonomous Gemini 2.5 Flash Chain-of-Thought Root Cause Analysis (RCA) and mitigation.

Architecture diagram

flowchart TD
    subgraph Ingestion ["Managed Open-Source Ingestion Broker"]
        A[Synthetic Telemetry Generator] -->|Streaming JSON Events| B[Managed Apache Kafka: telemetry-raw]
        B --> C[Dataproc Serverless PySpark Engine]
    end

    subgraph Processing ["Vectorized Spark Processing Engine"]
        C -->|C++ Velox / Gluten Acceleration| D[10-Second Tumbling Window Aggregator]
        D -->|Sub-millisecond State Writes| E[(Cloud Bigtable: aegis-bigtable)]
        D -->|Streaming Analytics Sink| F[(BigQuery: analytics.telemetry_events)]
        C -.->|DATAPROC_LINEAGE_ENABLED| G[Knowledge Graph / OpenLineage]
    end

    subgraph Agentic ["GEAP Agentic Operations and Security Shield"]
        E -->|SSE Stream / Alert| H[HUD Backend FastAPI Service]
        H -->|Telemetry Payload| I[Model Armor Security Shield]
        I -->|Sanitized Prompt| J[GEAP: Gemini 2.5 Flash Agent]
        J -->|Chain-of-Thought RCA and Action| K[Mitigation Engine]
        J -->|Token Spend and ROI Logging| L[(BigQuery: analytics.rca_events)]
        J -->|Metrics and Logs| M[Cloud Monitoring and Cloud Logging]
    end

    subgraph Dashboard ["Executive Command HUD"]
        H -->|Server-Sent Events| N[HUD Next.js Frontend Dashboard]
        N -->|Interactive Control| H
    end

Business value and enterprise ROI

Business Pillar Value Proposition Measurable Impact
Prevented Downtime Autonomous AI agent detects thermal anomalies and issues mitigation commands before hardware shutdown. Reduces unmitigated failure costs (~$5,000 per incident) to near-zero.
C++ Vectorized Efficiency Dataproc Serverless utilizes the C++ Lightning Engine (Velox/Gluten) to eliminate JVM garbage collection pauses. Up to 4x execution speedup and 60% lower compute cost compared to standard Spark.
Financial Governance (Tokenomics) Every LLM call tracks exact input/output token counts, inference costs, and averted downtime value in BigQuery. Demonstrates >50,000% ROI per incident (for example, $0.0001 USD inference cost compared to $5,000 averted downtime).
Enterprise Security and Compliance Model Armor filters incoming prompts for PII and neutralizes prompt injection vectors before LLM execution. Prevents adversarial prompt attacks and sensitive data leakage.
Data Lineage and Provenance Dataproc integration with OpenLineage maps end-to-end data provenance into Google Cloud Knowledge Graph. Satisfies regulatory compliance and audit requirements out of the box.

Technical architecture and component breakdown

1. Ingestion broker: Managed Apache Kafka

  • Resource: google_managed_kafka_cluster (aegis-kafka-cluster) and google_managed_kafka_topic (telemetry-raw).
  • Role: Serves as the open-source messaging backbone for real-time telemetry ingestion without requiring customer-managed broker VMs.

2. Processing engine: Dataproc Serverless PySpark (Velox engine)

  • Script: data-ingestion/src/aegis_etl.py
  • Role: Consumes streaming Kafka events using PySpark Structured Streaming (org.apache.spark:spark-sql-kafka-0-10_2.12). Computes 10-second tumbling window aggregations with C++ Velox vectorized columnar acceleration.

3. Operational state database: Cloud Bigtable

  • Resource: google_bigtable_instance (aegis-bigtable), table telemetry_metrics.
  • Role: Dual-sink operational database storing live rolling averages, status flags (OK, WARNING, CRITICAL), and sub-second metrics for real-time dashboard visualization.

4. Data warehouse and tokenomics: BigQuery

  • Resource: google_bigquery_dataset (analytics), tables telemetry_events (partitioned by day) and rca_events.
  • Role: Persistent analytical data warehouse for trend SQL queries, historical reporting, and LLM token spend auditing.

5. AI agentic platform: GEAP and Gemini 2.5 Flash

  • Module: agent-service/src/agent.py and security.py
  • Role: Powers the cognitive operator. Sanitizes incoming payloads using Model Armor, executes Gemini 2.5 Flash Root Cause Analysis (RCA), and records tokenomics metrics.

6. Observability and Command HUD

  • Backend: FastAPI Python service (hud/backend/src/main.py) providing Server-Sent Events (SSE) and Cloud Monitoring and Logging integration.
  • Frontend: Next.js React Dashboard (hud/frontend/src/app/page.tsx) displaying live asset grids, real-time metrics, anomaly controls, and agent mitigation cards.

Intended audience

  • Chief Technology Officers (CTOs) and VPs of Engineering: Evaluate open-source integration (Kafka/Spark) on Google Cloud infrastructure.
  • VPs of Data Infrastructure and Data Architects: Inspect Dataproc Serverless PySpark execution, C++ Velox vectorization, and Bigtable/BigQuery dual-sink patterns.
  • Chief Security Officers (CSOs) and AI Leads: Review Model Armor prompt injection defense, PII masking, and GEAP tokenomics governance.
  • Customer Engineers (CEs) and Solution Architects: Deliver interactive 10-minute executive walkthroughs using DEMO_GUIDE.md.

Quick start and deployment guide

1. Provision infrastructure with modular Terraform

cd terraform/stacks/oss
terraform init
terraform apply -auto-approve

The modular Terraform structure includes:

  • terraform/modules/base_platform — Shared VPC, IAM, Cloud Bigtable, BigQuery, Model Armor, GEAP, and Cloud Run services
  • terraform/stacks/oss — Managed Apache Kafka and Dataproc Serverless PySpark stack
  • terraform/stacks/first-party — Cloud Pub/Sub and Cloud Dataflow Apache Beam stack
  • terraform/stacks/low-code — Cloud Pub/Sub and BigQuery Continuous Queries stack

2. Launch telemetry simulator

# Start telemetry stream using the HTTP API:
curl -X POST "https://hud-backend-xxxx-uc.a.run.app/api/start-stream" -H "Content-Type: application/json" -d '{"rate_msgs_per_sec": 100}'

# Or run simulator locally:
uvicorn main:app --app-dir telemetry-simulator/src --port 8080

3. Start Dataproc PySpark streaming pipeline

# Start Dataproc Serverless PySpark streaming pipeline through the HUD Backend:
curl -X POST "https://hud-backend-xxxx-uc.a.run.app/api/pipeline/start" -H "Content-Type: application/json"

4. Open Command HUD and interactive demo

Access the HUD Dashboard by running the automated authentication proxy helper script at project root:

./RUN_PROXY.sh

Note

If the cloud-run-proxy gcloud component is not yet installed on your system, gcloud prompts to install it. Type Y and press Enter. The script waits until the local proxy starts listening on http://localhost:8080 before opening your web browser.

Refer to DEMO_GUIDE.md for the complete interactive presentation script.