Operations 26 min read

Building an Automated Root Cause Analysis System for Alibaba Cloud Big Data Operations

This article details a comprehensive root cause analysis system for Alibaba Cloud big data platforms, covering a five-layer architecture, three-dimensional analysis (cluster health, task execution, data pipelines), knowledge graph-driven automation, and AI-enhanced future directions to shift from reactive firefighting to proactive defense.

Lakehouse Research Base
Lakehouse Research Base
Lakehouse Research Base
Building an Automated Root Cause Analysis System for Alibaba Cloud Big Data Operations

This article presents a complete construction plan for an operations root cause analysis system built on Alibaba Cloud's big data ecosystem, including EMR StarRocks, Paimon, Flink, MaxCompute, DataWorks, Kafka, and Elasticsearch. The system aims to solve core pain points such as complex component collaboration, inefficient fault localization, and delayed operational response, enabling a shift from passive firefighting to proactive defense.

1. Core Objectives and Design Principles

1.1 Core Objectives

Early Warning: Establish metric baselines for components and business chains to precisely identify hidden anomalies (e.g., slow resource creep, task duration deviations) and achieve preemptive alerts without business impact.

Precise Localization: Break monitoring silos by correlating component status, task execution, and data flow, enabling minute-level root cause identification rather than surface symptoms.

Efficient Closed Loop: Combine automated repair scripts with tiered manual response to accelerate resolution and accumulate operational knowledge, reducing costs and ensuring business continuity.

Value Upgrade: Feed root cause insights back into architecture optimization and business enablement, improving overall availability and performance of the big data platform.

1.2 Design Principles

Native Component Adaptation: Leverage built-in monitoring capabilities of Alibaba Cloud components (DataWorks task monitoring, EMR cluster monitoring, MaxCompute metadata APIs) to minimize third-party integrations and reduce complexity.

End-to-End Coverage: Cover the full lifecycle from data ingestion, processing, storage, query, to application, spanning Kafka transmission, Flink/MaxCompute processing, Paimon storage, StarRocks query, ES log storage, and DataWorks scheduling.

Intelligence-Driven: Fuse rule engines with AI algorithms to automate anomaly detection, root cause correlation, and repair recommendation, boosting operational efficiency.

Extensible and Practical: Adopt a layered architecture supporting horizontal component expansion (e.g., adding Lindorm, Hologres) and follow a phased implementation approach from simple to complex for quick wins.

2. Technical Architecture and Core Component Roles

The system adopts a five-layer architecture: Collection, Storage, Detection, Analysis, and Closed Loop, fully leveraging the Alibaba Cloud component ecosystem for end-to-end metric collection, intelligent analysis, and automated repair.

Technical architecture diagram
Technical architecture diagram
Component responsibility matrix
Component responsibility matrix

2.1 Core Component Operational Responsibilities

The responsibility matrix assigns specific monitoring and analysis duties to each component, ensuring clear ownership across the stack.

3. End-to-End Root Cause Analysis Implementation Practice

Focusing on three core dimensions—cluster health, task execution, and data pipelines—the system implements targeted root cause analysis solutions covering high-frequency operational failure scenarios.

3.1 Dimension 1: EMR Cluster Health Root Cause Analysis

EMR clusters host StarRocks, Paimon, and Flink; cluster anomalies directly impact end-to-end stability. Analysis focuses on resource levels, component collaboration, and dependent services.

3.1.1 Core Monitoring Metrics and Baseline Setting

Resource Levels: Node CPU/memory usage, disk I/O utilization, network bandwidth utilization. Baselines set using 95th percentile over 7 days, differentiated by business peak/off-peak periods (daytime query peaks, early morning migration peaks). Anomaly if continuously exceeding baseline by 20% for 5 minutes; disk I/O >85%; bandwidth >90%.

Component Status: StarRocks FE/BE liveness, Flink JobManager/TaskManager status, Paimon Catalog connection status. Baseline is normal running state. Anomaly if process exits, connection failures >3/min, FE metadata sync latency >10s.

Dependent Services: EMR-OSS/HDFS connectivity, metadata service availability, Alibaba Cloud VPC network connectivity. Baseline: 100% connection success rate, network latency <50ms. Anomaly if success rate <99%, latency >500ms, VPC port unreachable.

3.1.2 EMR StarRocks Root Cause Analysis Standardized Process (6-Step Closed Loop)

