Skip to content

Repository files navigation

Prism - Hadoop Multi-Cluster Architecture

Enterprise-grade dual-cluster Hadoop architecture with YARN queue management, multi-tenant processing, real-time analytics via Impala, and unified data governance through Collibra.

Architecture Overview

Data Sources --> A-Cluster (YARN/HDFS/Hive) --> B-Cluster (YARN/HDFS/Hive) --> API/Output
                  (Ingestion & ETL)     ^         (Analytics & ML)
                                        |
                              Direct HDFS File Read

Key Components

Component Role
A-Cluster Heavy data ingestion and batch/stream processing
B-Cluster Shared multi-tenant application processing, analytics, and ML
Collibra Unified data governance, catalog, and lineage
Impala Real-time SQL-on-Hadoop query engine
Apache NiFi 3-node HA data ingestion and routing
Apache Kafka Event streaming platform
Apache Ranger Fine-grained access control and authorization

Dual-Cluster Design

A-Cluster - Ingestion Zone

Handles heavy input data streams including batch file uploads and API-based ingestion (REST/Kafka).

YARN Queue Configuration:

Functional ID Type Purpose
func-id-A1 Dedicated Critical ingestion jobs
func-id-A2 Dedicated Time-sensitive processing
func-id-A3 Batch Large batch processing

VM Farm Allocation:

  • Dedicated VMs: VM-A1, VM-A2 (guaranteed resources)
  • Shared VMs: VM-A3, VM-A4, VM-A5 (elastic capacity)
  • Batch VMs: VM-A6 (low-priority workloads)

HDFS Storage:

  • NameNode with HA (JournalNodes)
  • DataNode 1-3 with replication factor 3
  • Encryption zones for sensitive data

Hive Databases (Section A):

/user/hive/warehouse/section_a/
  raw_events_db/      -- Raw incoming data
  staging_db/         -- Data validation staging
  ingestion_db/       -- Ingestion tracking
  archive_db/         -- Historical data archive
  external_tables/    -- External data references

B-Cluster - Processing Zone

Shared multi-tenant environment for application-specific compute, analytics, and ML workloads.

YARN Queue Configuration (Multi-Tenant):

Functional ID Team/Purpose
func-id-B1 Application Team 1
func-id-B2 Application Team 2
func-id-B3 Analytics Team
func-id-B4 ML/AI Workloads
func-id-B5 Reporting Team

Resource Isolation:

  • Fair Scheduler with preemption
  • Capacity limits per functional ID
  • Priority-based scheduling

Hive Databases (Section B):

/user/hive/warehouse/section_b/
  analytics_db/       -- Analytics results
  processed_db/       -- Processed/transformed data
  aggregates_db/      -- Aggregated metrics
  ml_features_db/     -- ML feature store
  reports_db/         -- Report datasets
  curated_db/         -- Curated data products
  exports_db/         -- Export-ready data

Cross-Cluster Data Flow

+-------------------+          +-------------------+
|   A-Cluster       |          |   B-Cluster       |
|   HDFS            |  ------> |   HDFS            |
|   (Source)        |  Direct  |   (Target)        |
|                   |  Read    |                   |
+-------------------+          +-------------------+

Integration Patterns:

  • Direct HDFS Read: B-Cluster reads A-Cluster HDFS files directly
  • No data duplication: Reference-based access via HDFS Federation (ViewFS)
  • Batch transfers: DistCp for scheduled data movement
  • Streaming: Kafka Connect for real-time event streams

Protocol Connectivity

Protocol Port Connection
HTTPS :443 External sources to ingestion layer
SFTP :22 Batch file uploads
Kafka :9092 Event streaming between components
HDFS RPC :8020 Cross-cluster HDFS communication
Hive Thrift :9083 Hive Metastore access
Impala JDBC :21050 Real-time SQL queries
HiveServer2 JDBC :10000 Hive query execution
YARN RPC :8032 Resource management
Kerberos :88 Authentication
Ranger REST :6080 Policy management

Impala Real-Time Query Engine

Component Function
Impala Daemon 1-3 Distributed query execution nodes
Catalog Service Metadata caching from Hive Metastore
StateStore Cluster membership and health management

Capabilities:

  • SQL-on-Hadoop with sub-second query response
  • Direct Hive Metastore integration
  • Parquet/ORC format optimization
  • Real-time dashboards, ad-hoc analytics, and API-driven queries

Collibra Data Governance

Centralized data governance layer spanning both clusters.

Component Description
Metadata Definitions Common data element definitions
Schema Registry Versioned schema management
Data Lineage Track data flow and transformations
Business Glossary Business term definitions
Data Quality Rules Quality validation rules
API Catalog API documentation and metadata
Access Policies Security and access control
Data Stewardship Ownership and accountability

Functional ID Management

Functional IDs provide workload isolation, resource quota management, cost allocation, and access control.

functional_id: func-id-B1
queue_name: app_team_1
capacity:
  min: 10%
  max: 40%
resource_limits:
  memory: 256GB
  vcores: 64
