简介:这是一份面向Java后端开发者与大数据入门者的Kafka实践示例包,聚焦分布式流处理平台与Web服务器场景的集成应用。内容围绕Kafka核心概念展开,涵盖主题、分区、副本、生产者、消费者及消费者组等基础机制,并延伸至日志聚合、API消息队列与事件驱动架构等典型用法,帮助读者理解如何用Java客户端完成消息发布与订阅。压缩包共30个文件,以19个jar依赖库、4个java源码、4个class编译文件为主,另含classpath、prefs与project等工程配置,整体约6.95MB,可直接导入IDE运行调试。示例代码演示了Bootstrap Servers、序列化器等连接参数配置,以及生产者发送消息、消费者订阅并处理消息的完整流程,同时涉及offset管理与幂等性设计等故障恢复思路。目前已有108人学习,适合希望快速上手Kafka与Web服务集成、优化并发处理能力的开发者参考。
1. 从 kafka-example.rar 说起:一个 Java Web 服务器项目到底该长什么样
拿到kafka-example.rar_Web服务器_Java_这个标题,很多人第一反应是「又一个 Kafka 教程 Demo」。但把关键词拆开看——Kafka、Web 服务器、Java——它其实指向一个非常具体的落地场景:用 Java 写一个 Web 服务,把 HTTP 请求产生的消息可靠地投递到 Kafka,再由消费端处理。这不是单纯的 Kafka 客户端示例,也不是单纯的 Servlet 项目,而是两者的接缝处,而接缝处恰恰是最容易翻车的地方。
我见过太多团队在这个接缝上踩坑:Web 层用同步发送,QPS 一上来线程池直接打满;消费端多线程拉高吞吐,结果消息顺序全乱;Kafka 集群三节点刚装好,InvalidReceiveException就糊在日志里。这篇笔记就按「Web 服务器 + Java + Kafka」这条线,把选型理由、最小可跑代码、参数怎么调、坑在哪,一层层讲清楚。适合正在做消息队列选型、或者已经选了 Kafka 但还没跑通生产级链路的 Java 后端。
2. 为什么 Web 服务器接 Kafka 要用异步生产者:选型与线程模型
2.1 同步发送和异步发送的真实差距
Java 里往 Kafka 发消息,KafkaProducer.send()返回的是一个Future<RecordMetadata>。如果你调.get(),那就是同步发送——每条消息都要等 broker 确认。在 Web 服务器场景下,这意味着一个 HTTP 请求线程被阻塞在网络上,Tomcat 默认 200 个线程,QPS 稍微上来就排队。
异步发送则是传一个Callback,请求线程把消息丢进生产者的缓冲区就返回。Kafka 生产者内部有一个Sender线程,负责把缓冲区里的批次发出去。这个设计的关键在于批量:batch.size和linger.ms决定了攒批的大小和等待时间。
我一般会这样配:
Properties props = new Properties(); props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092"); props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer"); props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer"); // 关键参数:攒批大小 32KB,等待 10ms props.put("batch.size", 32768); props.put("linger.ms", 10); // 重试与幂等 props.put("retries", Integer.MAX_VALUE); props.put("enable.idempotence", true); // 缓冲区上限,防止内存打爆 props.put("buffer.memory", 67108864); // 确认级别:all 保证不丢 props.put("acks", "all"); KafkaProducer<String, String> producer = new KafkaProducer<>(props);batch.size设 32KB 是个折中:太小攒不起来,太大延迟高。linger.ms=10意味着即使批次没满,最多等 10ms 也会发出去,这是延迟和吞吐的平衡点。enable.idempotence=true配合acks=all和retries,能保证单分区内不丢不重——这是 Web 服务器场景的底线,因为用户请求产生的消息丢了就是业务事故。
2.2 Web 层怎么封装生产者才不拖垮请求线程
直接在 Servlet 里new KafkaProducer是灾难,每次请求建一个生产者,连接和元数据开销直接压垮 broker。正确做法是单例生产者 + 应用级生命周期管理。
public class KafkaProducerHolder { private static volatile KafkaProducer<String, String> producer; public static KafkaProducer<String, String> get() { if (producer == null) { synchronized (KafkaProducerHolder.class) { if (producer == null) { producer = new KafkaProducer<>(buildProps()); } } } return producer; } public static void sendAsync(String topic, String key, String value) { ProducerRecord<String, String> record = new ProducerRecord<>(topic, key, value); get().send(record, (metadata, exception) -> { if (exception != null) { // 这里必须打日志,不能吞 log.error("send failed, topic={}, key={}", topic, key, exception); } }); } public static void close() { if (producer != null) { producer.flush(); producer.close(Duration.ofSeconds(10)); } } }KafkaProducer是线程安全的,多个请求线程共用一个实例完全没问题。send的回调里必须处理异常,否则消息丢了都不知道。close要注册到 ServletContextListener 的contextDestroyed里,先flush再close,保证缓冲区里的消息发完。
提示:
buffer.memory满了之后send()会阻塞最多max.block.ms(默认 60 秒)。Web 场景下这个值建议调到 5000ms 以内,宁可快速失败也不要拖死请求线程。
3. 三节点 Kafka 集群在 Web 项目里的最小可用配置
3.1 broker 端必须改的四个参数
Kafka 集群安装教程网上很多,但 Web 服务器项目接进来时,有几个参数和默认值不一样。假设三台机器kafka1/2/3,每台server.properties里:
broker.id=1 listeners=PLAINTEXT://0.0.0.0:9092 advertised.listeners=PLAINTEXT://kafka1:9092 log.dirs=/data/kafka-logs num.partitions=6 default.replication.factor=3 min.insync.replicas=2 log.retention.hours=72advertised.listeners是最容易配错的。如果 Web 服务器和 Kafka 不在同一台机器,客户端拿到的是advertised.listeners里的地址去连。配成localhost就会出现「能连上但发不出去」的玄学问题。min.insync.replicas=2配合生产者的acks=all,意味着至少两个副本写入成功才确认,一个 broker 挂了不影响写入。
num.partitions=6是给 Web 场景的:分区数决定了消费端并行度上限。三节点集群,每个 topic 6 个分区,副本因子 3,刚好每个 broker 承载 6 个副本,分布均匀。
3.2 用 AdminClient 在应用启动时自动建 topic
手动敲命令建 topic 容易漏,我一般会在 Web 应用启动时用AdminClient检查并创建:
public static void ensureTopic(String topic, int partitions, short replication) { Properties props = new Properties(); props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092"); try (AdminClient admin = AdminClient.create(props)) { Set<String> existing = admin.listTopics().names().get(); if (!existing.contains(topic)) { NewTopic newTopic = new NewTopic(topic, partitions, replication); admin.createTopics(Collections.singleton(newTopic)).all().get(); log.info("topic {} created", topic); } } catch (Exception e) { log.error("ensure topic failed", e); } }AdminClient是 Kafka 2.0 之后的标准管理接口,比调脚本可靠。createTopics返回的CreateTopicsResult调.all().get()会阻塞直到创建完成或超时。注意replication不能超过 broker 数量,三节点集群最大就是 3。
注意:如果 broker 端配了
auto.create.topics.enable=true,生产者往不存在的 topic 发消息会自动创建,但分区数和副本因子用的是 broker 默认值,往往不是你想要的。生产环境建议关掉自动创建,用 AdminClient 显式管理。
4. 消费端多线程如何保证消息顺序性:分区策略与位移提交
4.1 顺序性的边界:分区内有序,跨分区无序
Kafka 只保证单分区内消息有序。Web 服务器场景下,如果同一用户的请求要按顺序处理,就必须让同一用户的消息进同一个分区。做法是在生产者端指定 key:
// 用 userId 作为 key,Kafka 默认按 key hash 分区 ProducerRecord<String, String> record = new ProducerRecord<>("order-events", userId, payload);Kafka 默认的DefaultPartitioner对 key 做 murmur2 hash 再对分区数取模,同一个 key 永远进同一个分区。这样消费端只要单线程消费一个分区,顺序就有保证。
4.2 多线程消费的正确姿势:分区级并行
消费端要提吞吐,又不能乱序,常见做法是每个分区一个消费线程,而不是一个消费者里开线程池处理消息。
public class PartitionConsumer implements Runnable { private final String topic; private final int partition; public void run() { Properties props = new Properties(); props.put("bootstrap.servers", "kafka1:9092,kafka2:9092,kafka3:9092"); props.put("group.id", "web-consumer-group"); props.put("key.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("value.deserializer", "org.apache.kafka.common.serialization.StringDeserializer"); props.put("enable.auto.commit", "false"); props.put("max.poll.records", "100"); try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) { TopicPartition tp = new TopicPartition(topic, partition); consumer.assign(Collections.singletonList(tp)); consumer.seekToBeginning(Collections.singletonList(tp)); while (running) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500)); for (ConsumerRecord<String, String> record : records) { process(record); // 单线程顺序处理 } consumer.commitSync(); // 处理完再提交 } } } }这里用assign而不是subscribe,手动指定分区,每个线程负责一个分区,天然有序。enable.auto.commit=false加上处理完再commitSync,保证至少一次语义。max.poll.records=100控制单次拉取量,避免处理太久触发 rebalance。
如果分区数多于线程数,可以一个线程负责多个分区,但每个分区仍然是顺序处理。关键是不要把一个分区的消息丢给线程池并发处理,那样顺序必乱。
提示:
commitSync会阻塞,如果处理逻辑慢,可以改成commitAsync加定期同步提交,但要注意异步提交失败时的重试逻辑,否则可能重复消费。
5. 避坑:Web 服务器接 Kafka 最常见的五类翻车
5.1 InvalidReceiveException:连接上了但握手就断
现象:日志里刷org.apache.kafka.common.network.InvalidReceiveException: Invalid receive (size = 369296129)。
原因:客户端连到了非 Kafka 端口。最常见的是 Web 服务器和 Kafka 混部,客户端配的bootstrap.servers指向了 Web 服务的端口,或者advertised.listeners配错导致客户端拿到错误地址。
解决:确认bootstrap.servers和advertised.listeners端口一致,用telnet kafka1 9092确认端口是 Kafka 在监听。如果混部,给 Kafka 单独绑一个网卡或端口。
5.2 消息延迟高:不是 Kafka 慢,是攒批参数不对
现象:监控显示生产端 P99 延迟几百毫秒甚至秒级。
原因:linger.ms设太大,或者batch.size太大导致批次迟迟不满。也有可能是acks=all加上min.insync.replicas过高,等待副本同步。
解决:Web 场景linger.ms控制在 5~20ms,batch.size32KB~64KB。如果延迟还是高,检查 broker 的num.replica.fetchers和磁盘 IO。用kafka-producer-perf-test压一下,看是网络还是磁盘瓶颈。
5.3 消费端 rebalance 导致重复消费
现象:消费组频繁 rebalance,消息重复处理。
原因:max.poll.interval.ms默认 5 分钟,如果单次poll拉取的消息处理超过这个时间,消费者被踢出组,触发 rebalance,位移没提交,新消费者从上次提交位置重新消费。
解决:调小max.poll.records,或者把处理逻辑异步化但保证提交前处理完。也可以调大max.poll.interval.ms,但根本办法是控制单批处理时间。
5.4 生产者缓冲区满导致请求线程阻塞
现象:Web 接口偶尔超时,日志显示send卡住。
原因:buffer.memory满了,send阻塞等待max.block.ms。
解决:调大buffer.memory,或者调小max.block.ms让快速失败。更根本的是检查消费端是不是挂了导致生产端积压。加监控看record-queue-time-avg。
5.5 三节点集群挂一个后写入失败
现象:一个 broker 宕机,生产端报NotEnoughReplicasException。
原因:min.insync.replicas=2,acks=all,挂一个后只剩两个副本,如果其中一个还没同步完,就不满足最小同步副本数。
解决:这是设计行为,保证不丢消息。如果业务能容忍少量丢失,可以降到min.insync.replicas=1,但要在业务层做补偿。或者加快副本同步,检查replica.lag.time.max.ms。
6. 进阶:用幂等生产者加事务把「至少一次」做成「精确一次」
前面消费端用commitSync做到的是至少一次,重复消费要靠业务去重。如果业务对重复零容忍,可以上 Kafka 事务。
生产者端开幂等和事务:
props.put("enable.idempotence", true); props.put("transactional.id", "web-tx-1"); // 每个生产者实例唯一 producer.initTransactions(); try { producer.beginTransaction(); producer.send(new ProducerRecord<>("topic-a", key, value)); producer.send(new ProducerRecord<>("topic-b", key, value)); producer.commitTransaction(); } catch (Exception e) { producer.abortTransaction(); }消费端配isolation.level=read_committed,只读已提交的事务消息。这样生产端多条消息要么全成功要么全失败,配合消费端的位移提交,能做到端到端的精确一次。
但事务有代价:transactional.id要唯一且稳定,生产者重启后要用同一个 id 才能恢复未完成事务;事务超时transaction.timeout.ms默认 60 秒,处理太慢会超时回滚。我一般只在跨 topic 写且要求原子性的场景用,普通 Web 埋点用幂等生产者加业务去重就够了。
验证方法:用kafka-console-consumer --isolation-level read_committed消费,对比read_uncommitted,看未提交事务的消息是否被过滤。再模拟生产者中途 kill,重启后看事务是否被正确 abort。
我自己的习惯是,任何接 Kafka 的 Web 项目,先把acks=all、enable.idempotence=true、min.insync.replicas=2这三条写进配置模板,再根据业务调linger.ms和batch.size。顺序性和精确一次不是靠背八股文,是靠这些参数一个个对齐出来的。希望帮到你。
本文还有配套的精品资源,点击获取