StarRocks anomalies concentrate in four scenarios: query latency/failure, import throughput drop, data inconsistency, node anomalies. The standardized process integrates Alibaba Cloud native tools for efficient localization.

Preparation: Tool and Permission Configuration

Monitoring: Enable EMR cluster monitoring (including StarRocks FE/BE dedicated metrics) and Alibaba Cloud ARMS application monitoring, ensuring metric collection frequency ≥1 minute.

Logging: Configure StarRocks FE/BE logs to Alibaba Cloud ES with structured indexing (log level, error type, task ID).

Permissions: Grant Ops accounts EMR cluster management permissions (ECS login, component operations), StarRocks SUPER USER, DataWorks task view, ARMS metric query.

Execution Flow: From Anomaly Detection to Root Cause Localization (6 Steps)

Step 1-2 diagram
Step 1-2 diagram
Step 3-4 diagram
Step 3-4 diagram
Step 5-6 diagram
Step 5-6 diagram

Typical Scenario Case: StarRocks Import Throughput Plunge (100MB/s → 20MB/s)

Step 1 Anomaly Detection: ARMS alerts on import throughput anomaly, affecting downstream report generation. Baseline comparison confirms 50% threshold breach, impacting all StarRocks import tasks.

Step 2 Cluster Inspection: FE/BE nodes alive, metadata sync latency 5s. BE node CPU 92% (baseline 60%), disk I/O 88% (baseline 40%) — resource level anomaly.

Step 3 Task Localization: SHOW LOAD shows all import tasks running, no failures. Correlating with DataWorks reveals Flink real-time import task scheduling cycle changed from 5 minutes to 1 minute, causing task surge.

Step 4 Log Analysis: ES search of BE logs finds numerous "resource limit exceeded" errors matching Flink import task IDs.

Step 5 Correlation Verification: EMR resource queue shows StarRocks and Flink sharing default queue; Flink tasks occupying 70% CPU, causing resource contention.

Step 6 Repair Execution: Restore Flink scheduling cycle via DataWorks; create dedicated resource queue for StarRocks imports (40% quota). Post-repair, import throughput recovers to 100MB/s, resource usage returns to baseline.

Operational Best Practices

Avoid direct component restarts; prioritize log and metric analysis to prevent metadata loss or data inconsistency.

Prioritize core business: during repair, safeguard core StarRocks tasks (core reports, critical imports), optionally pause non-core tasks to free resources.

Record configuration changes: document StarRocks config modifications and EMR resource queue changes with timestamps for traceability.

Regular metadata backup: daily backup of StarRocks FE metadata to prevent corruption during root cause analysis.

3.1.3 Anomaly Root Cause Correlation Logic (Case Study)

Case: StarRocks query task latency spikes, business reports timeout.

Surface Symptom: Query latency rises from 500ms to 5s, multiple report loading failures.

Metric Anomalies: EMR BE node CPU 95% (baseline 60%), disk I/O 88% (baseline 40%); concurrent Flink real-time import task volume doubled.

Correlation Analysis:

EMR monitoring shows Flink import tasks consuming 70% CPU, sharing default queue with StarRocks queries — resource contention.

DataWorks scheduling records reveal Flink import cycle mistakenly adjusted to 1 minute (baseline 5 minutes), causing dense execution.

Network and storage verified: EMR-OSS connection normal, Paimon table reads no latency, ruling out storage/network issues.

Root Cause Conclusion: Abnormal Flink import scheduling cycle adjustment led to resource contention with StarRocks queries, causing query latency.

Solution: Restore Flink task scheduling cycle via DataWorks; configure dedicated EMR resource queues for StarRocks queries and Flink imports with isolation ratio (query:import = 6:4); schedule non-core imports to off-peak early morning.

3.2 Dimension 2: Task Execution Anomaly Root Cause Analysis (DataWorks/Flink/MaxCompute)

Task anomalies are high-frequency issues; core root causes include resource shortage, dependency blocking, data anomalies, configuration errors. Leveraging DataWorks task dependency graphs and native component monitoring, a baseline-comparison, chain-tracing, data-profiling analysis flow is built.

3.2.1 Task Anomaly Types and Root Cause Mapping

Task anomaly type mapping
Task anomaly type mapping

3.2.2 Root Cause Analysis Practice (Case Study)

Case: DataWorks-scheduled MaxCompute offline processing task timeout failure.

