【Flink反压排查】别再一上来就加机器了
生产环境遇到反压,别急着加机器。按下面的顺序来,先治下游、再控上游,最后才是扩容。这篇文章按这个优先级组织,看完直接就能用。
一、先搞清楚:反压是怎么一回事
Flink是一条流水线,数据从上往下流。下游某个环节处理不过来了,Buffer被填满,于是告诉上游"我满了别发了",这个信号一层一层往上传,最后连Source端从Kafka消费的速度都被迫降下来——这就是反压。
用堵车来理解:前面路口堵了,车排起长队,这个拥堵往回传,最后你连小区都出不来。Flink的反压就是这个理,堵点在下游,影响往上传。
现在生产环境主流是Flink 1.17+,反压基于Credit-Based流控。简单说:下游给上游发"信用值"(Credit),告诉上游自己还能收多少条。下游满了就给0,上游就不发了,等着。这个机制是自动的,不需要你配置。
| 算子 | 角色 | Credit 状态 | 行为 |
|---|---|---|---|
| Source | 消费 Kafka | 收到 Credit=0 | 停止拉取,等待 |
| Window | 聚合计算 | Credit=0 传递给上游 | 自身 Buffer 满了,不再接收 |
| Sink | 写入 ClickHouse | Buffer 满 | 触发反压的起点 |
Credit 的流向:Sink Buffer 满 → Sink 给 Window 发 Credit=0 → Window 给 Source 发 Credit=0 → Source 降速。反压从下游往上游逐级传导。
反压会带来的连锁反应
反压本身是保护机制,但持续时间长了会触发一系列问题:
| 问题 | 具体表现 | 严重性 |
|---|---|---|
| Checkpoint超时 | Barrier流动慢,Checkpoint做不完,最终失败 | 致命 |
| Kafka消费滞后 | Source消费慢,Consumer Lag越来越大,端到端延迟增大 | 高 |
| 状态膨胀 | 窗口数据积压在RocksDB,最终OOM | 致命 |
| 数据延迟 | 端到端延迟从秒级飙到分钟级甚至小时级 | 高 |
| TM掉线 | 严重GC导致心跳超时,TaskManager被JobManager踢掉 | 致命 |
最要命的是Checkpoint失败。Checkpoint连续超时后,Flink无法生成新的快照,一旦作业挂了,只能从很久之前的Checkpoint恢复,期间数据全部重放,雪崩。
二、发现反压:Flink Web UI
2.1 看颜色
打开Web UI → Job Graph,每个算子方块的颜色直接告诉你状态:
| 颜色 | 状态 | 含义 |
|---|---|---|
| 蓝色 | IDLE | 没活干,闲置 |
| 红色 | BUSY | 在干活,正常 |
| 黑色 | BACKPRESSURED | 被堵住了,反压中 |
实际生产中看到的多半是混合色。比如红黑色,说明这个算子一部分时间干活、一部分时间等下游消化。
定位瓶颈口诀:从Source往Sink找,第一个出现BUSY/红色且下游是黑色的算子,就是瓶颈。
2.2 Backpressure Tab
点进具体算子 → Backpressure Tab,看采样结果:
| 状态 | 反压比例 | 处理建议 |
|---|---|---|
| OK | < 10% | 正常,不用管 |
| LOW | 10%~50% | 偶尔出现没关系,持续观察 |
| HIGH | > 50% | 持续存在就要排查 |
踩坑点:反压采样基于线程堆栈分析,本身有性能开销。大促高峰期别频繁点击采样,本来就卡,一点更卡。
三、解决优先级总览(核心!按这个顺序来)
生产环境遇到反压,处理顺序很重要。不能上来就加机器,也不能只盯着下游不管上游。完整顺序如下:
Step 1: 定位瓶颈(先找到谁慢了)
│
▼
Step 2: 下游处理侧优化(治本——优先做)
├── 2.1 数据倾斜治理
├── 2.2 算子代码优化
├── 2.3 状态/内存优化
└── 2.4 外部系统调优
│
▼
Step 3: Source端控制(加保险——配合做)
├── 3.1 Kafka Consumer参数调优
└── 3.2 启动追数限速
│
▼
Step 4: 作业配置调优(锦上添花)
├── 4.1 网络Buffer调优
├── 4.2 Checkpoint参数调整
└── 4.3 Operator Chain策略
│
▼
Step 5: 扩容(最后手段)
├── 5.1 增加TaskManager资源
└── 5.2 提升Source并行度
为什么要这个顺序?
- Step 2是治本:下游处理能力上不去,你光控上游有啥用?数据迟早要处理,瓶颈永远在下游。
- Step 3是保护:下游治得差不多了,但如果Kafka积压几个亿的数据,一重启又冲垮了。上游限速就是防止这种情况。
- Step 4是调优:前两步做完了还有余量,再调网络Buffer、Checkpoint这些参数,榨干性能。
- Step 5是兜底:以上全做完了还撑不住,才是真需要加资源。
四、Step 1:定位瓶颈——四步排查法
发现反压后,按下面四步找到根因:
4.1 数据倾斜排查
看Web UI里每个SubTask的Records Received/Sent。如果某个SubTask的数据量是其他的5倍以上,就是数据倾斜。
SubTask 0: Received 10,000,000 records ◀── 热点Key
SubTask 1: Received 8,000 records
SubTask 2: Received 7,500 records
SubTask 3: Received 8,200 records
倾斜的本质:keyBy之后某个Key的数据量特别大,全路由到一个Task上了。比如按用户ID聚合,有个大V的数据量是普通用户的1000倍。
4.2 火焰图排查
Flink 1.13+ Web UI内置JVM CPU火焰图。横轴是方法名,宽度=CPU占用比例。重点关注你自己的UDF函数。
常见代码层性能杀手:
| 问题 | 火焰图表现 | 解决 |
|---|---|---|
| 正则表达式灾难 | Pattern.match占很宽 |
预编译Pattern,简化正则 |
| 循环里查数据库 | JDBC相关方法占很宽 |
改Async I/O + 旁路缓存 |
| 大对象频繁创建 | new和GC占很宽 |
对象复用 |
| 复杂JSON解析 | Jackson/Fastjson占很宽 |
简化结构或用Avro/Protobuf |
| 同步HTTP调用 | Socket.read阻塞占很宽 |
必须用Async I/O |
4.3 GC日志排查
Flink默认G1 GC,看GC日志关注这几个指标:
| 指标 | 健康值 | 异常表现 |
|---|---|---|
| Young GC频率 | < 5秒/次 | 太频繁说明Eden区太小 |
| Full GC频率 | 尽量不出现 | 出现基本就会反压 |
| 单次GC停顿 | < 200ms | 太长会阻塞数据处理 |
4.4 外部系统排查
代码没问题、GC正常,瓶颈就在Flink之外:
| 外部系统 | 排查命令/指标 |
|---|---|
| Kafka | 看Broker requestLatency,Partition是否太少 |
| HBase | 看Region flush频率、compaction队列深度 |
| ClickHouse | 看system.merges,Merge是否跟不上 |
| MySQL | 看show processlist(锁等待)、表索引数量(索引多写入慢)、实例CPU/IO负载 |
| Redis | 看slowlog、大Key淘汰情况 |
踩坑点:外部系统的问题是间歇性的。HBase Compaction、ClickHouse Merge都是后台异步的,你看的时候好好的,过会儿又不行了。
五、Step 2:下游处理侧优化(治本,优先做)
5.1 数据倾斜治理
数据倾斜是生产环境反压的头号原因。
方案A:两阶段聚合(最常用)
先加随机前缀打散,局部聚合后再全局聚合:
// 第一步:加随机前缀打散
SingleOutputStreamOperator<Order> prefixed = orders
.map(new RichMapFunction<Order, Order>() {
private Random random;
@Override
public void open(Configuration parameters) {
random = new Random();
}
@Override
public Order map(Order order) {
// 10个随机前缀,把热点Key拆成10份
String prefix = String.valueOf(random.nextInt(10));
order.setCategoryId(prefix + "#" + order.getCategoryId());
return order;
}
});
// 第二步:按"前缀+原Key"局部聚合
DataStream<Order> localAgg = prefixed
.keyBy(Order::getCategoryId)
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.aggregate(new SumAggFunction());
// 第三步:去掉前缀,按原Key全局聚合
DataStream<Order> globalAgg = localAgg
.map(order -> {
String realKey = order.getCategoryId().split("#")[1];
order.setCategoryId(realKey);
return order;
})
.keyBy(Order::getCategoryId)
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.aggregate(new SumAggFunction());
方案B:rebalance强制重分区
倾斜不严重时,在keyBy之前加rebalance():
stream.rebalance().keyBy(Order::getCategoryId)
方案C:热点Key单独处理
已经知道哪个Key是热点,单独捞出来:
// 拆成两条流处理
OutputTag<Order> hotTag = new OutputTag<Order>("hot"){};
SingleOutputStreamOperator<Order> normalStream = orders
.process(new ProcessFunction<Order, Order>() {
@Override
public void processElement(Order order, Context ctx, Collector<Order> out) {
if ("cat_1001".equals(order.getCategoryId())) {
ctx.output(hotTag, order); // 热点数据旁路
} else {
out.collect(order); // 普通数据正常走
}
}
});
// 热点数据单独用更高的并行度处理
DataStream<Order> hotStream = normalStream
.getSideOutput(hotTag)
.keyBy(Order::getCategoryId)
.window(TumblingProcessingTimeWindows.of(Time.seconds(5)))
.aggregate(new SumAggFunction())
.setParallelism(8); // 给热点单独加并行度
踩坑点:
rebalance()会触发全量数据Shuffle,网络开销大。数据量巨大时慎用,可能反压没解决网络先挂了。两阶段聚合是多一次计算,但能从根本上解决倾斜,生产环境最常用。
5.2 算子代码优化
Async I/O(查维表、调外部API必用)
// 异步查询Redis维表
AsyncFunction<Order, EnrichedOrder> asyncFunc =
new RichAsyncFunction<Order, EnrichedOrder>() {
private transient JedisPool jedisPool;
@Override
public void open(Configuration parameters) {
jedisPool = new JedisPool(new JedisPoolConfig(), "redis", 6379);
}
@Override
public void asyncInvoke(Order order, ResultFuture<EnrichedOrder> resultFuture) {
CompletableFuture.supplyAsync(() -> {
try (Jedis jedis = jedisPool.getResource()) {
String info = jedis.get("user:" + order.getUserId());
return new EnrichedOrder(order, info);
}
}).thenAccept(r -> resultFuture.complete(Collections.singletonList(r)));
}
};
DataStream<EnrichedOrder> enriched = AsyncDataStream
.unorderedWait(
orders,
asyncFunc,
1000, // 超时1秒
TimeUnit.MILLISECONDS,
100 // 最大并发100
);
踩坑点:
unorderedWait会乱序,orderedWait保证顺序但性能差。大部分业务可以容忍短暂乱序,用unorderedWait。并发数别设太大,100起步,根据外部系统QPS能力慢慢调,别把外部系统打挂。
对象复用
// 别这样:每来一条数据new一个大对象
.map(json -> new ObjectMapper().readValue(json, Order.class)) // 灾难!
// 要这样:提前创建,复用
.map(new RichMapFunction<String, Order>() {
private transient ObjectMapper mapper;
@Override
public void open(Configuration parameters) {
mapper = new ObjectMapper(); // 只创建一次
}
@Override
public Order map(String json) {
return mapper.readValue(json, Order.class); // 复用
}
})
5.3 状态/内存优化
RocksDB调优(大状态必做)
# flink-conf.yaml
state.backend: rocksdb
state.backend.rocksdb.memory.managed: true
taskmanager.memory.managed.fraction: 0.35
| 参数 | 作用 | 建议值 |
|---|---|---|
managed.fraction |
托管内存占TM总内存比例 | 0.3 ~ 0.5 |
rocksdb.memory.managed |
RocksDB自动从托管内存分配 | true |
state.backend.incremental |
增量Checkpoint,只改差异 | true |
Watermark和窗口优化
// 别设太大,够用就行
WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5)) // 5秒够了
// .forBoundedOutOfOrderness(Duration.ofMinutes(30)) // 30分钟太长了!
// 窗口粒度别太细,数据量大用滚动窗口
.window(TumblingProcessingTimeWindows.of(Time.minutes(5))) // 5分钟窗口
踩坑点:Watermark设太大,窗口一直不触发,数据积压在State里,RocksDB压力山大。一般5~30秒够了。窗口粒度太细(比如1秒一个窗口),State频繁增删,性能差。
5.4 外部系统调优
| 外部系统 | 优化手段 |
|---|---|
| ClickHouse | 批量写入5000~10000条/批次;关注Parts数量和Merge进度 |
| HBase | 批量Put,启用WAL异步刷写,预分区避免Region热点 |
| MySQL | 连接池(HikariCP),批量Insert,考虑异步写入或消息队列缓冲 |
| Kafka | 确保Partition数 >= Source并行度,Producer端批量发送 |
ClickHouse Sink批量参数参考:
new JdbcExecutionOptions.Builder()
.withBatchSize(5000) // 每5000条写一次
.withBatchIntervalMs(200) // 最多等200ms
.withMaxRetries(3)
.build()
踩坑点:ClickHouse单次写入太大(比如10万条)会生成大量Parts,后台Merge跟不上,查询性能暴跌。5000~10000是经验值。
六、Step 3:Source端控制(加保险,配合做)
先搞清楚:反压治理的边界在哪里
很多链路是 CDC → Kafka → Flink,反压了要不要治理 CDC?
不需要。Flink 反压的治理边界就到 Flink 的 Source 端(Kafka Consumer),Kafka 再上游不需要管。
因为 Kafka 的写入和消费是解耦的:
| 组件 | 角色 | 反压时行为 | 是否需要治理 |
|---|---|---|---|
| CDC | Kafka Producer | 不受影响,继续写入 | ❌ |
| Kafka | 存储缓冲 | 消息堆积,磁盘占用由写入量和 retention 决定 | ❌ |
| Flink Source | Kafka Consumer | Credit=0,自动降速 | ✅ |
| Flink Sink | 外部系统写入 | 处理不过来,反压起点 | ✅ |
关键认知:Kafka 写入和消费是解耦的。Flink 消费慢不影响 CDC 写入,Flink 反压的治理边界到 Source 算子为止。
Flink 消费慢了 → Kafka Consumer 降速 → Kafka 消息堆积。
但 Kafka 的 Producer(CDC)不受影响,继续往里写。所以 CDC 不在 Flink 反压的治理范围内。
Source 端需要控制的三种场景
下游优化是治本,但生产环境还要防止上游"冲垮"下游。特别是以下三种场景:
| 场景 | 问题 | 后果 |
|---|---|---|
| Kafka积压后重启 | 停了2小时,堆积上亿条,一重启全速消费 | 下游瞬间被打爆,反压传导,Checkpoint超时 |
| 业务高峰突发流量 | 大促开始,订单量瞬间翻10倍 | Source拉取远快于处理,全链路反压 |
| 冷启动追历史数据 | 新作业第一次启动要追数 | 和积压重启类似,大量数据瞬间涌入 |
6.1 Kafka Consumer参数限速
最基础、最安全的Source端限速手段。调整这几个参数:
Properties props = new Properties();
props.setProperty("bootstrap.servers", "kafka:9092");
props.setProperty("group.id", "flink-order-group");
// ▼▼▼ 核心限速参数 ▼▼▼
// 每次poll最多拉多少条。反压严重时调小,Source自然慢下来
props.setProperty("max.poll.records", "200"); // 默认500,反压时降到100~200
// 单次拉取的最大字节数,防止大包
props.setProperty("max.partition.fetch.bytes", "1048576"); // 1MB
// 最少攒多少字节才返回,配合等待时间降低拉取频率
props.setProperty("fetch.min.bytes", "1024");
props.setProperty("fetch.max.wait.ms", "500");
FlinkKafkaConsumer<Order> source = new FlinkKafkaConsumer<>(
"order-topic",
new OrderDeserializationSchema(),
props
);
max.poll.records是最核心的参数。设200就是每次最多拉200条,设500就是500条。下游撑不住的时候,把这个调小,立竿见影。
踩坑点:
max.poll.records别设太小(比如10条),会导致Consumer频繁空轮询,CPU飙高。反压紧急时降到100~200,等恢复了再调回500。
6.2 启动追数时动态限速(生产必做)
作业重启后Kafka积压了几亿条数据,如果全速消费,下游直接被冲垮。启动时先慢速,等延迟降了再放开。
/**
* 自适应限速Source:启动时低速,延迟降下来后逐步提速
*/
public class AdaptiveRateKafkaSource extends RichParallelSourceFunction<Order> {
private volatile boolean isRunning = true;
// 速率范围
private static final long MIN_RATE = 1000; // 最低1000条/秒
private static final long MAX_RATE = 50000; // 最高50000条/秒
private static final long TARGET_LAG_MS = 5000; // 目标延迟5秒
// 检查间隔
private static final long CHECK_INTERVAL_MS = 30000; // 30秒检查一次
private transient RateLimiter rateLimiter;
private transient KafkaConsumer<String, String> consumer;
private long currentRate;
private long lastCheckTime;
@Override
public void open(Configuration parameters) {
currentRate = MIN_RATE; // 启动时从最低速开始
rateLimiter = RateLimiter.create(currentRate);
Properties props = new Properties();
props.setProperty("bootstrap.servers", "kafka:9092");
props.setProperty("group.id", "flink-order-group");
props.setProperty("max.poll.records", "500");
props.setProperty("auto.offset.reset", "latest");
props.setProperty("enable.auto.commit", "false");
consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("order-topic"));
}
@Override
public void run(SourceContext<Order> ctx) throws Exception {
lastCheckTime = System.currentTimeMillis();
while (isRunning) {
// 限速:获取令牌,没有就阻塞等待
rateLimiter.acquire();
ConsumerRecords<String, String> records =
consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
ctx.collectWithTimestamp(
parseOrder(record.value()),
record.timestamp()
);
}
// 定期调整速率
long now = System.currentTimeMillis();
if (now - lastCheckTime > CHECK_INTERVAL_MS) {
adjustRate();
lastCheckTime = now;
}
}
}
private void adjustRate() {
// 获取当前消费延迟
long currentLag = getMaxLagMs();
if (currentLag < TARGET_LAG_MS && currentRate < MAX_RATE) {
// 延迟达标,提速(每次翻倍,但不超过上限)
currentRate = Math.min(currentRate * 2, MAX_RATE);
rateLimiter.setRate(currentRate);
LOG.info("Lag is good ({}ms), increasing rate to {}/s", currentLag, currentRate);
} else if (currentLag > TARGET_LAG_MS * 3 && currentRate > MIN_RATE) {
// 延迟太大,降速
currentRate = Math.max(currentRate / 2, MIN_RATE);
rateLimiter.setRate(currentRate);
LOG.warn("Lag is high ({}ms), decreasing rate to {}/s", currentLag, currentRate);
}
}
private long getMaxLagMs() {
// 通过consumer.metrics()获取当前消费延迟
// 实际生产可以用Kafka AdminClient查询
Map<MetricName, ? extends Metric> metrics = consumer.metrics();
long maxLag = 0;
for (Map.Entry<MetricName, ? extends Metric> entry : metrics.entrySet()) {
if (entry.getKey().name().equals("records-lag-max")) {
maxLag = ((Number) entry.getValue().metricValue()).longValue();
}
}
return maxLag;
}
@Override
public void cancel() {
isRunning = false;
if (consumer != null) {
consumer.close();
}
}
}
更简单的实操方案:不用写自定义Source,直接分阶段调整max.poll.records:
作业重启时:
1. 先把 max.poll.records 设成 100(低速消费)
2. 观察 Flink Web UI,看延迟是否在下降
3. 延迟降到目标值后,逐步调大到 200 → 300 → 500
4. 恢复正常值
这个方案不用改代码,直接改配置重启,线上操作最安全。
6.3 Source并行度与Kafka Partition匹配
| 情况 | 问题 | 解决方案 |
|---|---|---|
| Partition数 < Source并行度 | 部分SubTask永远空闲,消费上限被Partition卡死 | 降低Source并行度,或增加Partition |
| Partition数 > Source并行度 | 一个SubTask处理多个Partition,可能处理不过来 | 增加Source并行度 |
| Partition数 = Source并行度 | 最理想,一对一消费 | 保持 |
最佳实践:Partition数是Source并行度的整数倍。比如Source并行度=4,Partition=4/8/12。
// 设置Source并行度
FlinkKafkaConsumer<Order> source = new FlinkKafkaConsumer<>(...);
source.setStartFromGroupOffsets();
env.addSource(source)
.setParallelism(4); // 跟Kafka Partition数匹配
踩坑点:很多新手Source并行度设得特别大(比如16),但Kafka Partition只有4。结果是12个SubTask空转,另外4个拼命消费也上不去。先看看Partition数再设并行度。
七、Step 4:作业配置调优
前两步做完,如果还有余量,再调这些参数榨取性能。
7.1 网络Buffer调优
# flink-conf.yaml
# 网络内存占TM总内存的比例,默认0.1
taskmanager.network.memory.fraction: 0.15
# Flink 1.14+ 开启Buffer自动调节,让系统自己根据反压情况调整
taskmanager.network.memory.buffer-debloat.enabled: true
# Buffer Debloat的目标时间——数据应该在Buffer里停留多久(毫秒)
taskmanager.network.memory.buffer-debloat.period: 500
buffer-debloat开启后,Flink会自动调节Buffer大小:反压严重时减小Buffer让数据更快流动,畅通时增大Buffer提高吞吐。生产环境建议开启,省得手动调。
7.2 Checkpoint参数
# Checkpoint间隔,根据业务容忍度设置
execution.checkpointing.interval: 30s
# Checkpoint超时时间,至少间隔的10倍
execution.checkpointing.timeout: 600s
# 最大同时进行中的Checkpoint数
execution.checkpointing.max-concurrent-checkpoints: 1
# 两次Checkpoint之间的最小间隔(避免Checkpoint太密集)
execution.checkpointing.min-pause-between-checkpoints: 30s
# 不对齐Checkpoint(Flink 1.11+,反压严重时用)
execution.checkpointing.unaligned.enabled: true
踩坑点:反压严重时对齐Checkpoint的Barrier很难往下游传,导致Checkpoint一直超时。这时候开
unaligned Checkpoint,Barrier绕过数据Buffer直接往下走,能大幅提升Checkpoint成功率。但会多占点内存,因为要把在途数据也快照了。
7.3 Operator Chain策略
排查反压时先禁掉Chain,看清每个算子的状态:
// 排查阶段禁用Chain
env.disableOperatorChaining();
定位完根因后,恢复Chain提升性能:
// 默认开启Chain(Flink默认就是开的)
// 只在需要断开的地方断开
stream.map(...)
.keyBy(...)
.window(...)
.aggregate(...) // 这里会自动和后面的断开
.map(...) // 需要单独看性能,前面断开
.startNewChain() // 强制从这里开始新Chain
.addSink(...);
八、Step 5:扩容(最后手段)
以上四步全做完了还撑不住,才是真需要加资源。
8.1 增加TaskManager资源
# flink-conf.yaml
jobmanager.memory.process.size: 2048m
taskmanager.memory.process.size: 8192m # 从4G加到8G
taskmanager.numberOfTaskSlots: 4 # 每个TM的Slot数
踩坑点:加内存前先确认是内存不够。很多情况下代码优化完就能省一半内存,盲目加内存GC时间反而更长。先治代码,再加资源。
8.2 提升Source并行度
env.addSource(kafkaSource).setParallelism(8) // 从4提升到8
提升前确认:
- Kafka Partition数 >= 新的并行度
- 下游算子并行度也能匹配(不然Source快了下游还是堵)
九、完整实战案例
背景
电商实时GMV计算作业:Kafka → Flink聚合 → ClickHouse。上线后反压严重,延迟飙到15分钟。
排查过程
Step 1:定位
Web UI看到WindowAggregate黑色,ClickHouseSink红色。Sink是瓶颈。
Step 2:看SubTask
Sink 4个SubTask中,SubTask 0收了800万条,其他SubTask各收了不到10万条。数据倾斜。
Step 3:分析
keyBy用的是品类ID。“食品生鲜”(cat_1001)占全站60%订单量,全到一个SubTask。
解决过程(按优先级顺序)
① 下游处理侧优化
采用两阶段聚合改造代码,重新上线:
// 加随机前缀打散 → 局部聚合 → 去前缀 → 全局聚合
orders
.map(order -> addRandomPrefix(order, 10)) // 10个前缀打散
.keyBy(Order::getPrefixedCategoryId)
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.aggregate(new LocalSumAgg())
.map(order -> removePrefix(order))
.keyBy(Order::getCategoryId)
.window(TumblingProcessingTimeWindows.of(Time.seconds(10)))
.aggregate(new GlobalSumAgg())
.addSink(clickHouseSink);
改造后SubTask数据量均衡,反压从HIGH降到OK。
② Source端加保险
配置Kafka Consumer参数:
props.setProperty("max.poll.records", "300"); // 从500降到300,留余量
props.setProperty("max.partition.fetch.bytes", "524288"); // 512KB
③ Checkpoint调优
execution.checkpointing.interval: 30s
execution.checkpointing.timeout: 600s
execution.checkpointing.unaligned.enabled: true # 开非对齐Checkpoint
最终结果
| 指标 | 优化前 | 优化后 |
|---|---|---|
| 端到端延迟 | 15分钟 | 3秒 |
| 反压状态 | HIGH持续 | OK |
| Checkpoint成功率 | 30% | 99.5% |
| ClickHouse写入QPS | 不稳定 | 稳定5万/秒 |
十、生产踩坑汇总
| 坑 | 表现 | 正确做法 |
|---|---|---|
| 上来就加机器 | 加了资源还是反压 | 先定位根因,治代码最后才扩容 |
| 只看下游不管上游 | 治好后一重启又崩 | Source端也要限速保护 |
| max.poll.records设太小 | Consumer CPU飙高 | 100~200是底线,别低于100 |
| Watermark设太大 | State膨胀,Checkpoint超时 | 5~30秒够用就好 |
| 忽略短暂反压 | Checkpoint/重启时正常反压 | 持续5分钟以上才需要介入 |
| Operator Chain干扰排查 | UI上看不清哪个算子慢 | 排查时disableOperatorChaining() |
| Kafka Partition < 并行度 | 部分SubTask空转 | Partition数 = 并行度的整数倍 |
| RocksDB托管内存配太大 | 容器OOM | 0.3~0.5根据实际调 |
| 写自定义限速Source | Bug多、难维护 | 优先用Kafka参数控制 |
| Checkpoint间隔设太短 | Checkpoint太密集反压更严重 | 至少10s,一般30s |
十一、一张图记住完整流程
下游决定能处理多快,上游决定能进来多快。先治下游、再控上游、最后扩容。按这个顺序来,反压问题基本都能解决。觉得有用就点个赞吧~
更多推荐
所有评论(0)