全部笔记All notes

Flink 监控体系完全指南

阅读 8m 52s8m 52s read

概述

在生产环境中运行 Flink 应用时,完善的监控体系是保证系统稳定性和性能的关键。Flink 提供了丰富的内置指标和灵活的监控框架,支持与各种监控系统集成。

💡 提示: 监控不仅是发现问题的手段,更是优化系统性能的重要依据。

监控体系组成

组件作用常用工具
指标收集收集运行时指标Flink Metrics API
指标存储持久化存储指标数据Prometheus、InfluxDB
可视化展示图表展示和分析Grafana、Kibana
告警通知异常情况及时通知AlertManager、PagerDuty
日志分析日志收集和分析ELK Stack、Loki

监控架构设计

整体架构

┌─────────────────────────────────────────────────────────────┐
│                     Flink 监控架构                           │
├─────────────────────────────────────────────────────────────┤
│                                                             │
│  Flink 集群                                                  │
│  ┌─────────────┐  ┌─────────────┐  ┌─────────────┐        │
│  │ JobManager  │  │TaskManager 1│  │TaskManager 2│        │
│  │  ┌───────┐  │  │  ┌───────┐  │  │  ┌───────┐  │        │
│  │  │Metrics│  │  │  │Metrics│  │  │  │Metrics│  │        │
│  │  └───┬───┘  │  │  └───┬───┘  │  │  └───┬───┘  │        │
│  └──────┼──────┘  └──────┼──────┘  └──────┼──────┘        │
│         │                 │                 │                │
│         └─────────────────┴─────────────────┘                │
│                           │                                  │
│                    ┌──────▼──────┐                          │
│                    │  Reporter   │                          │
│                    │(Prometheus) │                          │
│                    └──────┬──────┘                          │
│                           │                                  │
├───────────────────────────┼─────────────────────────────────┤
│                           │                                  │
│  监控系统                  ▼                                  │
│  ┌─────────────────────────────────────────┐               │
│  │           Prometheus Server              │               │
│  │  ┌─────────┐  ┌─────────┐  ┌─────────┐ │               │
│  │  │ Storage │  │  Rules  │  │   API   │ │               │
│  │  └─────────┘  └─────────┘  └─────────┘ │               │
│  └────────┬───────────┬───────────┬────────┘               │
│           │           │           │                         │
│           ▼           ▼           ▼                         │
│  ┌─────────────┐ ┌─────────┐ ┌─────────────┐              │
│  │   Grafana   │ │  Alert  │ │   其他系统   │              │
│  │ Dashboard   │ │ Manager │ │ (ELK, etc)  │              │
│  └─────────────┘ └─────────┘ └─────────────┘              │
│                                                             │
└─────────────────────────────────────────────────────────────┘

监控层次

public class MonitoringLayers {
    // 1. 应用层监控
    public static class ApplicationMonitoring {
        // 业务指标
        private Counter ordersProcessed;
        private Meter transactionRate;
        private Histogram orderAmount;
        
        // 数据质量
        private Counter invalidRecords;
        private Gauge dataLag;
    }
    
    // 2. 框架层监控
    public static class FrameworkMonitoring {
        // Flink 内置指标
        private Gauge checkpointDuration;
        private Counter numRestarts;
        private Gauge backPressure;
    }
    
    // 3. 系统层监控
    public static class SystemMonitoring {
        // 资源使用
        private Gauge cpuUsage;
        private Gauge memoryUsage;
        private Gauge diskIO;
        private Gauge networkIO;
    }
}

核心监控指标

作业级别指标

1. 作业运行状态

public class JobMetrics {
    // 作业状态监控
    public static void monitorJobStatus(MetricGroup jobMetricGroup) {
        // 作业运行时间
        jobMetricGroup.gauge("uptime", new Gauge<Long>() {
            @Override
            public Long getValue() {
                return System.currentTimeMillis() - jobStartTime;
            }
        });
        
        // 重启次数
        jobMetricGroup.counter("numRestarts");
        
        // 失败次数
        jobMetricGroup.counter("numFailures");
        
        // 最后一次失败时间
        jobMetricGroup.gauge("lastFailureTime", new Gauge<Long>() {
            @Override
            public Long getValue() {
                return lastFailureTimestamp;
            }
        });
    }
}

2. Checkpoint 指标

public class CheckpointMetrics {
    // Checkpoint 监控指标
    public static class CheckpointMonitor {
        // 持续时间
        private Histogram checkpointDuration;
        // 大小
        private Histogram checkpointSize;
        // 成功率
        private Meter checkpointSuccessRate;
        // 失败率
        private Meter checkpointFailureRate;
        // 最后成功时间
        private Gauge<Long> lastSuccessfulCheckpoint;
        
        public void registerMetrics(MetricGroup metricGroup) {
            checkpointDuration = metricGroup.histogram(
                "checkpoint.duration",
                new DescriptiveStatisticsHistogram(1000)
            );
            
            checkpointSize = metricGroup.histogram(
                "checkpoint.size",
                new DescriptiveStatisticsHistogram(1000)
            );
            
            checkpointSuccessRate = metricGroup.meter(
                "checkpoint.success.rate",
                new MeterView(60)
            );
            
            checkpointFailureRate = metricGroup.meter(
                "checkpoint.failure.rate",
                new MeterView(60)
            );
            
            metricGroup.gauge("checkpoint.last.success", 
                () -> lastSuccessfulCheckpointTime);
        }
    }
}

任务级别指标

1. 数据处理指标

public class TaskMetrics extends RichMapFunction<Event, ProcessedEvent> {
    // 吞吐量指标
    private transient Counter recordsProcessed;
    private transient Meter recordsPerSecond;
    
    // 延迟指标
    private transient Histogram processingLatency;
    private transient Gauge<Long> currentLatency;
    
    // 错误指标
    private transient Counter errors;
    private transient Meter errorRate;
    
    @Override
    public void open(Configuration parameters) {
        MetricGroup metrics = getRuntimeContext().getMetricGroup();
        
        // 注册计数器
        recordsProcessed = metrics.counter("records.processed");
        errors = metrics.counter("records.errors");
        
        // 注册速率计
        recordsPerSecond = metrics.meter("records.rate", new MeterView(60));
        errorRate = metrics.meter("error.rate", new MeterView(60));
        
        // 注册直方图
        processingLatency = metrics.histogram("processing.latency",
            new DescriptiveStatisticsHistogram(1000));
        
        // 注册仪表
        metrics.gauge("current.latency", () -> lastProcessingTime);
    }
    
    @Override
    public ProcessedEvent map(Event event) throws Exception {
        long startTime = System.currentTimeMillis();
        
        try {
            ProcessedEvent result = processEvent(event);
            
            // 更新成功指标
            recordsProcessed.inc();
            recordsPerSecond.markEvent();
            
            // 记录处理延迟
            long latency = System.currentTimeMillis() - startTime;
            processingLatency.update(latency);
            lastProcessingTime = latency;
            
            return result;
        } catch (Exception e) {
            // 更新错误指标
            errors.inc();
            errorRate.markEvent();
            throw e;
        }
    }
}

2. 算子级别指标