Step 1 Baseline Comparison: DataWorks monitoring shows historical baseline runtime 3 hours; current run 6 hours incomplete, judged anomalous.

Step 2 Chain Tracing: DataWorks dependency graph shows upstream ODS data integration task completed normally. MaxCompute job monitoring reveals source table data volume tripled (100GB → 300GB).

Step 3 Data Profiling: MaxCompute metadata API shows source table added one column with massive long strings (10KB/row vs baseline 100B/row), reducing processing efficiency. Job resource quota (CPU/memory) remained at baseline, not scaled with data growth.

Root Cause Conclusion: Upstream source data volume surge + new column format anomaly, combined with insufficient task resource quota, caused timeout.

Solution: Temporarily double task resource quota via DataWorks; add data filtering rule in MaxCompute job to cleanse long strings; coordinate upstream business to optimize data format; configure DataWorks elastic scaling rule to auto-scale when data volume exceeds 2x baseline.

3.3 Dimension 3: Data Pipeline Anomaly Root Cause Analysis (Kafka→Flink→Paimon→StarRocks)

Data pipelines are the core of data flow; anomalies cause latency or loss. End-to-end monitoring across producer, transport, and consumer enables bottleneck localization and root cause analysis.

3.3.1 Key Pipeline Monitoring Metrics

Producer (Data Ingestion): Kafka produce rate, send success rate, MaxCompute export rate. Anomaly if produce rate drops 30%, send failure >1%, export rate lags plan by 20%. Responsible: Kafka Producer, MaxCompute.

Transport (Data Transfer): Kafka replica sync latency, network bandwidth utilization, OSS/Paimon write rate. Anomaly if sync latency >1s, bandwidth >90%, write rate <50% baseline. Responsible: Kafka Broker, Alibaba Cloud VPC, OSS/Paimon.

Consumer (Data Processing): Flink consume rate, StarRocks import success rate, data backlog. Anomaly if consume rate < produce rate, import failure >5%, backlog grows 50GB+/hour. Responsible: Flink Sink, StarRocks Load, DataWorks Data Integration.

3.3.2 Root Cause Analysis Practice (Case Study)

Case: Kafka→Flink→Paimon pipeline blockage, severe data backlog.

Surface Symptom: Kafka topic backlog grows 80GB in 1 hour; Flink consume rate drops from 100MB/s to 20MB/s; Paimon table write latency exceeds 2 hours.

End-to-End Metric Investigation:

Producer: Kafka produce rate stable at 100MB/s, send success 100% — producer excluded.

Transport: VPC bandwidth 40%, Kafka replica sync <100ms — transport excluded.

Consumer: Flink task CPU/memory 100%, operator checkpoint failure 80%. Paimon monitoring shows snapshot count >100 (baseline 20), snapshot cleanup job not running, degrading storage read/write performance.

Deep Analysis: ES query of Flink logs shows checkpoint failures due to Paimon write timeout. DataWorks scheduling records reveal Paimon snapshot cleanup job failed due to temporary metadata service unavailability, without auto-retry.

Root Cause Conclusion: Paimon snapshot cleanup job failure caused excessive snapshots, degrading storage performance, triggering Flink write timeouts and checkpoint failures, ultimately blocking the pipeline.

Solution: Manually clean Paimon historical snapshots to restore storage performance; trigger Flink task rerun via DataWorks to accelerate consumption; configure auto-retry and alerting for Paimon snapshot cleanup; optimize metadata service HA to avoid single point of failure.

4. Engineering Implementation Guarantees

4.1 Build Alibaba Cloud Dedicated Operations Knowledge Graph

Based on Alibaba Cloud component failure cases, accumulate a knowledge graph of "Anomaly Phenomenon → Related Metrics → Root Cause → Solution → Prevention Strategy" stored and managed via Neo4j. Examples:

Phenomenon: StarRocks import throughput plunge → Metrics: BE CPU/IO, resource queue usage → Root Cause: Resource contention → Solution: Dedicated queue + off-peak scheduling → Prevention: Resource isolation config inspection.

Phenomenon: Flink checkpoint failure → Metrics: Paimon write latency, snapshot count → Root Cause: Excessive Paimon snapshots → Solution: Cleanup snapshots + fix cleanup job → Prevention: Snapshot count monitoring alert.

