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) andgoogle_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), tabletelemetry_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), tablestelemetry_events(partitioned by day) andrca_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.pyandsecurity.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¶
The modular Terraform structure includes:
terraform/modules/base_platform— Shared VPC, IAM, Cloud Bigtable, BigQuery, Model Armor, GEAP, and Cloud Run servicesterraform/stacks/oss— Managed Apache Kafka and Dataproc Serverless PySpark stackterraform/stacks/first-party— Cloud Pub/Sub and Cloud Dataflow Apache Beam stackterraform/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:
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.