public class OperatorMetrics {
    // 窗口算子监控
    public static class WindowOperatorMetrics<T> 
            extends ProcessWindowFunction<T, WindowResult<T>, String, TimeWindow> {
        
        private transient Counter windowsTriggered;
        private transient Histogram windowSize;
        private transient Gauge<Integer> activeWindows;
        
        @Override
        public void open(Configuration parameters) {
            MetricGroup metrics = getRuntimeContext().getMetricGroup();
            
            windowsTriggered = metrics.counter("windows.triggered");
            windowSize = metrics.histogram("window.size",
                new DescriptiveStatisticsHistogram(1000));
            metrics.gauge("windows.active", () -> currentActiveWindows.get());
        }
        
        @Override
        public void process(String key, Context context, 
                          Iterable<T> elements, Collector<WindowResult<T>> out) {
            List<T> windowElements = new ArrayList<>();
            elements.forEach(windowElements::add);
            
            // 更新指标
            windowsTriggered.inc();
            windowSize.update(windowElements.size());
            
            // 处理窗口数据
            WindowResult<T> result = processWindow(key, windowElements, context.window());
            out.collect(result);
        }
    }
}

系统级别指标

1. 资源使用监控

public class SystemMetrics {
    public static void registerSystemMetrics(MetricGroup metricGroup) {
        // CPU 使用率
        metricGroup.gauge("system.cpu.usage", new Gauge<Double>() {
            private final OperatingSystemMXBean osBean = 
                ManagementFactory.getOperatingSystemMXBean();
            
            @Override
            public Double getValue() {
                if (osBean instanceof com.sun.management.OperatingSystemMXBean) {
                    return ((com.sun.management.OperatingSystemMXBean) osBean)
                        .getProcessCpuLoad() * 100;
                }
                return -1.0;
            }
        });
        
        // 内存使用
        metricGroup.gauge("system.memory.heap.used", new Gauge<Long>() {
            private final MemoryMXBean memoryBean = 
                ManagementFactory.getMemoryMXBean();
            
            @Override
            public Long getValue() {
                return memoryBean.getHeapMemoryUsage().getUsed();
            }
        });
        
        metricGroup.gauge("system.memory.heap.max", new Gauge<Long>() {
            private final MemoryMXBean memoryBean = 
                ManagementFactory.getMemoryMXBean();
            
            @Override
            public Long getValue() {
                return memoryBean.getHeapMemoryUsage().getMax();
            }
        });
        
        // GC 监控
        for (GarbageCollectorMXBean gc : 
             ManagementFactory.getGarbageCollectorMXBeans()) {
            
            metricGroup.gauge("system.gc." + gc.getName() + ".count", 
                new Gauge<Long>() {
                    @Override
                    public Long getValue() {
                        return gc.getCollectionCount();
                    }
                });
            
            metricGroup.gauge("system.gc." + gc.getName() + ".time", 
                new Gauge<Long>() {
                    @Override
                    public Long getValue() {
                        return gc.getCollectionTime();
                    }
                });
        }
    }
}

2. 网络 I/O 监控

public class NetworkMetrics {
    public static class NetworkMonitor {
        private final MetricGroup metrics;
        
        public NetworkMonitor(MetricGroup metrics) {
            this.metrics = metrics;
            registerNetworkMetrics();
        }
        
        private void registerNetworkMetrics() {
            // 输入缓冲区使用率
            metrics.gauge("network.input.buffers.usage", new Gauge<Float>() {
                @Override
                public Float getValue() {
                    return getInputBufferUsage();
                }
            });
            
            // 输出缓冲区使用率
            metrics.gauge("network.output.buffers.usage", new Gauge<Float>() {
                @Override
                public Float getValue() {
                    return getOutputBufferUsage();
                }
            });
            
            // 网络吞吐量
            metrics.meter("network.input.bytes.rate", new MeterView(60));
            metrics.meter("network.output.bytes.rate", new MeterView(60));
            
            // 缓冲池状态
            metrics.gauge("network.buffers.available", 
                () -> getAvailableBuffers());
            metrics.gauge("network.buffers.total", 
                () -> getTotalBuffers());
        }
    }
}

监控系统集成

Prometheus 集成

1. Prometheus Reporter 配置

# flink-conf.yaml
metrics.reporters: prom
metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
metrics.reporter.prom.port: 9249
metrics.reporter.prom.host: 0.0.0.0

# 高级配置
metrics.scope.jm: flink.jobmanager
metrics.scope.tm: flink.taskmanager.<host>
metrics.scope.jm.job: flink.jobmanager.job.<job_name>
metrics.scope.tm.job: flink.taskmanager.<host>.job.<job_name>
metrics.scope.task: flink.task.<task_id>.<subtask_index>
metrics.scope.operator: flink.operator.<operator_id>.<subtask_index>

2. Prometheus 配置

# prometheus.yml
global:
  scrape_interval: 15s
  evaluation_interval: 15s

scrape_configs:
  - job_name: 'flink'
    static_configs:
      # JobManager
      - targets: ['jobmanager:9249']
        labels:
          component: 'jobmanager'
      
      # TaskManagers
      - targets: ['taskmanager-1:9249', 'taskmanager-2:9249']
        labels:
          component: 'taskmanager'
    
    # 动态服务发现(Kubernetes)
  - job_name: 'flink-k8s'
    kubernetes_sd_configs:
      - role: pod
    relabel_configs:
      - source_labels: [__meta_kubernetes_pod_label_component]
        regex: (jobmanager|taskmanager)
        action: keep
      - source_labels: [__meta_kubernetes_pod_annotation_prometheus_io_port]
        action: replace
        target_label: __address__
        regex: ([^:]+)(?::\d+)?;(\d+)
        replacement: $1:$2

3. PromQL 查询示例

# 作业吞吐量
sum(rate(flink_taskmanager_job_task_operator_records_out[1m])) by (job_name)

# Checkpoint 持续时间 P99
histogram_quantile(0.99, 
  sum(rate(flink_jobmanager_job_checkpoint_duration_bucket[5m])) by (job_name, le)
)

# 任务管理器内存使用率
(flink_taskmanager_Status_JVM_Memory_Heap_Used / 
 flink_taskmanager_Status_JVM_Memory_Heap_Max) * 100

# 反压监控
avg(flink_taskmanager_job_task_backPressureTimeMsPerSecond) by (task_name) / 10

# GC 时间占比
rate(flink_taskmanager_Status_JVM_GarbageCollector_G1_Old_Generation_Time[5m]) / 60000

Grafana 可视化

1. Dashboard 设计

{
  "dashboard": {
    "title": "Flink Monitoring Dashboard",
    "panels": [
      {
        "title": "Job Overview",
        "targets": [
          {
            "expr": "up{job='flink'}",
            "legendFormat": "{{instance}} - {{component}}"
          }
        ],
        "type": "stat"
      },
      {
        "title": "Throughput",
        "targets": [
          {
            "expr": "sum(rate(flink_taskmanager_job_task_operator_records_out[1m])) by (job_name)",
            "legendFormat": "{{job_name}}"
          }
        ],
        "type": "graph"
      },
      {
        "title": "Checkpoint Duration",
        "targets": [
          {
            "expr": "histogram_quantile(0.99, sum(rate(flink_jobmanager_job_checkpoint_duration_bucket[5m])) by (job_name, le))",
            "legendFormat": "{{job_name}} - P99"
          }
        ],
        "type": "graph"
      },
      {
        "title": "Memory Usage",
        "targets": [
          {
            "expr": "(flink_taskmanager_Status_JVM_Memory_Heap_Used / flink_taskmanager_Status_JVM_Memory_Heap_Max) * 100",
            "legendFormat": "{{host}} - Heap Usage %"
          }
        ],
        "type": "graph"
      }
    ]
  }
}

