| name | flink-optimizer |
| description | Optimize Flink job code and configuration based on runtime metrics and execution flow graph. Use when user asks to "optimize Flink job", "analyze Flink job performance", "check for backpressure", "review Flink DAG", "find Flink configuration issues", or provides a Flink cluster URL with job ID. Analyzes running jobs from Flink UI (via REST API), reviews YAML jobspec configuration, examines Java DAG construction code, and identifies optimization opportunities including backpressure, data skew, dead code/wiring, incorrect stream connections, checkpoint issues, memory configuration, and inefficient aggregations. |
Flink Job Optimizer
Comprehensive Flink job optimization skill for analyzing and improving Apache Flink jobs built with the jobspec-based runtime.
Overview
This skill helps optimize Flink jobs by:
- Fetching runtime metrics from Flink cluster REST API (backpressure, checkpoints, parallelism)
- Analyzing YAML jobspec for configuration issues (dead streams, missing wiring, memory settings)
- Reviewing Java code for DAG construction problems (operator chaining, inefficient wiring)
- Generating actionable recommendations with specific fixes
Typical Usage
User starts Claude session in the Flink job codebase and says:
- "Optimize the payment-processing Flink job using cluster https://flink.de.razorpay.com"
- "Analyze backpressure issues in job abc123 on https://flink.de.razorpay.com"
- "Optimize job {job_id} on cluster {url}, config at configs/local/jobs/sample-jobspec.yaml"
- "Review my Flink job for performance issues"
Workflow
Step 1: Gather Information
Ask the user for:
- Flink cluster URL (e.g.,
https://flink.de.razorpay.com)
- Job ID (optional - can list running jobs if not provided)
- Jobspec YAML path (e.g.,
configs/local/jobs/sample-jobspec_eventtime.yaml)
If job ID not provided: Run the analysis script with --list-jobs to show running jobs, then ask user which job to analyze.
Step 2: Run Analysis Script
Activate Pyhton Virtual Environment if not already active.
Execute scripts/analyze_flink_job.py:
python3 ~/.claude/skills/flink-optimizer/scripts/analyze_flink_job.py \
--cluster-url https://flink.de.razorpay.com \
--job-id <job-id> \
--jobspec-path configs/local/jobs/sample-jobspec.yaml
Script Output: Comprehensive report with:
- Runtime issues (backpressure, checkpoint failures)
- Configuration issues (dead streams, missing wiring, low parallelism)
- DAG summary (sources, operators, streams)
Step 3: Deep Dive Code Analysis
Based on issues found, examine the Java code:
-
For wiring issues: Check Main.java buildJobDAG() method
- Look for how
inputStreams and outputStream are wired in the streams HashMap
- Verify all outputStreams are consumed by downstream operators or sinks
-
For parallelism issues: Check operator configurations
- Review
applyOperatorParallelism() logic in Main.java
- Verify parallelism settings in YAML for heavy operators
-
For aggregation issues: Review operator implementations
- Check
GenericWindowAggregateFunction for efficiency
- Look for repeated CASE WHEN patterns in YAML aggregations
Step 4: Generate Recommendations
Provide specific, actionable fixes with file paths and line numbers.
Example format:
## Optimization Recommendations
### 1. Fix Dead Output Stream (CRITICAL)
**Issue:** Operator 'merchant-filter' produces 'filtered-payments' but no operator consumes it.
**Fix in YAML (configs/local/jobs/sample-jobspec.yaml:82):**
```yaml
operators:
- name: "window-aggregator"
inputStreams: ["filtered-payments"] # ← Add this
outputStream: "aggregated-payments"
Expected Impact: Fixes broken DAG, enables proper data flow.
### Step 5: Validate Fixes (Optional)
If user requests, help validate the fixes:
1. Review modified YAML for syntax correctness
2. Check that all stream references are now valid
3. Verify parallelism settings are reasonable
4. Run ConfigValidator logic manually if needed
## Common Issues & Fixes
### Issue 1: Dead Output Streams
**Detection:** Script reports "Dead output stream 'X' from operator 'Y'"
**Root Cause:** Operator produces an outputStream that no downstream operator consumes.
**Fix:**
- Either wire it to a downstream operator's inputStreams
- Or remove the operator if it's not needed
**Files to check:**
- YAML jobspec: Look for operators with matching outputStream
- `Main.java`: Check `buildJobDAG()` - outputStream stored but never retrieved
### Issue 2: Missing Input Streams
**Detection:** Script reports "Missing input stream 'X' required by operators: Y"
**Root Cause:** Operator references an inputStream that doesn't exist.
**Fix:**
- Correct the inputStream name to match existing outputStream or source
- Or add the missing upstream operator/source
**Files to check:**
- YAML jobspec: Verify source names and operator outputStreams
- `Main.java`: Check if source name matches YAML definition
### Issue 3: Backpressure from Low Parallelism
**Detection:**
- Script reports "Low parallelism (1) for heavy operator"
- Flink UI shows high backpressure on specific vertices
**Root Cause:** Expensive operators (window aggregation, joins) run with parallelism=1.
**Fix:**
- Add explicit parallelism to the operator in YAML
- Or increase default job parallelism
**Recommended Values:**
- Window aggregators: 4-16
- Rule evaluators: 2-8
- RCA analyzers: 4-8
- Filters: 1-2 (lightweight)
**Files to modify:**
- YAML jobspec: Add `parallelism: N` to operator config
### Issue 4: Data Skew
**Detection:**
- Some subtasks much busier than others
- Per-subtask metrics show large variance (check Flink UI metrics)
**Root Cause:** Uneven key distribution in keyBy operations.
**Fix Options:**
1. **Add key salting** (requires code change):
```java
// In operator implementation
String saltedKey = merchantId + "-" + (hash(merchantId) % numSalts);
- Two-phase aggregation (YAML):
operators:
- name: "pre-agg"
config:
keyField: "salted_key"
- name: "final-agg"
config:
keyField: "original_key"
- Increase parallelism (partial fix):
- Higher parallelism spreads skew across more instances
Files to check:
- YAML: Look for keyField in WINDOW_AGGREGATOR config
- Java: Check WindowAggregateMetricsWrapper keyBy logic
Issue 5: Inefficient Aggregations
Detection: Script reports "Operator has 50+ aggregations" or "similar CASE WHEN conditions"
Root Cause: Too many aggregation expressions or repeated logic.
Fix Options:
- Split into multiple operators:
operators:
- name: "upi-aggregator"
config:
aggregations:
- name: "card-aggregator"
config:
aggregations:
- Simplify expressions:
- Pre-compute complex fields before aggregation
- Reduce repeated CASE WHEN patterns
Files to modify:
- YAML jobspec: Split aggregations config
- May need to add intermediate operators
Issue 6: Checkpoint Issues
Detection:
- Script reports "High checkpoint failure rate" or "Long checkpoint duration"
- Flink UI shows frequent checkpoint timeouts
Root Causes & Fixes:
High Failure Rate:
checkpointing:
interval: 60000
checkpointTimeout: 300000
Long Duration:
- Reduce state size (simplify aggregations, shorter TTL)
- Increase managed memory
- Enable incremental checkpointing
state:
config:
rocksdb:
execution.checkpointing.incremental: true
Interval Too Frequent:
checkpointing:
interval: 60000
Files to modify:
- YAML jobspec: Update checkpointing config
Issue 7: Memory Configuration
Detection: Script warns about low memory allocation
Root Cause: Insufficient memory for state-heavy operators (RocksDB).
Fix:
resources:
memory: "8gb"
managedMemory: "2gb"
Best Practices:
- Total memory: 4-16GB per task slot
- Managed memory: 25-40% of total for RocksDB jobs
Files to modify:
- YAML jobspec: Update resources config
Reference Materials
Flink REST API Endpoints
See references/flink-rest-api.md for detailed API documentation.
Quick Reference:
- List jobs:
GET /jobs
- Job details:
GET /jobs/:jobid
- Checkpoint stats:
GET /jobs/:jobid/checkpoints
- Metrics:
GET /jobs/:jobid/metrics
Optimization Patterns
See references/optimization-patterns.md for comprehensive patterns and anti-patterns.
Key Sections:
- DAG Wiring Issues
- Parallelism & Backpressure
- Data Skew
- Checkpoint Configuration
- Memory Management
- Window Aggregation Optimization
- Operator Chaining
Code Structure Reference
Main.java Key Methods
buildJobDAG(env, jobSpec) - Constructs the DAG:
Map<String, DataStream<?>> streams = new HashMap<>();
for (SourceConfig source : jobSpec.getSources()) {
streams.put(source.getName(), buildSource(env, source));
}
for (OperatorConfig operator : jobSpec.getOperators()) {
DataStream<?> result = buildOperator(streams, operator, jobSpec, env);
streams.put(operator.getOutputStream(), result);
}
for (SinkConfig sink : jobSpec.getSinks()) {
buildSink(streams, sink);
}
buildOperator(streams, operatorConfig, jobSpec, env) - Builds specific operator:
String inputStreamName = operatorConfig.getInputStreams().get(0);
DataStream<?> inputStream = streams.get(inputStreamName);
DataStream<?> result = inputStream.filter(...).map(...);
return result;
Jobspec YAML Structure
jobName: "payment-detection"
parallelism: 2
sources:
- name: "payments-kafka-source"
type: "KAFKA"
config:
topic: "payments"
operators:
- name: "merchant-filter"
type: "FILTER"
inputStreams: ["payments-kafka-source"]
outputStream: "filtered-payments"
parallelism: 2
- name: "window-aggregator"
type: "WINDOW_AGGREGATOR"
inputStreams: ["filtered-payments"]
outputStream: "aggregated-payments"
parallelism: 8
config:
windowType: "TUMBLING"
keyField: "merchant_id"
aggregations:
total_count: "COUNT(*)"
- name: "kafka-sink"
type: "ANOMALY_SINK"
inputStreams: ["aggregated-payments"]
config:
sinks:
- type: "KAFKA"
topic: "alerts"
Tips for Effective Optimization
- Start with the script - Don't manually inspect without data
- Fix critical issues first - Dead streams and missing wiring before performance tuning
- Measure before and after - Note current backpressure, throughput metrics
- Test incrementally - Fix one issue at a time, validate
- Use Flink UI - Visual representation helps understand flow
- Check job logs - ConfigValidator errors provide specific line numbers
- Validate YAML syntax - Use YAML linter before deploying
When to Read Reference Files
- Before first optimization: Skim
optimization-patterns.md for overview
- For API questions: Check
flink-rest-api.md for endpoint details
- For specific issues: Jump to relevant section in
optimization-patterns.md
- DAG issues → "DAG Wiring Issues"
- Slow processing → "Parallelism & Backpressure"
- Uneven load → "Data Skew"
- Checkpoint failures → "Checkpoint Configuration"
- OOM errors → "Memory Management"
- Slow aggregations → "Window Aggregation Optimization"