
一、Kafka 环境搭建单机版1. 下载 Kafka前往 Apache Kafka 官网 下载最新稳定版例如 3.6.0。解压后目录结构如下kafka_2.13-3.6.0/ ├── bin/ # 启动脚本 ├── config/ # 配置文件 ├── libs/ # 依赖库 └── ...2. 启动 Kafka早期 Kafka 依赖 ZooKeeper新版本推荐使用KRaft模式无需 ZooKeeper。以下分别介绍两种方式。方式一KRaft 模式Kafka 3.3 推荐# 1. 生成集群 ID KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) # 2. 格式化存储目录使用默认配置 bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 3. 启动 Kafka bin/kafka-server-start.sh config/kraft/server.properties方式二ZooKeeper 模式传统# 1. 启动 ZooKeeper bin/zookeeper-server-start.sh config/zookeeper.properties # 2. 启动 Kafka bin/kafka-server-start.sh config/server.properties默认端口Kafka:9092ZooKeeper:21813. 创建 Topicbin/kafka-topics.sh --create --topic my-topic --bootstrap-server localhost:9092 --partitions 3 --replication-factor 1查看 Topic 列表bin/kafka-topics.sh --list --bootstrap-server localhost:90924. 命令行测试发送消息bin/kafka-console-producer.sh --topic my-topic --bootstrap-server localhost:9092 hello kafka消费消息bin/kafka-console-consumer.sh --topic my-topic --from-beginning --bootstrap-server localhost:9092二、Java 原生客户端使用在 Maven 项目中添加依赖pom.xmldependency groupIdorg.apache.kafka/groupId artifactIdkafka-clients/artifactId version3.6.0/version /dependency1. 生产者Producer基本配置项配置项说明bootstrap.serversKafka 集群地址多个用逗号分隔key.serializer键的序列化器如StringSerializervalue.serializer值的序列化器acks确认机制0不等待确认1仅 leader 确认all或-1所有副本确认retries发送失败重试次数batch.size批量发送大小字节linger.ms等待更多消息加入批次的时间buffer.memory生产者缓冲区大小compression.type压缩类型none、gzip、snappy、lz4、zstd生产者代码示例import org.apache.kafka.clients.producer.*; import org.apache.kafka.common.serialization.StringSerializer; import java.util.Properties; import java.util.concurrent.Future; public class MyProducer { public static void main(String[] args) { Properties props new Properties(); props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); props.put(ProducerConfig.ACKS_CONFIG, all); props.put(ProducerConfig.RETRIES_CONFIG, 3); props.put(ProducerConfig.BATCH_SIZE_CONFIG, 16384); props.put(ProducerConfig.LINGER_MS_CONFIG, 1); props.put(ProducerConfig.BUFFER_MEMORY_CONFIG, 33554432); props.put(ProducerConfig.COMPRESSION_TYPE_CONFIG, snappy); // 创建生产者 ProducerString, String producer new KafkaProducer(props); // 1. 发送消息异步不关心结果 producer.send(new ProducerRecord(my-topic, key1, value1)); // 2. 发送消息并获取 Future可阻塞等待结果 FutureRecordMetadata future producer.send( new ProducerRecord(my-topic, key2, value2) ); try { RecordMetadata metadata future.get(); System.out.println(发送成功offset metadata.offset() , partition metadata.partition()); } catch (Exception e) { e.printStackTrace(); } // 3. 发送消息并带回调异步 producer.send(new ProducerRecord(my-topic, key3, value3), new Callback() { Override public void onCompletion(RecordMetadata metadata, Exception exception) { if (exception null) { System.out.println(发送成功: metadata.offset()); } else { exception.printStackTrace(); } } }); // 关闭生产者会等待所有缓冲消息发送完成 producer.close(); } }2. 消费者Consumer基本配置项配置项说明bootstrap.serversKafka 集群地址group.id消费者组 ID相同组内的消费者共同消费key.deserializer键的反序列化器value.deserializer值的反序列化器enable.auto.commit是否自动提交 offsetauto.offset.reset初始消费位置earliest从头开始latest从最新开始max.poll.records一次 poll 最多返回的记录数session.timeout.ms会话超时时间消费者代码示例import org.apache.kafka.clients.consumer.*; import org.apache.kafka.common.serialization.StringDeserializer; import java.time.Duration; import java.util.Collections; import java.util.Properties; public class MyConsumer { public static void main(String[] args) { Properties props new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:9092); props.put(ConsumerConfig.GROUP_ID_CONFIG, my-group); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, true); // 自动提交 props.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, 1000); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest); // 从头消费 ConsumerString, String consumer new KafkaConsumer(props); consumer.subscribe(Collections.singletonList(my-topic)); try { while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, String record : records) { System.out.printf(offset%d, key%s, value%s%n, record.offset(), record.key(), record.value()); } } } finally { consumer.close(); } } }手动提交 offset推荐生产环境使用props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false); // 关闭自动提交 while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(100)); for (ConsumerRecordString, String record : records) { // 处理消息 System.out.println(record.value()); } // 处理完后手动提交 offset同步提交 consumer.commitSync(); // 或者异步提交 // consumer.commitAsync(); }三、Spring Boot 集成 KafkaSpring Kafka 对 Kafka 进行了封装大大简化了开发。下面演示完整流程。1. 添加依赖在pom.xml中添加dependency groupIdorg.springframework.kafka/groupId artifactIdspring-kafka/artifactId /dependencySpring Boot 会自动管理版本通常不需要显式指定版本号2. 配置文件application.ymlspring: kafka: bootstrap-servers: localhost:9092 producer: key-serializer: org.apache.kafka.common.serialization.StringSerializer value-serializer: org.apache.kafka.common.serialization.StringSerializer acks: all retries: 3 consumer: group-id: my-group key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer auto-offset-reset: earliest enable-auto-commit: false listener: ack-mode: manual # 手动提交 offset3. 生产者服务import org.springframework.kafka.core.KafkaTemplate; import org.springframework.stereotype.Service; Service public class KafkaProducerService { private final KafkaTemplateString, String kafkaTemplate; public KafkaProducerService(KafkaTemplateString, String kafkaTemplate) { this.kafkaTemplate kafkaTemplate; } public void sendMessage(String topic, String message) { kafkaTemplate.send(topic, message); } // 带 key 和回调 public void sendMessageWithCallback(String topic, String key, String message) { kafkaTemplate.send(topic, key, message).addCallback( result - System.out.println(发送成功: result.getRecordMetadata().offset()), ex - System.err.println(发送失败: ex.getMessage()) ); } }4. 消费者监听器import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; Component public class KafkaConsumer { // 自动确认配置中 enable-auto-commit: true 时 KafkaListener(topics my-topic, groupId my-group) public void listen(String message) { System.out.println(收到消息: message); } // 手动提交 offset配置中 enable-auto-commit: false 时 KafkaListener(topics my-topic, groupId my-group) public void listenManual(String message, Acknowledgment ack) { System.out.println(收到消息: message); // 处理完成后提交 offset ack.acknowledge(); } }5. 发送消息测试在 Controller 或测试类中注入KafkaProducerService调用即可。四、进阶特性与最佳实践1. 自定义序列化器如果消息是 Java 对象可以使用 JSON 或 Avro 序列化。常用的是 Spring Kafka 提供的JsonSerializer/JsonDeserializer或者使用StringSerializer配合 Jackson 手动转换。2. 分区策略默认分区器如果指定了 key则根据 key 的 hash 选择分区如果 key 为 null则使用轮询round-robin。自定义分区器实现Partitioner接口并在配置中指定partitioner.class。3. 消息顺序性Kafka 只能保证同一个分区内消息有序。因此对于需要严格顺序的业务可以将需要有序的消息发送到同一个分区例如使用相同的 key。4. 幂等性Kafka 生产者默认开启幂等enable.idempotencetrue可以避免因网络重试导致的重复消息。对于消费者端需要业务逻辑本身具备幂等性。5. 事务Kafka 支持事务可以保证“发送多条消息要么全部成功要么全部失败”。配置props.put(ProducerConfig.TRANSACTIONAL_ID_CONFIG, my-transactional-id); producer.initTransactions(); producer.beginTransaction(); producer.send(record1); producer.send(record2); producer.commitTransaction();6. 错误处理在消费者中如果处理消息失败可以抛出异常Spring Kafka 会进行重试。也可以配置ErrorHandler或SeekToCurrentErrorHandler来控制失败后的行为例如重试一定次数后跳过或发送到死信队列。五、常见问题解答Q1: 如何保证消息不丢失生产者设置acksallretries0开启幂等。Broker设置min.insync.replicas至少为 2保证至少有一个副本同步。消费者禁用自动提交处理完消息后手动提交 offset。Q2: 如何保证消息不重复Kafka 本身无法绝对避免重复需要消费者端做幂等如数据库唯一约束、Redis 去重等。可以使用事务或幂等生产者减少重复发送。Q3: 消费进度如何回溯使用kafka-consumer-groups.sh工具重置 offsetbin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --reset-offsets --to-earliest --topic my-topic --executeQ4: 如何监控 Kafka使用 Kafka 自带的 JMX 指标 Prometheus Grafana。使用 Kafka Manager、Kafka Eagle、Burrow 等第三方工具。六、总结使用 Kafka 的基本步骤启动 Kafka 集群单机或集群。创建 Topic。编写生产者配置序列化器、acks 等发送消息。编写消费者配置反序列化器、group.id订阅 Topic 并处理消息。生产环境考虑消息可靠性不丢失、不重复、顺序性、分区策略、监控等。