乐途乐途
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
主页
  • 计算机基础

    • TCP/IP
    • Linux
    • HTTP
  • 数据库

    • SQL
    • MySQL 5.7
  • 编程语言

    • C
    • C++
    • Java SE
    • Python2
    • Python3
  • 数据格式

    • JSON
    • XML
  • 认证与安全

    • JWT
  • 工具

    • Markdown
  • Git

    • GitFlow
  • Quartz

    • Quartz
  • Java

    • Maven 入门
    • Maven 进阶
    • MyBatis
    • Spring
    • Spring MVC
  • Java

    • Spring Boot
    • Spring Cloud
    • Spring Cloud Alibaba
    • Spring Security
    • Spring AI
    • Spring Batch
    • Kafka
    • Java 设计模式
  • 缓存

    • Redis
  • 搜索引擎

    • Elasticsearch
  • 分布式协调

    • ZooKeeper
联系
阿里云
  • 学习路径
  • 第1章 消息队列与 Kafka 概述

    • 章节导读
    • 消息队列基础
    • Kafka 概述
    • Kafka 为什么快
  • 第2章 快速上手:单机环境搭建

    • 章节导读
    • 环境准备与安装
    • 快速启动
    • Topic 管理
  • 第3章 核心概念:主题、分区与日志

    • 章节导读
    • Topic
    • Partition
    • Offset
    • Segment 与存储结构
  • 第4章 生产者详解

    • 章节导读
    • Producer 概述
    • Producer 发送机制
    • Producer 分区策略
    • Producer 幂等与事务
    • Producer 配置
  • 第5章 消费者与消费组

    • 章节导读
    • Consumer 概述
    • Consumer Group 消费组
    • Consumer 分区分配策略
    • Consumer Offset 提交
    • Consumer 多线程
    • Consumer 配置
  • 第6章 Broker 与控制器

    • 章节导读
    • Broker 概述
    • Broker 配置
    • Controller
    • KRaft 共识协议
  • 第7章 副本与数据可靠性

    • 章节导读
    • 副本机制
    • ISR 与副本同步
    • Leader 选举
    • ACK 与一致性保证
    • 高水位与 Leader Epoch
  • 第8章 存储与性能优化

    • 章节导读
    • 存储架构
    • 日志清理与压缩
    • Page Cache 与零拷贝
    • Producer 性能优化
    • Consumer 性能优化
    • Broker 性能优化
  • 第9章 生产环境运维与监控

    • 章节导读
    • Topic 管理
    • Kafka 运维工具
    • 监控
    • 常见故障排查
  • 第10章 Kafka生态与面试考点

    • 章节导读
    • Kafka 生态全景
    • 面试高频 30 题

Producer 发送机制

定义与作用

Producer 发送机制描述了消息从 send() 调用到 Broker 确认写入的完整链路。这条链路决定了消息的可靠性和延迟——理解它,才能正确配置 ACK、重试和批处理参数。

核心原理

完整发送流程

ACK 机制的三级保证

acks 是 Producer 最重要的可靠性参数:

acks 值发送延迟可靠性适用场景
0最低消息可能丢失指标监控(丢几条无所谓)
1中等Leader 崩溃时可能丢失日志收集(可容忍少量丢失)
all / -1最高不丢失(配合 min.insync.replicas)交易、订单等关键数据

重试机制

重要:retries 配合 enable.idempotence=true 使用时,Kafka 会自动将 max.in.flight.requests.per.connection 限制为 5,保证消息顺序。

完整示例

示例一:三种 ACK 策略的吞吐对比

场景:测试 acks 对吞吐量的影响。

操作前环境:Topic ack-test,3 分区,3 副本。

# acks=0
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
  --record-size 100 --throughput -1 \
  --producer-props acks=0 linger.ms=5 bootstrap.servers=localhost:9092
# 100000 records sent, 142857 records/sec (13.62 MB/sec)

# acks=1
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
  --record-size 100 --throughput -1 \
  --producer-props acks=1 linger.ms=5
# 100000 records sent, 89285 records/sec (8.51 MB/sec)

# acks=all
bin/kafka-producer-perf-test.sh --topic ack-test --num-records 100000 \
  --record-size 100 --throughput -1 \
  --producer-props acks=all linger.ms=5
# 100000 records sent, 35714 records/sec (3.41 MB/sec)

操作后对比:

acks吞吐量延迟(平均)可靠性
0142,857/s0.05ms可能丢失
189,285/s1.2msLeader 故障可能丢失
all35,714/s3.8ms不丢(含副本)

示例二:配置重试避免消息丢失

场景:网络抖动导致 Producer 发送失败,需要自动重试。

Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");

// ==== 可靠性优先配置 ====
props.put("acks", "all");                              // 等待所有 ISR
props.put("retries", Integer.MAX_VALUE);               // 无限重试(实际上由 delivery.timeout.ms 控制)
props.put("max.in.flight.requests.per.connection", 5); // 幂等时默认5,防止乱序
props.put("enable.idempotence", true);                 // 开启幂等,防止重复
props.put("delivery.timeout.ms", 120000);              // 2分钟超时

try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {
    ProducerRecord<String, String> record = new ProducerRecord<>("orders", "order-1", "data");
    producer.send(record, (metadata, exception) -> {
        if (exception != null) {
            // delivery.timeout.ms 耗尽后才会进入这里
            System.err.println("最终发送失败: " + exception.getMessage());
        }
    });
}

操作前后对比:

配置无重试(retries=0)有重试(retries=MAX)
网络闪断 100ms消息丢失自动重试成功
Leader 切换消息丢失刷新元数据后重试成功
连续故障 >2min消息丢失delivery.timeout.ms 超时后失败

易错场景

易错 1:acks=all 但 min.insync.replicas=1

场景:Producer 配置 acks=all,但没有设置 Broker 端的 min.insync.replicas。

后果:当 ISR 只有 Leader 一个节点时(Follower 故障),acks=all 退化为 acks=1,不具备真正的可靠性。

正确做法:Topic 级别设置 min.insync.replicas=2,确保至少 2 个副本在 ISR 中才允许写入。

易错 2:max.in.flight.requests.per.connection 与顺序性

场景:设置 max.in.flight.requests.per.connection=5(允许 5 个未确认请求并发),且未开启幂等。

后果:如果 batch1 发送失败被重试,而 batch2-batch5 已成功写入,重试成功后的 batch1 会写在 batch5 之后——消息乱序。

规则:

  • 需要严格有序且未开启幂等 → max.in.flight.requests.per.connection=1
  • 开启幂等 → Kafka 自动限制为 5,且保证顺序

面试高频考点

Q:acks=all 是否绝对保证消息不丢失?

A:不绝对。acks=all 保证的是"消息被 ISR 中所有副本确认后才认为发送成功"。但如果:

  1. ISR 中的所有副本在确认后、返回 Producer 响应前的瞬间全部同时崩溃(极小概率)
  2. min.insync.replicas=1 且 ISR 只有 Leader,退化为 acks=1

真正的"绝对不丢"需要在 Producer(acks=all + retries=MAX + 幂等)、Broker(min.insync.replicas >= 2 + unclean.leader.election.enable=false)和 Consumer(手动提交 Offset、处理完再提交)三个层面同时保证。

上一页
Producer 概述
下一页
Producer 分区策略