2. 告警面板

{
  "alert": {
    "name": "High Checkpoint Duration",
    "conditions": [
      {
        "evaluator": {
          "params": [300000],
          "type": "gt"
        },
        "query": {
          "model": {
            "expr": "flink_jobmanager_job_checkpoint_duration",
            "intervalMs": 1000,
            "maxDataPoints": 43200,
            "refId": "A"
          }
        }
      }
    ],
    "frequency": "60s",
    "handler": 1,
    "message": "Checkpoint duration is too high",
    "name": "High Checkpoint Duration",
    "noDataState": "no_data",
    "notifications": []
  }
}

其他监控系统

1. InfluxDB 集成

public class InfluxDBReporter extends AbstractReporter {
    private InfluxDB influxDB;
    private String database = "flink_metrics";
    
    @Override
    public void open(MetricConfig config) {
        String url = config.getString("url", "http://localhost:8086");
        String username = config.getString("username", "admin");
        String password = config.getString("password", "admin");
        
        influxDB = InfluxDBFactory.connect(url, username, password);
        influxDB.setDatabase(database);
        influxDB.enableBatch(BatchOptions.DEFAULTS);
    }
    
    @Override
    public void report() {
        long timestamp = System.currentTimeMillis();
        BatchPoints batchPoints = BatchPoints.database(database)
            .retentionPolicy("autogen")
            .build();
        
        // 报告所有指标
        for (Map.Entry<Gauge<?>, String> entry : gauges.entrySet()) {
            Point point = Point.measurement(entry.getValue())
                .time(timestamp, TimeUnit.MILLISECONDS)
                .addField("value", entry.getKey().getValue())
                .build();
            batchPoints.point(point);
        }
        
        influxDB.write(batchPoints);
    }
}

2. Datadog 集成

# flink-conf.yaml
metrics.reporters: datadog
metrics.reporter.datadog.class: org.apache.flink.metrics.datadog.DatadogHttpReporter
metrics.reporter.datadog.apikey: <your-api-key>
metrics.reporter.datadog.host: https://api.datadoghq.com
metrics.reporter.datadog.tags: env:prod,team:data
metrics.reporter.datadog.interval: 60 SECONDS

自定义指标开发

Metrics API 使用

1. 基础指标类型

public class CustomMetrics extends RichMapFunction<Order, EnrichedOrder> {
    // Counter - 计数器
    private transient Counter orderCounter;
    
    // Gauge - 仪表
    private transient Gauge<Double> avgOrderValue;
    private double sumOrderValue = 0;
    private long orderCount = 0;
    
    // Histogram - 直方图
    private transient Histogram orderValueDistribution;
    
    // Meter - 速率计
    private transient Meter orderRate;
    
    @Override
    public void open(Configuration parameters) throws Exception {
        MetricGroup metricGroup = getRuntimeContext()
            .getMetricGroup()
            .addGroup("business")
            .addGroup("orders");
        
        // 注册计数器
        orderCounter = metricGroup.counter("total");
        
        // 注册仪表
        metricGroup.gauge("average_value", new Gauge<Double>() {
            @Override
            public Double getValue() {
                return orderCount > 0 ? sumOrderValue / orderCount : 0.0;
            }
        });
        
        // 注册直方图
        orderValueDistribution = metricGroup.histogram(
            "value_distribution",
            new DescriptiveStatisticsHistogram(1000)
        );
        
        // 注册速率计
        orderRate = metricGroup.meter("rate", new MeterView(60));
    }
    
    @Override
    public EnrichedOrder map(Order order) throws Exception {
        // 更新指标
        orderCounter.inc();
        orderRate.markEvent();
        
        sumOrderValue += order.getValue();
        orderCount++;
        
        orderValueDistribution.update((long) order.getValue());
        
        return enrichOrder(order);
    }
}

2. 复杂指标实现

public class AdvancedMetrics {
    // 自定义 Gauge 实现
    public static class CacheHitRateGauge implements Gauge<Double> {
        private final AtomicLong hits = new AtomicLong(0);
        private final AtomicLong misses = new AtomicLong(0);
        
        public void recordHit() {
            hits.incrementAndGet();
        }
        
        public void recordMiss() {
            misses.incrementAndGet();
        }
        
        @Override
        public Double getValue() {
            long totalHits = hits.get();
            long totalMisses = misses.get();
            long total = totalHits + totalMisses;
            
            return total > 0 ? (double) totalHits / total : 0.0;
        }
    }
    
    // 自定义 Histogram 实现
    public static class PercentileHistogram implements Histogram {
        private final DescriptiveStatistics stats = 
            new DescriptiveStatistics(10000);
        
        @Override
        public void update(long value) {
            synchronized (stats) {
                stats.addValue(value);
            }
        }
        
        @Override
        public long getCount() {
            return stats.getN();
        }
        
        @Override
        public HistogramStatistics getStatistics() {
            synchronized (stats) {
                return new HistogramStatistics() {
                    @Override
                    public double getQuantile(double quantile) {
                        return stats.getPercentile(quantile * 100);
                    }
                    
                    @Override
                    public long[] getValues() {
                        return new long[0];
                    }
                    
                    @Override
                    public int size() {
                        return (int) stats.getN();
                    }
                    
                    @Override
                    public long getMin() {
                        return (long) stats.getMin();
                    }
                    
                    @Override
                    public long getMax() {
                        return (long) stats.getMax();
                    }
                    
                    @Override
                    public double getMean() {
                        return stats.getMean();
                    }
                    
                    @Override
                    public double getStdDev() {
                        return stats.getStandardDeviation();
                    }
                };
            }
        }
    }
}

业务指标实现

1. 业务监控指标

