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

    • 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 概述

定义与作用

Kafka Producer 是向 Kafka Topic 发送消息的客户端程序。它的职责不限于"把数据发出去"——它承担了分区路由、序列化、压缩、批处理攒批、重试容错等一系列复杂职责,而将这些复杂性对应用开发者透明化。

Producer 在 Kafka 架构中的位置:数据入口,是所有数据管道的第一步。

核心原理

Producer 的内部架构

双线程设计:

组件线程职责
KafkaProducer(主线程)应用程序线程拦截 → 序列化 → 分区 → 放入 RecordAccumulator
Sender(后台线程)独立 I/O 线程从 RecordAccumulator 取批次 → 发送到 Broker → 处理 ACK

ProducerRecord 消息结构

ProducerRecord<String, String> record = new ProducerRecord<>(
    "orders",           // topic: 目标 Topic
    0,                  // partition: 指定分区(可选)
    1718208000000L,     // timestamp: 时间戳(可选)
    "order-1001",       // key: 消息 Key(可选,用于分区路由)
    "{\"amount\":99}"   // value: 消息体
);
字段必选用途
topic是目标 Topic
value是消息负载(payload)
key否用于分区路由,相同 Key 的消息进入同一分区
partition否显式指定分区,优先级高于 Key 路由
timestamp否消息时间戳,不指定则使用 System.currentTimeMillis()
headers否键值对形式的自定义元数据(可用于链路追踪等)

完整示例

示例一:最小可用的 Producer

场景:向 orders Topic 发送一条订单消息。

操作前环境:Kafka 已启动,Topic orders 已创建(1 分区)。

import org.apache.kafka.clients.producer.*;
import java.util.Properties;

public class SimpleProducer {
    public static void main(String[] args) throws Exception {
        // 1. 配置 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");

        // 2. 创建 Producer 实例
        try (KafkaProducer<String, String> producer = new KafkaProducer<>(props)) {

            // 3. 构造消息
            ProducerRecord<String, String> record = new ProducerRecord<>(
                "orders", "order-1001", "{\"amount\": 99.90, \"item\": \"book\"}"
            );

            // 4. 发送(异步)
            producer.send(record, (metadata, exception) -> {
                if (exception == null) {
                    System.out.printf("发送成功 → topic=%s, partition=%d, offset=%d%n",
                        metadata.topic(), metadata.partition(), metadata.offset());
                } else {
                    exception.printStackTrace();
                }
            });
        }
        // 5. try-with-resources 自动调用 close(),等待缓冲区消息全部发送
    }
}

执行结果:

发送成功 → topic=orders, partition=0, offset=0

操作后状态:

指标发送前发送后
orders Partition 0 的 LEO01
消息内容—{"amount": 99.90, "item": "book"}

示例二:同步发送 vs 异步发送

场景:对比 get() 同步阻塞和 Callback 异步两种方式。

// ===== 方式一:同步发送(阻塞等待结果) =====
try {
    RecordMetadata metadata = producer.send(record).get();
    System.out.println("同步发送成功, offset=" + metadata.offset());
} catch (Exception e) {
    System.err.println("同步发送失败: " + e.getMessage());
}

// ===== 方式二:异步发送(Callback 非阻塞) =====
producer.send(record, new Callback() {
    @Override
    public void onCompletion(RecordMetadata metadata, Exception e) {
        if (e != null) {
            System.err.println("异步发送失败: " + e.getMessage());
        } else {
            System.out.println("异步发送成功, offset=" + metadata.offset());
        }
    }
});

操作前后对比:

方式吞吐量延迟适用场景
同步 get()低(~100 QPS)每条 ~1-5ms强顺序要求、少量消息
异步 Callback高(~100K QPS)应用层几乎零等待高吞吐场景(默认选择)

注意:异步方式下,producer.close() 会阻塞等待所有未完成的发送确认,所以上例中 try-with-resources 保证了消息不会因进程退出而丢失。

易错场景

易错 1:忘记调用 close()

现象:main 方法结束后 JVM 退出,发现 Producer 已发送的消息在 Kafka 中丢失。

原因:消息在 RecordAccumulator 的缓冲区中尚未被 Sender 线程实际发送到 Broker。

解决:始终使用 try-with-resources 或 finally { producer.close(); },close() 会等待缓冲区清空。

易错 2:每条消息创建新的 Producer 实例

// 错误:每条消息创建 Producer(开销巨大)
for (int i = 0; i < 10000; i++) {
    KafkaProducer<String, String> p = new KafkaProducer<>(props);
    p.send(new ProducerRecord<>("topic", "msg" + i));
    p.close();
}

后果:每次 new KafkaProducer 都会创建 Sender 线程、建立 TCP 连接、获取元数据,性能极差。

正确做法:Producer 是线程安全的,一个进程通常只需一个 Producer 实例。

面试高频考点

Q:Kafka Producer 为什么设计为异步发送?

A:异步发送 + 批处理是 Kafka 高吞吐的核心设计之一:

  1. 解耦应用线程和 I/O:应用线程放入缓冲区即返回,Sender 线程独立管理网络 I/O
  2. 批处理攒批:RecordAccumulator 按分区聚合消息,将小消息合并为大网络包,减少网络往返
  3. 压缩:批次级别的压缩远优于逐条压缩
  4. 容错:Sender 线程可独立处理重试,不阻塞业务线程
上一页
章节导读
下一页
Producer 发送机制