priority: 2
preemption: true
users:
  - app_team_1_svc
  - analyst_1

YARN Queue XML Example:

<queue name="root">
  <queue name="a_cluster">
    <queue name="func_id_a1">
      <capacity>30</capacity>
      <maximum-capacity>50</maximum-capacity>
      <user-limit-factor>2</user-limit-factor>
    </queue>
    <queue name="func_id_a2">
      <capacity>30</capacity>
      <maximum-capacity>50</maximum-capacity>
    </queue>
    <queue name="func_id_a3">
      <capacity>40</capacity>
      <maximum-capacity>60</maximum-capacity>
    </queue>
  </queue>
</queue>

Metadata Schemas

File Upload Ingestion

{
  "file_name": "transactions_2024.parquet",
  "file_size": 1073741824,
  "file_format": "PARQUET",
  "partition_keys": ["year", "month", "day"],
  "schema_version": "v2.1",
  "compression_type": "SNAPPY",
  "upload_timestamp": "2024-01-15T10:30:00Z",
  "source_system": "CORE_BANKING",
  "functional_id": "func-id-A1",
  "data_classification": "CONFIDENTIAL",
  "retention_policy": "7_YEARS",
  "owner_team": "DATA_ENGINEERING",
  "checksum_md5": "abc123..."
}

API-Based Ingestion

{
  "api_endpoint": "/api/v1/transactions",
  "payload_schema": "transaction_event_v2",
  "event_type": "TRANSACTION_CREATED",
  "correlation_id": "uuid-1234-5678",
  "batch_id": "BATCH-2024-001",
  "ingestion_timestamp": "2024-01-15T10:30:00Z",
  "source_functional_id": "func-id-A1",
  "target_hive_table": "raw_events_db.transactions",
  "data_classification": "PII",
  "record_count": 50000,
  "processing_priority": "HIGH"
}

API & Output Layer

Service Purpose
REST API Gateway Synchronous data access
GraphQL Endpoint Flexible queries
Data Export Service Scheduled exports

Consumers:

Consumer Data Source
BI Tools (Tableau) Impala queries
Dashboards Real-time Impala
Reports Scheduled exports

Security Architecture

Authentication

  • Kerberos for cluster-wide authentication
  • LDAP/Active Directory integration
  • Service accounts per functional ID

Authorization

  • Apache Ranger for fine-grained access control
  • Collibra-managed policies synced to Ranger
  • Row/Column level security in Hive

Encryption

  • HDFS encryption zones for data at rest
  • TLS for all data in transit
  • Key management with Hadoop KMS

API Security

  • OAuth2/JWT authentication
  • Role-based access control (RBAC)
  • Data masking for PII fields

High Availability

Per-Cluster HA

  • YARN ResourceManager: Active-Standby failover with ZooKeeper
  • HDFS NameNode: HA with JournalNodes
  • Hive Metastore: Load-balanced instances

B-Cluster Specific

  • Impala: Multiple daemons with load balancing
  • Catalog/StateStore: Automatic failover

Cross-Cluster

  • Federated HDFS namespace (ViewFS)
  • Backup paths for data access
  • DR replication strategy

Monitoring & Operations

Stack: Prometheus + Grafana dashboards, YARN ResourceManager UI, HDFS NameNode UI, Impala Query Profile

Metric Threshold
Queue utilization < 80%
HDFS capacity < 75%
Query latency (P99) < 5s
Job failure rate < 1%

Alerting: PagerDuty, Slack notifications, Email escalation


Hive External Table Example

CREATE EXTERNAL TABLE section_a.raw_events (
    event_id STRING,
    event_type STRING,
    payload STRING,
    created_at TIMESTAMP
)
PARTITIONED BY (year INT, month INT, day INT)
STORED AS PARQUET
LOCATION 'hdfs://a-cluster/data/raw_events';

Diagrams

File Description
hadoop-architecture.drawio Ecosystem-level architecture diagram with 5 layers
hadoop-layered-architecture.drawio Component layered architecture (Page 1) + Data flow sequence diagram (Page 2)
hadoop-architecture.svg SVG export of the ecosystem diagram
hadoop-architecture.pptx PowerPoint presentation export
hadoop-architecture-slides.md Full architecture documentation in slide format

Architecture Highlights

  1. Dual-cluster design separates ingestion from processing for workload isolation
  2. YARN queues with functional IDs enable multi-tenancy with resource guarantees
  3. Direct HDFS read minimizes data movement between clusters
  4. Collibra governance ensures data quality, lineage, and compliance
  5. Impala enables real-time sub-second analytics
  6. Metadata schemas standardize file and API ingestion patterns
  7. Petabyte-scale with elastic resource allocation per team

About

Hadoop Multi-Cluster Architecture - Dual-cluster design with YARN queue management, multi-tenant processing, Impala real-time analytics, and Collibra data governance. Includes component layered diagrams, sequence diagrams, and full architecture documentation.

Resources

Stars

Watchers

Forks

Releases

Packages

Contributors

Languages