public class BusinessMetrics {
    // 订单处理监控
    public static class OrderProcessingMetrics 
            extends KeyedProcessFunction<String, Order, ProcessedOrder> {
        
        // 业务指标
        private transient ValueState<OrderStats> orderStatsState;
        private transient MapState<String, Long> productSalesState;
        
        // 监控指标
        private transient Counter successfulOrders;
        private transient Counter failedOrders;
        private transient Histogram orderProcessingTime;
        private transient Gauge<Double> revenuePerMinute;
        
        @Override
        public void open(Configuration parameters) throws Exception {
            // 初始化状态
            orderStatsState = getRuntimeContext().getState(
                new ValueStateDescriptor<>("order-stats", OrderStats.class));
            productSalesState = getRuntimeContext().getMapState(
                new MapStateDescriptor<>("product-sales", String.class, Long.class));
            
            // 初始化指标
            MetricGroup metrics = getRuntimeContext()
                .getMetricGroup()
                .addGroup("business")
                .addGroup("orders");
            
            successfulOrders = metrics.counter("successful");
            failedOrders = metrics.counter("failed");
            orderProcessingTime = metrics.histogram("processing_time",
                new DescriptiveStatisticsHistogram(1000));
            
            metrics.gauge("revenue_per_minute", new Gauge<Double>() {
                @Override
                public Double getValue() {
                    try {
                        OrderStats stats = orderStatsState.value();
                        if (stats != null) {
                            long timeDiff = System.currentTimeMillis() - stats.startTime;
                            double minutes = timeDiff / 60000.0;
                            return minutes > 0 ? stats.totalRevenue / minutes : 0.0;
                        }
                    } catch (Exception e) {
                        LOG.error("Error calculating revenue rate", e);
                    }
                    return 0.0;
                }
            });
        }
        
        @Override
        public void processElement(Order order, Context ctx, 
                                 Collector<ProcessedOrder> out) throws Exception {
            long startTime = System.currentTimeMillis();
            
            try {
                // 处理订单
                ProcessedOrder processed = processOrder(order);
                
                // 更新业务状态
                OrderStats stats = orderStatsState.value();
                if (stats == null) {
                    stats = new OrderStats();
                    stats.startTime = System.currentTimeMillis();
                }
                stats.totalRevenue += order.getAmount();
                stats.orderCount++;
                orderStatsState.update(stats);
                
                // 更新产品销售统计
                Long productSales = productSalesState.get(order.getProductId());
                productSalesState.put(order.getProductId(), 
                    (productSales == null ? 0L : productSales) + 1);
                
                // 更新监控指标
                successfulOrders.inc();
                orderProcessingTime.update(System.currentTimeMillis() - startTime);
                
                out.collect(processed);
                
            } catch (Exception e) {
                failedOrders.inc();
                LOG.error("Order processing failed", e);
                throw e;
            }
        }
    }
}

2. 数据质量监控

public class DataQualityMetrics extends ProcessFunction<RawData, ValidatedData> {
    // 数据质量指标
    private transient Counter totalRecords;
    private transient Counter validRecords;
    private transient Counter invalidRecords;
    private transient Map<String, Counter> errorTypeCounters;
    
    // 质量评分
    private transient Gauge<Double> dataQualityScore;
    private final SlidingWindow qualityWindow = new SlidingWindow(1000);
    
    @Override
    public void open(Configuration parameters) throws Exception {
        MetricGroup metrics = getRuntimeContext()
            .getMetricGroup()
            .addGroup("data_quality");
        
        totalRecords = metrics.counter("total");
        validRecords = metrics.counter("valid");
        invalidRecords = metrics.counter("invalid");
        
        errorTypeCounters = new HashMap<>();
        for (ErrorType errorType : ErrorType.values()) {
            errorTypeCounters.put(errorType.name(), 
                metrics.counter("error." + errorType.name().toLowerCase()));
        }
        
        metrics.gauge("score", new Gauge<Double>() {
            @Override
            public Double getValue() {
                return qualityWindow.getAverageQuality();
            }
        });
    }
    
    @Override
    public void processElement(RawData data, Context ctx, 
                             Collector<ValidatedData> out) throws Exception {
        totalRecords.inc();
        
        ValidationResult result = validateData(data);
        
        if (result.isValid()) {
            validRecords.inc();
            qualityWindow.addSample(1.0);
            out.collect(new ValidatedData(data));
        } else {
            invalidRecords.inc();
            qualityWindow.addSample(0.0);
            
            // 记录错误类型
            for (ErrorType errorType : result.getErrors()) {
                errorTypeCounters.get(errorType.name()).inc();
            }
            
            // 输出到侧输出流
            ctx.output(INVALID_DATA_TAG, 
                new InvalidData(data, result.getErrors()));
        }
    }
}

告警体系建设

告警规则设计

1. 告警规则配置

# alerting-rules.yml
groups:
  - name: flink_alerts
    interval: 30s
    rules:
      # 作业失败告警
      - alert: FlinkJobFailed
        expr: flink_jobmanager_job_uptime == 0
        for: 1m
        labels:
          severity: critical
          team: data-platform
        annotations:
          summary: "Flink job {{ $labels.job_name }} has failed"
          description: "Job {{ $labels.job_name }} is not running for more than 1 minute"
      
      # Checkpoint 超时告警
      - alert: CheckpointTimeout
        expr: flink_jobmanager_job_checkpoint_duration > 300000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Checkpoint duration too high"
          description: "Checkpoint duration for {{ $labels.job_name }} is {{ $value }}ms"
      
      # 反压告警
      - alert: BackPressureHigh
        expr: avg(flink_taskmanager_job_task_backPressureTimeMsPerSecond) by (task_name) > 100
        for: 10m
        labels:
          severity: warning
        annotations:
          summary: "High back pressure detected"
          description: "Task {{ $labels.task_name }} has high back pressure"
      
      # 内存使用告警
      - alert: HighMemoryUsage
        expr: |
          (flink_taskmanager_Status_JVM_Memory_Heap_Used / 
           flink_taskmanager_Status_JVM_Memory_Heap_Max) > 0.9
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High heap memory usage"
          description: "TaskManager {{ $labels.host }} heap usage is {{ $value | humanizePercentage }}"
      
      # GC 时间过长告警
      - alert: HighGCTime
        expr: |
          rate(flink_taskmanager_Status_JVM_GarbageCollector_G1_Old_Generation_Time[5m]) > 0.5
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High GC time"
          description: "TaskManager {{ $labels.host }} spending too much time in GC"

2. 自定义告警逻辑

public class CustomAlertManager {
    private final AlertingService alertingService;
    private final Map<String, AlertRule> alertRules = new ConcurrentHashMap<>();
    
    public void registerAlertRule(AlertRule rule) {
        alertRules.put(rule.getName(), rule);
    }
    
    // 告警规则定义
    public static class AlertRule {
        private final String name;
        private final String expression;
        private final Duration duration;
        private final AlertSeverity severity;
        private final Map<String, String> labels;
        private final ThresholdEvaluator evaluator;
        
        public boolean evaluate(MetricSnapshot snapshot) {
            return evaluator.evaluate(snapshot, expression);
        }
    }
    
    // 告警评估器
    public static class AlertEvaluator extends ProcessFunction<MetricEvent, Alert> {
        private final Map<String, AlertState> alertStates = new ConcurrentHashMap<>();
        
        @Override
        public void processElement(MetricEvent event, Context ctx, 
                                 Collector<Alert> out) throws Exception {
            
            for (AlertRule rule : alertRules.values()) {
                String alertKey = rule.getName() + "_" + event.getMetricName();
                AlertState state = alertStates.computeIfAbsent(alertKey, 
                    k -> new AlertState());
                
                if (rule.evaluate(event.getSnapshot())) {
                    // 条件满足
                    if (state.firstTriggerTime == 0) {
                        state.firstTriggerTime = ctx.timestamp();
                    }
                    
                    if (ctx.timestamp() - state.firstTriggerTime >= rule.getDuration().toMillis()) {
                        if (!state.fired) {
                            // 触发告警
                            Alert alert = new Alert()
                                .setName(rule.getName())
                                .setSeverity(rule.getSeverity())
                                .setLabels(rule.getLabels())
                                .setValue(event.getValue())
                                .setTimestamp(ctx.timestamp());
                            
                            out.collect(alert);
                            state.fired = true;
                        }
                    }
                } else {
                    // 条件不满足,重置状态
                    state.reset();
                }
            }
        }
    }
}

