概述
在生产环境中运行 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: 监控对性能的影响如何控制?
解答:
- 使用采样:对高频指标进行采样
- 异步上报:避免阻塞主流程
- 批量发送:减少网络开销
- 本地聚合:减少数据量
- 合理的上报间隔:平衡实时性和性能
Q2: 如何设计告警规则避免告警风暴?
解答:
- 告警聚合:相同类型告警合并
- 告警抑制:避免重复告警
- 告警升级:分级处理
- 告警静默:维护期间静默
- 智能告警:基于机器学习的异常检测
Q3: 监控数据存储和查询优化?
解答:
# 数据保留策略
- 原始数据:保留 7 天
- 5 分钟聚合:保留 30 天
- 1 小时聚合:保留 90 天
- 1 天聚合:保留 1 年
# 查询优化
- 使用时间索引
- 预聚合常用指标
- 缓存查询结果
- 使用采样查询
Q4: 如何监控 Flink SQL 作业?
解答:
// 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 性能优化 - 性能调优指南
- Checkpoint & Savepoint - 容错监控
- Flink CDC - CDC 监控方案
基础知识
- Flink 架构与工作流程 - 理解监控点
- Flink DataStream API - API 使用
- Flink DataStream API 高级用法 - 高级特性
集成方案
- Flink + Kafka、HBase、Iceberg - 集成监控
- Flink Table & SQL API - SQL 监控
相关技术
- Prometheus 官方文档 - Prometheus 使用
- Grafana 官方文档 - Grafana 配置
- Elasticsearch + Kibana - ELK 监控
总结
Flink 监控体系是保障流处理应用稳定运行的基础设施,需要从多个维度进行建设:
🎯 核心要点
- 全面的指标体系:覆盖基础设施、平台、应用和业务层面
- 实时的监控能力:及时发现问题,快速定位原因
- 智能的告警机制:准确告警,避免告警风暴
- 直观的可视化:清晰展示系统状态和趋势
- 深入的分析能力:支持问题诊断和性能优化
✅ 最佳实践
- 建立分层的监控体系
- 使用标准的命名规范
- 实施合理的采样策略
- 配置智能的告警规则
- 保持监控系统的高可用
- 定期审查和优化监控指标
🚀 实施建议
- 循序渐进:先建立基础监控,再逐步完善
- 重点突出:优先监控关键业务指标
- 持续优化:根据实际情况调整监控策略
- 自动化运维:结合监控实现自动化响应
🚫 注意事项:
- 避免过度监控影响性能
- 合理设置数据保留期限
- 注意监控数据的安全性
- 定期清理无用的监控指标
通过完善的监控体系,可以确保 Flink 应用在生产环境中稳定、高效地运行,为业务提供可靠的实时数据处理服务。