生产环境遇到反压,别急着加机器。按下面的顺序来,先治下游、再控上游,最后才是扩容。这篇文章按这个优先级组织,看完直接就能用。


一、先搞清楚:反压是怎么一回事

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:定位瓶颈——四步排查法

发现反压后,按下面四步找到根因:

数据均衡

某个SubTask数据量>其他5倍

CPU正常

UDF方法占很宽

GC正常

Full GC频繁 / GC停顿>200ms

① 看SubTask数据量

② 看火焰图CPU

数据倾斜

两阶段聚合 / rebalance

③ 看GC日志

代码瓶颈

Async I/O / 对象复用

④ 查外部系统

内存问题

调托管内存 / RocksDB配置

ClickHouse批量 / HBase预分区 / MySQL连接池

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 + 旁路缓存
大对象频繁创建 newGC占很宽 对象复用
复杂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

提升前确认:

  1. Kafka Partition数 >= 新的并行度
  2. 下游算子并行度也能匹配(不然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

十一、一张图记住完整流程

发现反压
Web UI黑色/红色

排查四步法

数据倾斜?
两阶段聚合 / rebalance / 热点Key分离

代码瓶颈?
Async I/O / 对象复用 / 正则优化

GC频繁?
内存调优 / RocksDB配置 / 减少State

外部系统?
ClickHouse批量 / HBase预分区 / MySQL连接池

Source端保护

Kafka max.poll.records调小

启动追数时分阶段放开

Source并行度匹配Partition数

作业配置调优

开buffer-debloat自动调节

反压严重时开unaligned Checkpoint

合理设置Checkpoint间隔

扩容
最后手段

加TM资源

提升Source并行度


下游决定能处理多快,上游决定能进来多快。先治下游、再控上游、最后扩容。按这个顺序来,反压问题基本都能解决。觉得有用就点个赞吧~

更多推荐