告警通知集成

1. 多渠道通知

public class NotificationService {
    private final List<NotificationChannel> channels = new ArrayList<>();
    
    public void registerChannel(NotificationChannel channel) {
        channels.add(channel);
    }
    
    public void sendAlert(Alert alert) {
        for (NotificationChannel channel : channels) {
            try {
                if (channel.shouldNotify(alert)) {
                    channel.send(alert);
                }
            } catch (Exception e) {
                LOG.error("Failed to send alert via " + channel.getName(), e);
            }
        }
    }
    
    // 邮件通知
    public static class EmailNotificationChannel implements NotificationChannel {
        private final JavaMailSender mailSender;
        private final String[] recipients;
        
        @Override
        public void send(Alert alert) {
            SimpleMailMessage message = new SimpleMailMessage();
            message.setTo(recipients);
            message.setSubject("[Flink Alert] " + alert.getName());
            message.setText(formatAlertMessage(alert));
            
            mailSender.send(message);
        }
    }
    
    // 钉钉通知
    public static class DingTalkNotificationChannel implements NotificationChannel {
        private final String webhook;
        
        @Override
        public void send(Alert alert) {
            DingTalkMessage message = new DingTalkMessage();
            message.setMsgtype("markdown");
            message.setMarkdown(new Markdown()
                .setTitle("Flink 告警")
                .setText(formatMarkdownMessage(alert)));
            
            HttpPost post = new HttpPost(webhook);
            post.setEntity(new StringEntity(JSON.toJSONString(message), "UTF-8"));
            post.setHeader("Content-Type", "application/json");
            
            try (CloseableHttpClient client = HttpClients.createDefault()) {
                client.execute(post);
            }
        }
    }
    
    // Slack 通知
    public static class SlackNotificationChannel implements NotificationChannel {
        private final SlackWebhook webhook;
        
        @Override
        public void send(Alert alert) {
            SlackMessage message = SlackMessage.builder()
                .text("Flink Alert: " + alert.getName())
                .attachments(Arrays.asList(
                    SlackAttachment.builder()
                        .color(getSeverityColor(alert.getSeverity()))
                        .fields(Arrays.asList(
                            new SlackField("Severity", alert.getSeverity().name(), true),
                            new SlackField("Time", formatTime(alert.getTimestamp()), true),
                            new SlackField("Value", String.valueOf(alert.getValue()), false)
                        ))
                        .build()
                ))
                .build();
            
            webhook.send(message);
        }
    }
}

2. 告警聚合和抑制

public class AlertAggregator {
    // 告警聚合
    public static class AlertAggregationFunction 
            extends KeyedProcessFunction<String, Alert, AggregatedAlert> {
        
        private transient MapState<String, Alert> pendingAlerts;
        private final Duration aggregationWindow = Duration.ofMinutes(5);
        
        @Override
        public void processElement(Alert alert, Context ctx, 
                                 Collector<AggregatedAlert> out) throws Exception {
            String alertKey = alert.getName() + "_" + alert.getSeverity();
            
            // 添加到待处理告警
            pendingAlerts.put(alert.getId(), alert);
            
            // 注册定时器
            ctx.timerService().registerProcessingTimeTimer(
                ctx.timerService().currentProcessingTime() + aggregationWindow.toMillis()
            );
        }
        
        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, 
                          Collector<AggregatedAlert> out) throws Exception {
            List<Alert> alerts = new ArrayList<>();
            pendingAlerts.values().forEach(alerts::add);
            
            if (!alerts.isEmpty()) {
                AggregatedAlert aggregated = new AggregatedAlert()
                    .setName(alerts.get(0).getName())
                    .setSeverity(getMaxSeverity(alerts))
                    .setCount(alerts.size())
                    .setFirstOccurrence(getEarliestTimestamp(alerts))
                    .setLastOccurrence(getLatestTimestamp(alerts))
                    .setAlerts(alerts);
                
                out.collect(aggregated);
                pendingAlerts.clear();
            }
        }
    }
    
    // 告警抑制
    public static class AlertSuppressor {
        private final Map<String, Long> suppressedUntil = new ConcurrentHashMap<>();
        private final Duration suppressionDuration = Duration.ofHours(1);
        
        public boolean shouldSuppress(Alert alert) {
            String key = alert.getName() + "_" + alert.getLabels().toString();
            Long suppressUntil = suppressedUntil.get(key);
            
            if (suppressUntil != null && System.currentTimeMillis() < suppressUntil) {
                return true;
            }
            
            // 记录抑制时间
            suppressedUntil.put(key, 
                System.currentTimeMillis() + suppressionDuration.toMillis());
            return false;
        }
    }
}

性能分析与诊断

反压监控

1. 反压检测

public class BackPressureMonitor {
    // 反压检测器
    public static class BackPressureDetector extends ProcessFunction<Event, BackPressureInfo> {
        private final long checkInterval = 10000; // 10秒
        private long lastCheckTime = 0;
        private long lastProcessedCount = 0;
        
        @Override
        public void processElement(Event event, Context ctx, 
                                 Collector<BackPressureInfo> out) throws Exception {
            long currentTime = System.currentTimeMillis();
            
            if (currentTime - lastCheckTime >= checkInterval) {
                // 计算处理速率
                long processedCount = ctx.currentWatermark();
                double processingRate = (processedCount - lastProcessedCount) / 
                                       ((currentTime - lastCheckTime) / 1000.0);
                
                // 获取输入队列大小
                int inputQueueSize = getInputQueueSize();
                
                // 计算反压比率
                double backPressureRatio = calculateBackPressureRatio(
                    processingRate, inputQueueSize);
                
                BackPressureInfo info = new BackPressureInfo()
                    .setTaskName(getRuntimeContext().getTaskName())
                    .setSubtaskIndex(getRuntimeContext().getIndexOfThisSubtask())
                    .setBackPressureRatio(backPressureRatio)
                    .setProcessingRate(processingRate)
                    .setInputQueueSize(inputQueueSize)
                    .setTimestamp(currentTime);
                
                out.collect(info);
                
                lastCheckTime = currentTime;
                lastProcessedCount = processedCount;
            }
        }
    }
    
    // 反压分析
    public static class BackPressureAnalyzer {
        public static BackPressureAnalysis analyze(List<BackPressureInfo> infos) {
            BackPressureAnalysis analysis = new BackPressureAnalysis();
            
            // 识别瓶颈算子
            Map<String, Double> avgBackPressure = infos.stream()
                .collect(Collectors.groupingBy(
                    BackPressureInfo::getTaskName,
                    Collectors.averagingDouble(BackPressureInfo::getBackPressureRatio)
                ));
            
            analysis.setBottleneckOperators(
                avgBackPressure.entrySet().stream()
                    .filter(e -> e.getValue() > 0.5)
                    .sorted(Map.Entry.<String, Double>comparingByValue().reversed())
                    .map(Map.Entry::getKey)
                    .collect(Collectors.toList())
            );
            
            // 计算整体反压水平
            double overallBackPressure = avgBackPressure.values().stream()
                .mapToDouble(Double::doubleValue)
                .average()
                .orElse(0.0);
            
            analysis.setOverallBackPressureLevel(overallBackPressure);
            
            // 生成优化建议
            analysis.setOptimizationSuggestions(generateSuggestions(analysis));
            
            return analysis;
        }
    }
}

