简介这是一份面向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()返回的是一个FutureRecordMetadata。如果你调.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); KafkaProducerString, String producer new KafkaProducer(props);batch.size设 32KB 是个折中太小攒不起来太大延迟高。linger.ms10意味着即使批次没满最多等 10ms 也会发出去这是延迟和吞吐的平衡点。enable.idempotencetrue配合acksall和retries能保证单分区内不丢不重——这是 Web 服务器场景的底线因为用户请求产生的消息丢了就是业务事故。2.2 Web 层怎么封装生产者才不拖垮请求线程直接在 Servlet 里new KafkaProducer是灾难每次请求建一个生产者连接和元数据开销直接压垮 broker。正确做法是单例生产者 应用级生命周期管理。public class KafkaProducerHolder { private static volatile KafkaProducerString, String producer; public static KafkaProducerString, 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) { ProducerRecordString, 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.id1 listenersPLAINTEXT://0.0.0.0:9092 advertised.listenersPLAINTEXT://kafka1:9092 log.dirs/data/kafka-logs num.partitions6 default.replication.factor3 min.insync.replicas2 log.retention.hours72advertised.listeners是最容易配错的。如果 Web 服务器和 Kafka 不在同一台机器客户端拿到的是advertised.listeners里的地址去连。配成localhost就会出现「能连上但发不出去」的玄学问题。min.insync.replicas2配合生产者的acksall意味着至少两个副本写入成功才确认一个 broker 挂了不影响写入。num.partitions6是给 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)) { SetString 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.enabletrue生产者往不存在的 topic 发消息会自动创建但分区数和副本因子用的是 broker 默认值往往不是你想要的。生产环境建议关掉自动创建用 AdminClient 显式管理。4. 消费端多线程如何保证消息顺序性分区策略与位移提交4.1 顺序性的边界分区内有序跨分区无序Kafka 只保证单分区内消息有序。Web 服务器场景下如果同一用户的请求要按顺序处理就必须让同一用户的消息进同一个分区。做法是在生产者端指定 key// 用 userId 作为 keyKafka 默认按 key hash 分区 ProducerRecordString, 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 (KafkaConsumerString, String consumer new KafkaConsumer(props)) { TopicPartition tp new TopicPartition(topic, partition); consumer.assign(Collections.singletonList(tp)); consumer.seekToBeginning(Collections.singletonList(tp)); while (running) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(500)); for (ConsumerRecordString, String record : records) { process(record); // 单线程顺序处理 } consumer.commitSync(); // 处理完再提交 } } } }这里用assign而不是subscribe手动指定分区每个线程负责一个分区天然有序。enable.auto.commitfalse加上处理完再commitSync保证至少一次语义。max.poll.records100控制单次拉取量避免处理太久触发 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太大导致批次迟迟不满。也有可能是acksall加上min.insync.replicas过高等待副本同步。解决Web 场景linger.ms控制在 5~20msbatch.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.replicas2acksall挂一个后只剩两个副本如果其中一个还没同步完就不满足最小同步副本数。解决这是设计行为保证不丢消息。如果业务能容忍少量丢失可以降到min.insync.replicas1但要在业务层做补偿。或者加快副本同步检查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.levelread_committed只读已提交的事务消息。这样生产端多条消息要么全成功要么全失败配合消费端的位移提交能做到端到端的精确一次。但事务有代价transactional.id要唯一且稳定生产者重启后要用同一个 id 才能恢复未完成事务事务超时transaction.timeout.ms默认 60 秒处理太慢会超时回滚。我一般只在跨 topic 写且要求原子性的场景用普通 Web 埋点用幂等生产者加业务去重就够了。验证方法用kafka-console-consumer --isolation-level read_committed消费对比read_uncommitted看未提交事务的消息是否被过滤。再模拟生产者中途 kill重启后看事务是否被正确 abort。我自己的习惯是任何接 Kafka 的 Web 项目先把acksall、enable.idempotencetrue、min.insync.replicas2这三条写进配置模板再根据业务调linger.ms和batch.size。顺序性和精确一次不是靠背八股文是靠这些参数一个个对齐出来的。希望帮到你。本文还有配套的精品资源点击获取