The knowledge graph is regularly updated through historical incident retrospectives, powering the root cause analysis engine's intelligent recommendation capability.

4.2 Establish Tiered Closed-Loop Operations Process

Anomaly Alerting: Detection layer pushes alerts via Alibaba Cloud Ops platform to DingTalk/WeCom, tagged with severity (P1 core business interruption, P2 non-core impact, P3 no business impact).

Auto Repair: For P2/P3 common anomalies (task resource shortage, excessive snapshots), trigger automated scripts via Alibaba Cloud Function Compute (resource scaling, snapshot cleanup, task rerun).

Effect Verification: Post-repair, automatically collect metrics and compare with baseline; if resolved, close loop; else escalate to manual handling.

Manual Fallback: P1 requires ops + dev team intervention within 15 minutes; P2 within 2 hours. Post-resolution retrospective updates knowledge graph.

4.3 Tool Platform Integration

Integrate root cause analysis capabilities into native Alibaba Cloud platforms (DataWorks, ARMS, EMR Monitoring) for "one-stop operations":

Data Integration: API integration with component monitoring data for centralized metric and log display.

Visual Operations: Provide anomaly traceback, root cause analysis report viewing, one-click repair.

Report Output: Auto-generate daily/weekly ops reports covering anomaly count, root cause distribution, repair timeliness, supporting optimization decisions.

4.4 Team Capability and Process Building

Specialized Training: Conduct training on Alibaba Cloud component monitoring, root cause analysis methods, tool usage to enhance team ops capability.

Retrospective Mechanism: Monthly incident retrospective meetings to capture lessons, update knowledge graph and ops standards.

Inspection Regime: Regular inspection of core configs (resource isolation, snapshot cleanup, task rerun) to proactively detect potential risks.

5. AI Empowerment and Future Optimization Directions

5.1 Intelligent Baseline Prediction

Leverage Alibaba Cloud AI PAI platform to train ML models (e.g., LSTM) on historical metric data, intelligently predicting metric baselines for different business scenarios and time periods (e.g., Kafka produce rate baseline during promotions, EMR resource baseline during early morning migrations), improving anomaly detection accuracy.

5.2 Cross-Component Intelligent Correlation Analysis

Utilize Large Language Models (LLM) to understand multi-component log semantics, uncover hidden cross-component correlations (e.g., EMR metadata service anomaly linked to Paimon cleanup failure and Flink checkpoint failure), enabling precise root cause localization for complex failures.

5.3 End-to-End Intelligent Self-Healing

Based on reinforcement learning algorithms, let AI models learn from historical repair cases to automatically generate repair scripts adapted to new scenarios; achieve full automation from anomaly detection, root cause analysis, to repair verification, drastically reducing manual intervention costs.

6. Summary

This solution, built on Alibaba Cloud core big data components, establishes an end-to-end operations root cause analysis system covering cluster, task, and data pipeline dimensions. Core advantages lie in full adaptation to Alibaba Cloud native capabilities, enabling seamless monitoring data integration and analysis; through knowledge graph and automation scripts, root cause localization efficiency and closed-loop resolution speed are improved.

Implementation follows a phased approach:

Phase 1: Complete core component monitoring integration and basic root cause analysis capability setup.

Phase 2: Enhance knowledge graph and automated repair processes.

Phase 3: Introduce AI capabilities for intelligent upgrade. Ultimately realize the transformation from passive firefighting to proactive defense, providing solid assurance for stable operation of Alibaba Cloud big data systems.

Original Source

Signed-in readers can open the original source through BestHub's protected redirect.

Sign in to view source
Republication Notice

This article has been distilled and summarized from source material, then republished for learning and reference. If you believe it infringes your rights, please contactadmin@besthub.devand we will review it promptly.

Knowledge GraphAlibaba Cloudroot cause analysisAI OperationsBig Data OperationsAutomated Remediation
Lakehouse Research Base
Written by

Lakehouse Research Base

Focused on technical sharing in the data field, covering a tech stack that includes Hadoop, Spark, Flink, Kafka, Fluss, Paimon, Iceberg, StarRocks, ClickHouse, ES, Milvus, and more. Welcome to follow.

0 followers
Reader feedback

How this landed with the community

Sign in to like

Rate this article

Was this worth your time?

Sign in to rate
Discussion

0 Comments

Thoughtful readers leave field notes, pushback, and hard-won operational detail here.