延迟分析

1. 端到端延迟监控

public class LatencyMonitor {
    // 延迟追踪
    public static class LatencyTracker extends ProcessFunction<Event, Event> {
        private transient Histogram endToEndLatency;
        private transient Histogram processingLatency;
        private transient Counter lateEvents;
        
        @Override
        public void open(Configuration parameters) {
            MetricGroup metrics = getRuntimeContext().getMetricGroup();
            
            endToEndLatency = metrics.histogram("latency.end_to_end",
                new DescriptiveStatisticsHistogram(10000));
            processingLatency = metrics.histogram("latency.processing",
                new DescriptiveStatisticsHistogram(10000));
            lateEvents = metrics.counter("events.late");
        }
        
        @Override
        public void processElement(Event event, Context ctx, 
                                 Collector<Event> out) throws Exception {
            long currentTime = System.currentTimeMillis();
            long eventTime = event.getTimestamp();
            long processingStartTime = event.getProcessingStartTime();
            
            // 端到端延迟
            long e2eLatency = currentTime - eventTime;
            endToEndLatency.update(e2eLatency);
            
            // 处理延迟
            long procLatency = currentTime - processingStartTime;
            processingLatency.update(procLatency);
            
            // 迟到事件
            if (e2eLatency > 60000) { // 1分钟
                lateEvents.inc();
            }
            
            // 添加延迟信息
            event.setEndToEndLatency(e2eLatency);
            event.setProcessingLatency(procLatency);
            
            out.collect(event);
        }
    }
    
    // 延迟分析报告
    public static class LatencyAnalysisReport {
        private double avgEndToEndLatency;
        private double p99EndToEndLatency;
        private double maxEndToEndLatency;
        private Map<String, Double> operatorLatencies;
        private List<LatencyBottleneck> bottlenecks;
        
        public void generateReport(List<LatencyMetric> metrics) {
            // 计算统计信息
            DoubleSummaryStatistics stats = metrics.stream()
                .mapToDouble(LatencyMetric::getEndToEndLatency)
                .summaryStatistics();
            
            avgEndToEndLatency = stats.getAverage();
            maxEndToEndLatency = stats.getMax();
            
            // 计算 P99
            List<Double> latencies = metrics.stream()
                .map(LatencyMetric::getEndToEndLatency)
                .sorted()
                .collect(Collectors.toList());
            
            int p99Index = (int) (latencies.size() * 0.99);
            p99EndToEndLatency = latencies.get(p99Index);
            
            // 识别瓶颈
            identifyBottlenecks(metrics);
        }
    }
}

资源使用分析

1. 资源使用报告

public class ResourceAnalyzer {
    // 资源使用收集器
    public static class ResourceUsageCollector extends TimerTask {
        private final MetricRegistry metrics;
        private final List<ResourceSnapshot> snapshots = new ArrayList<>();
        
        @Override
        public void run() {
            ResourceSnapshot snapshot = new ResourceSnapshot();
            snapshot.setTimestamp(System.currentTimeMillis());
            
            // CPU 使用率
            OperatingSystemMXBean osBean = ManagementFactory.getOperatingSystemMXBean();
            if (osBean instanceof com.sun.management.OperatingSystemMXBean) {
                com.sun.management.OperatingSystemMXBean sunOsBean = 
                    (com.sun.management.OperatingSystemMXBean) osBean;
                snapshot.setCpuUsage(sunOsBean.getProcessCpuLoad() * 100);
                snapshot.setSystemCpuUsage(sunOsBean.getSystemCpuLoad() * 100);
            }
            
            // 内存使用
            MemoryMXBean memoryBean = ManagementFactory.getMemoryMXBean();
            MemoryUsage heapUsage = memoryBean.getHeapMemoryUsage();
            MemoryUsage nonHeapUsage = memoryBean.getNonHeapMemoryUsage();
            
            snapshot.setHeapUsed(heapUsage.getUsed());
            snapshot.setHeapMax(heapUsage.getMax());
            snapshot.setNonHeapUsed(nonHeapUsage.getUsed());
            snapshot.setNonHeapMax(nonHeapUsage.getMax());
            
            // 线程信息
            ThreadMXBean threadBean = ManagementFactory.getThreadMXBean();
            snapshot.setThreadCount(threadBean.getThreadCount());
            snapshot.setDaemonThreadCount(threadBean.getDaemonThreadCount());
            
            // GC 信息
            List<GarbageCollectorMXBean> gcBeans = 
                ManagementFactory.getGarbageCollectorMXBeans();
            for (GarbageCollectorMXBean gcBean : gcBeans) {
                snapshot.addGcInfo(gcBean.getName(), 
                    gcBean.getCollectionCount(), 
                    gcBean.getCollectionTime());
            }
            
            snapshots.add(snapshot);
            
            // 生成报告
            if (snapshots.size() >= 60) { // 1小时的数据
                generateResourceReport();
                snapshots.clear();
            }
        }
        
        private void generateResourceReport() {
            ResourceUsageReport report = new ResourceUsageReport();
            
            // CPU 分析
            double avgCpu = snapshots.stream()
                .mapToDouble(ResourceSnapshot::getCpuUsage)
                .average()
                .orElse(0.0);
            double maxCpu = snapshots.stream()
                .mapToDouble(ResourceSnapshot::getCpuUsage)
                .max()
                .orElse(0.0);
            
            report.setCpuAnalysis(new CpuAnalysis(avgCpu, maxCpu));
            
            // 内存分析
            long avgHeap = (long) snapshots.stream()
                .mapToLong(ResourceSnapshot::getHeapUsed)
                .average()
                .orElse(0.0);
            long maxHeap = snapshots.stream()
                .mapToLong(ResourceSnapshot::getHeapUsed)
                .max()
                .orElse(0L);
            
            report.setMemoryAnalysis(new MemoryAnalysis(avgHeap, maxHeap));
            
            // GC 分析
            Map<String, GcAnalysis> gcAnalysis = analyzeGC(snapshots);
            report.setGcAnalysis(gcAnalysis);
            
            // 发送报告
            sendReport(report);
        }
    }
}

实战案例

案例1:电商实时监控系统

public class ECommerceMonitoring {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 配置监控
        configureMetrics(env);
        
        // 订单流
        DataStream<Order> orders = env
            .addSource(new KafkaSource<Order>())
            .name("Order Source");
        
        // 业务监控
        DataStream<OrderMetrics> orderMetrics = orders
            .keyBy(Order::getUserId)
            .process(new OrderMonitoringFunction())
            .name("Order Monitoring");
        
        // 实时大屏数据
        orderMetrics
            .windowAll(TumblingProcessingTimeWindows.of(Time.seconds(10)))
            .aggregate(new DashboardAggregator())
            .addSink(new WebSocketSink())
            .name("Dashboard Sink");
        
        // 异常检测
        DataStream<Alert> alerts = orderMetrics
            .keyBy(OrderMetrics::getCategory)
            .process(new AnomalyDetector())
            .name("Anomaly Detection");
        
        // 告警通知
        alerts
            .filter(alert -> alert.getSeverity() == Severity.CRITICAL)
            .addSink(new AlertNotificationSink())
            .name("Alert Notification");
        
        env.execute("E-Commerce Monitoring");
    }
    
    // 订单监控函数
    public static class OrderMonitoringFunction 
            extends KeyedProcessFunction<String, Order, OrderMetrics> {
        
        // 业务指标
        private transient ValueState<UserOrderStats> userStatsState;
        
        // 监控指标
        private transient Counter orderCounter;
        private transient Meter orderRate;
        private transient Histogram orderAmount;
        private transient Gauge<Double> avgOrderValue;
        
        @Override
        public void open(Configuration parameters) throws Exception {
            // 初始化状态
            userStatsState = getRuntimeContext().getState(
                new ValueStateDescriptor<>("user-stats", UserOrderStats.class));
            
            // 初始化指标
            MetricGroup metrics = getRuntimeContext()
                .getMetricGroup()
                .addGroup("business")
                .addGroup("orders");
            
            orderCounter = metrics.counter("count");
            orderRate = metrics.meter("rate", new MeterView(60));
            orderAmount = metrics.histogram("amount",
                new DescriptiveStatisticsHistogram(1000));
            
            metrics.gauge("avg_value", new Gauge<Double>() {
                @Override
                public Double getValue() {
                    try {
                        UserOrderStats stats = userStatsState.value();
                        if (stats != null && stats.orderCount > 0) {
                            return stats.totalAmount / stats.orderCount;
                        }
                    } catch (Exception e) {
                        LOG.error("Error getting average order value", e);
                    }
                    return 0.0;
                }
            });
        }
        
        @Override
        public void processElement(Order order, Context ctx, 
                                 Collector<OrderMetrics> out) throws Exception {
            // 更新用户统计
            UserOrderStats stats = userStatsState.value();
            if (stats == null) {
                stats = new UserOrderStats();
            }
            stats.orderCount++;
            stats.totalAmount += order.getAmount();
            stats.lastOrderTime = order.getTimestamp();
            userStatsState.update(stats);
            
            // 更新监控指标
            orderCounter.inc();
            orderRate.markEvent();
            orderAmount.update((long) (order.getAmount() * 100)); // 转为分
            
            // 生成监控数据
            OrderMetrics metrics = new OrderMetrics()
                .setUserId(order.getUserId())
                .setOrderId(order.getOrderId())
                .setAmount(order.getAmount())
                .setCategory(order.getCategory())
                .setTimestamp(order.getTimestamp())
                .setUserOrderCount(stats.orderCount)
                .setUserTotalAmount(stats.totalAmount);
            
            out.collect(metrics);
            
            // 定时清理过期用户数据
            ctx.timerService().registerProcessingTimeTimer(
                ctx.timerService().currentProcessingTime() + 86400000); // 24小时
        }
        
        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, 
                          Collector<OrderMetrics> out) throws Exception {
            UserOrderStats stats = userStatsState.value();
            if (stats != null && 
                timestamp - stats.lastOrderTime > 86400000) { // 24小时无订单
                userStatsState.clear();
            }
        }
    }
}

案例2:IoT 设备监控平台

public class IoTDeviceMonitoring {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // 设备数据流
        DataStream<DeviceData> deviceData = env
            .addSource(new MqttSource())
            .assignTimestampsAndWatermarks(
                WatermarkStrategy.<DeviceData>forBoundedOutOfOrderness(Duration.ofSeconds(5))
                    .withTimestampAssigner((data, timestamp) -> data.getTimestamp())
            );
        
        // 设备状态监控
        DataStream<DeviceStatus> deviceStatus = deviceData
            .keyBy(DeviceData::getDeviceId)
            .process(new DeviceStatusMonitor())
            .name("Device Status Monitor");
        
        // 异常检测
        DataStream<DeviceAlert> alerts = deviceStatus
            .keyBy(DeviceStatus::getDeviceId)
            .process(new DeviceAnomalyDetector())
            .name("Anomaly Detector");
        
        // 设备组聚合
        DataStream<GroupMetrics> groupMetrics = deviceStatus
            .keyBy(DeviceStatus::getGroupId)
            .window(TumblingEventTimeWindows.of(Time.minutes(5)))
            .aggregate(new GroupMetricsAggregator())
            .name("Group Metrics");
        
        // 监控大屏
        groupMetrics
            .addSink(new DashboardSink())
            .name("Dashboard Sink");
        
        // 告警处理
        alerts
            .process(new AlertProcessor())
            .addSink(new MultiChannelAlertSink())
            .name("Alert Sink");
        
        env.execute("IoT Device Monitoring");
    }
    
    // 设备状态监控
    public static class DeviceStatusMonitor 
            extends KeyedProcessFunction<String, DeviceData, DeviceStatus> {
        
        private transient ValueState<DeviceInfo> deviceInfoState;
        private transient ListState<DeviceData> recentDataState;
        
        // 监控指标
        private transient Gauge<Integer> onlineDevices;
        private transient Counter dataPoints;
        private transient Histogram temperature;
        private transient Histogram humidity;
        
        @Override
        public void open(Configuration parameters) throws Exception {
            // 状态初始化
            deviceInfoState = getRuntimeContext().getState(
                new ValueStateDescriptor<>("device-info", DeviceInfo.class));
            recentDataState = getRuntimeContext().getListState(
                new ListStateDescriptor<>("recent-data", DeviceData.class));
            
            // 指标初始化
            MetricGroup metrics = getRuntimeContext()
                .getMetricGroup()
                .addGroup("iot")
                .addGroup("devices");
            
            dataPoints = metrics.counter("data_points");
            temperature = metrics.histogram("temperature",
                new DescriptiveStatisticsHistogram(1000));
            humidity = metrics.histogram("humidity",
                new DescriptiveStatisticsHistogram(1000));
            
            // 在线设备数
            metrics.gauge("online_count", new Gauge<Integer>() {
                @Override
                public Integer getValue() {
                    // 实际实现需要跨所有并行实例统计
                    return getOnlineDeviceCount();
                }
            });
        }
        
        @Override
        public void processElement(DeviceData data, Context ctx, 
                                 Collector<DeviceStatus> out) throws Exception {
            // 更新设备信息
            DeviceInfo info = deviceInfoState.value();
            if (info == null) {
                info = new DeviceInfo();
                info.deviceId = data.getDeviceId();
                info.firstSeen = data.getTimestamp();
            }
            info.lastSeen = data.getTimestamp();
            info.dataCount++;
            deviceInfoState.update(info);
            
            // 保存最近数据
            recentDataState.add(data);
            
            // 更新指标
            dataPoints.inc();
            if (data.getTemperature() != null) {
                temperature.update((long) (data.getTemperature() * 10));
            }
            if (data.getHumidity() != null) {
                humidity.update((long) (data.getHumidity() * 10));
            }
            
            // 计算设备状态
            DeviceStatus status = calculateDeviceStatus(data, info);
            out.collect(status);
            
            // 注册心跳检测定时器
            ctx.timerService().registerEventTimeTimer(
                data.getTimestamp() + 300000); // 5分钟超时
        }
        
        @Override
        public void onTimer(long timestamp, OnTimerContext ctx, 
                          Collector<DeviceStatus> out) throws Exception {
            DeviceInfo info = deviceInfoState.value();
            if (info != null && timestamp - info.lastSeen > 300000) {
                // 设备离线
                DeviceStatus status = new DeviceStatus()
                    .setDeviceId(info.deviceId)
                    .setStatus(DeviceStatus.Status.OFFLINE)
                    .setTimestamp(timestamp);
                
                out.collect(status);
            }
        }
    }
}

最佳实践

1. 监控设计原则

public class MonitoringBestPractices {
    // 1. 分层监控
    public static class LayeredMonitoring {
        // 基础设施层
        private void monitorInfrastructure() {
            // CPU、内存、磁盘、网络
        }
        
        // 平台层
        private void monitorPlatform() {
            // Flink 集群状态、资源使用
        }
        
        // 应用层
        private void monitorApplication() {
            // 作业状态、数据流量、处理延迟
        }
        
        // 业务层
        private void monitorBusiness() {
            // 业务指标、数据质量、SLA
        }
    }
    
    // 2. 指标命名规范
    public static class MetricNamingConvention {
        // 格式:<namespace>.<component>.<metric>
        private static final String METRIC_PATTERN = 
            "flink.%s.%s.%s"; // flink.taskmanager.memory.heap_used
        
        public static String buildMetricName(String component, 
                                           String subComponent, 
                                           String metric) {
            return String.format(METRIC_PATTERN, 
                component.toLowerCase(),
                subComponent.toLowerCase(),
                metric.toLowerCase().replace(" ", "_"));
        }
    }
    
    // 3. 采样策略
    public static class SamplingStrategy {
        private final double samplingRate;
        private final Random random = new Random();
        
        public boolean shouldSample() {
            return random.nextDouble() < samplingRate;
        }
        
        public <T> void recordMetric(T value, Consumer<T> recorder) {
            if (shouldSample()) {
                recorder.accept(value);
            }
        }
    }
}

2. 监控优化

public class MonitoringOptimization {
    // 1. 异步指标上报
    public static class AsyncMetricReporter {
        private final BlockingQueue<MetricEvent> queue = 
            new LinkedBlockingQueue<>(10000);
        private final ExecutorService executor = 
            Executors.newSingleThreadExecutor();
        
        public void reportMetric(String name, double value) {
            MetricEvent event = new MetricEvent(name, value, 
                System.currentTimeMillis());
            
            if (!queue.offer(event)) {
                // 队列满,丢弃或告警
                LOG.warn("Metric queue is full, dropping metric: " + name);
            }
        }
        
        private void processMetrics() {
            List<MetricEvent> batch = new ArrayList<>();
            
            while (true) {
                try {
                    // 批量处理
                    queue.drainTo(batch, 100);
                    if (!batch.isEmpty()) {
                        sendBatch(batch);
                        batch.clear();
                    } else {
                        Thread.sleep(100);
                    }
                } catch (Exception e) {
                    LOG.error("Error processing metrics", e);
                }
            }
        }
    }
    
    // 2. 指标缓存
    public static class MetricCache {
        private final Map<String, CachedMetric> cache = 
            new ConcurrentHashMap<>();
        private final long cacheTimeout = 1000; // 1秒
        
        public double getMetricValue(String name, Supplier<Double> calculator) {
            CachedMetric cached = cache.get(name);
            
            if (cached == null || 
                System.currentTimeMillis() - cached.timestamp > cacheTimeout) {
                double value = calculator.get();
                cache.put(name, new CachedMetric(value, System.currentTimeMillis()));
                return value;
            }
            
            return cached.value;
        }
    }
}

常见问题

Q1: 监控对性能的影响如何控制?

解答:

  1. 使用采样:对高频指标进行采样
  2. 异步上报:避免阻塞主流程
  3. 批量发送:减少网络开销
  4. 本地聚合:减少数据量
  5. 合理的上报间隔:平衡实时性和性能

Q2: 如何设计告警规则避免告警风暴?

解答:

  1. 告警聚合:相同类型告警合并
  2. 告警抑制:避免重复告警
  3. 告警升级:分级处理
  4. 告警静默:维护期间静默
  5. 智能告警:基于机器学习的异常检测

Q3: 监控数据存储和查询优化?

解答:

# 数据保留策略
- 原始数据:保留 7 天
- 5 分钟聚合:保留 30 天
- 1 小时聚合:保留 90 天
- 1 天聚合:保留 1 年

# 查询优化
- 使用时间索引
- 预聚合常用指标
- 缓存查询结果
- 使用采样查询

解答:

// SQL 作业特殊指标
tableEnv.getConfig().getConfiguration().setString(
    "table.exec.mini-batch.enabled", "true");
tableEnv.getConfig().getConfiguration().setString(
    "table.exec.mini-batch.allow-latency", "5 s");
tableEnv.getConfig().getConfiguration().setString(
    "table.exec.mini-batch.size", "5000");

// 注册 SQL 专用指标
tableEnv.createTemporaryView("metrics_view", 
    "SELECT " +
    "  COUNT(*) as total_count, " +
    "  AVG(amount) as avg_amount, " +
    "  MAX(amount) as max_amount " +
    "FROM orders " +
    "GROUP BY TUMBLE(rowtime, INTERVAL '1' MINUTE)");

Q5: 分布式追踪如何实现?

解答:

// 集成 OpenTracing
public class DistributedTracing {
    private final Tracer tracer = GlobalTracer.get();
    
    @Override
    public void processElement(Event event, Context ctx, Collector<Event> out) {
        Span span = tracer.buildSpan("processEvent")
            .withTag("event.id", event.getId())
            .withTag("operator", getRuntimeContext().getTaskName())
            .start();
        
        try (Scope scope = tracer.activateSpan(span)) {
            // 处理逻辑
            processEvent(event);
            out.collect(event);
        } finally {
            span.finish();
        }
    }
}

相关文章

监控相关

基础知识

集成方案

相关技术


总结

Flink 监控体系是保障流处理应用稳定运行的基础设施,需要从多个维度进行建设:

🎯 核心要点

  1. 全面的指标体系:覆盖基础设施、平台、应用和业务层面
  2. 实时的监控能力:及时发现问题,快速定位原因
  3. 智能的告警机制:准确告警,避免告警风暴
  4. 直观的可视化:清晰展示系统状态和趋势
  5. 深入的分析能力:支持问题诊断和性能优化

✅ 最佳实践

  • 建立分层的监控体系
  • 使用标准的命名规范
  • 实施合理的采样策略
  • 配置智能的告警规则
  • 保持监控系统的高可用
  • 定期审查和优化监控指标

🚀 实施建议

  1. 循序渐进:先建立基础监控,再逐步完善
  2. 重点突出:优先监控关键业务指标
  3. 持续优化:根据实际情况调整监控策略
  4. 自动化运维:结合监控实现自动化响应

🚫 注意事项:

  • 避免过度监控影响性能
  • 合理设置数据保留期限
  • 注意监控数据的安全性
  • 定期清理无用的监控指标

通过完善的监控体系,可以确保 Flink 应用在生产环境中稳定、高效地运行,为业务提供可靠的实时数据处理服务。

上一章 / 下一章