简介:这份资源是面向Java后端开发者与大数据入门者的Kafka实践示例包,聚焦Kafka与Web服务器集成的典型场景,帮助读者理解分布式消息队列在实时数据处理中的作用。压缩包共30个文件,以19个jar依赖库、4个java源码、4个class编译文件为主,另含classpath、project等工程配置,整体约6.95MB,可直接导入IDE运行调试。内容围绕主题、分区、副本、生产者、消费者及消费者组等核心概念展开,示例代码演示了通过Java API配置Bootstrap Servers与序列化器、向指定主题发送消息、订阅主题并拉取处理消息的完整流程,同时涉及日志聚合、API异步通信与事件驱动架构等Web服务器结合方式。目前已有108人学习,适合希望从代码层面掌握Kafka基本用法、理解消息队列与Web服务集成思路的开发者参考。
1. 从 kafka-example.rar 拆开看:一个 0.7.1 时代的 Java 客户端全貌
手里这个kafka-example.rar,解压后第一眼看到的不是源码,而是一堆 jar 和目录:kafka-0.7.1.jar、zookeeper-3.3.4.jar、zkclient-0.1.jar、log4j-1.2.15.jar、snappy-java-1.0.4.1.jar,还有src、bin、.classpath、.project、.settings这些 Eclipse 工程文件。这不是一个 Maven 工程,而是一个典型的 Eclipse 时代 Java 项目快照,依赖全靠 lib 目录手工管理。它解决的核心问题很具体:让你在一个没有构建工具、没有容器编排的环境里,用最原始的 Java API 把 Kafka 生产者和消费者跑起来,理解 Web 服务器日志或 API 请求是怎么被塞进消息队列的。适合谁?适合那些拿到一份老项目源码、需要在本地快速验证 Kafka 收发逻辑,或者想搞清楚 Kafka 0.7.x 时代客户端到底长什么样的 Java 开发者。如果你平时用的是 Spring Kafka 或 Kafka 2.x 以上的客户端,这份代码会让你看到很多后来被封装掉的底层细节,比如 ZooKeeper 直连、分区元数据手工拉取、offset 自己存。
2. 环境搭建与依赖梳理:把散落的 jar 变成可运行的 classpath
2.1 为什么这个工程没有 pom.xml 也能跑
kafka-example.rar里的依赖是平铺在 lib 目录下的,没有 Maven 的传递依赖解析。这种做法的好处是版本完全锁定,坏处是你得自己保证所有 jar 都在 classpath 里。从文件列表看,核心依赖分四类:Kafka 本体(kafka-0.7.1.jar)、ZooKeeper 客户端(zookeeper-3.3.4.jar、zkclient-0.1.jar)、日志(log4j-1.2.15.jar)、序列化与压缩(snappy-java-1.0.4.1.jar)。另外还有一堆测试和构建辅助 jar,比如junit-4.1.jar、easymock-3.0.jar、scalatest-1.2.jar、apache-rat-*.jar,这些在跑示例时不是必须的,但如果你要改代码做单元测试,就得留着。常见做法是直接在 Eclipse 里右键工程 → Build Path → Configure Build Path → Add JARs,把 lib 下所有 jar 加进去。但如果你用命令行编译,就得手写 classpath。
2.2 手工编译与运行的最小命令集
假设你把压缩包解压到了D:\kafka-example,源码在src目录下,编译输出到bin。先确认 JDK 版本,Kafka 0.7.1 时代主流是 JDK 6 或 7,用 JDK 8 编译一般也能过,但别用 JDK 11 以上,sun.misc相关的类可能找不到。编译命令如下:
# Windows 下用分号分隔 classpath,Linux/macOS 用冒号 javac -cp "lib/*" -d bin src/kafka/example/*.java这里-cp "lib/*"是通配符写法,JDK 6 以上支持。-d bin把编译后的 class 文件按包结构输出到 bin 目录。如果源码里引用了kafka.javaapi.producer.Producer这类类,编译能过就说明 classpath 没问题。运行生产者示例:
java -cp "bin;lib/*" kafka.example.ProducerDemo注意 Windows 下 classpath 分隔符是分号,Linux/macOS 是冒号。kafka.example.ProducerDemo是假设的类名,实际类名以 src 下为准。参数方面,Kafka 0.7.1 的生产者配置里zk.connect是必填的,指向 ZooKeeper 地址,而不是后来 0.8 以后的bootstrap.servers。这是最容易被后来版本惯坏的地方——你拿一份 2.x 的配置去填,连不上是必然的。
2.3 ZooKeeper 在 0.7.1 里的角色:不是可选,是必须
Kafka 0.7.1 把消费者 offset 存在 ZooKeeper 里,broker 的元数据也注册在 ZooKeeper 上。所以你的本地环境必须先起一个 ZooKeeper。压缩包里带了zookeeper-3.3.4.jar,但没有带 ZooKeeper 服务端的启动脚本。常见做法是单独下载一个 ZooKeeper 3.3.x 的发行包,解压后改conf/zoo.cfg,指定dataDir,然后bin/zkServer.sh start(Windows 下是zkServer.cmd)。启动后确认 2181 端口在监听。如果你跳过这一步直接跑生产者,会看到org.apache.zookeeper.KeeperException$ConnectionLossException,这不是 Kafka 的错,是 ZooKeeper 没通。
提示:ZooKeeper 3.3.4 和 JDK 版本有绑定关系,JDK 8 下跑一般没问题,JDK 9 以上会因为模块化限制报错,建议用 JDK 8 做这个实验。
3. 生产者与消费者代码拆解:从配置项到消息流转
3.1 生产者配置里的三个关键参数
Kafka 0.7.1 的ProducerConfig和后来的ProducerConfig完全不是一回事。打开生产者示例代码,你会看到类似这样的配置:
// Kafka 0.7.1 生产者配置示例 Properties props = new Properties(); props.put("zk.connect", "127.0.0.1:2181"); // ZooKeeper 地址,不是 broker 地址 props.put("serializer.class", "kafka.serializer.StringEncoder"); // 消息序列化类 props.put("zk.connectiontimeout.ms", "10000"); // 连接 ZooKeeper 超时 props.put("producer.type", "sync"); // sync 或 async props.put("compression.codec", "0"); // 0 不压缩,1 gzip,2 snappy ProducerConfig config = new ProducerConfig(props); Producer<String, String> producer = new Producer<String, String>(config);zk.connect是核心,它告诉生产者去哪里找 broker 列表和分区元数据。serializer.class在 0.7.1 里是写全限定类名的,不像后来用value.serializer和key.serializer分开。producer.type选sync时每条消息同步发送,吞吐低但可靠;选async时批量发送,需要额外配queue.buffering.max.ms和batch.size。compression.codec如果选 snappy,必须确保snappy-java-1.0.4.1.jar在 classpath 里,否则运行时会报UnsatisfiedLinkError,因为 snappy 依赖本地库。
3.2 消费者代码里的 offset 管理逻辑
消费者示例通常长这样:
// Kafka 0.7.1 消费者配置示例 Properties props = new Properties(); props.put("zk.connect", "127.0.0.1:2181"); props.put("groupid", "test-group"); // 消费者组 ID props.put("zk.sessiontimeout.ms", "10000"); props.put("zk.synctime.ms", "200"); props.put("autocommit.interval.ms", "1000"); // 自动提交 offset 间隔 ConsumerConfig config = new ConsumerConfig(props); ConsumerConnector connector = Consumer.createJavaConsumerConnector(config); // 指定主题和分区数,注意 0.7.1 的 API 和 0.8 以后不同 Map<String, Integer> topicCountMap = new HashMap<String, Integer>(); topicCountMap.put("test-topic", new Integer(1)); Map<String, List<KafkaStream<byte[], byte[]>>> consumerMap = connector.createMessageStreams(topicCountMap); KafkaStream<byte[], byte[]> stream = consumerMap.get("test-topic").get(0); // 消费消息 ConsumerIterator<byte[], byte[]> it = stream.iterator(); while (it.hasNext()) { System.out.println("收到消息: " + new String(it.next().message())); }groupid决定 offset 存在 ZooKeeper 的哪个路径下,默认是/consumers/<groupid>/offsets/<topic>/<partition>。autocommit.interval.ms控制自动提交频率,设太短会增加 ZooKeeper 写压力,设太长则消费者崩溃后重复消费的消息多。createMessageStreams的第二个参数是每个主题要开几个流,0.7.1 里一个流对应一个分区,如果你有 3 个分区但只开 1 个流,另外两个分区的消息就没人消费。这是新手最容易翻车的地方——以为消费者组会自动负载均衡,实际上 0.7.1 的流数必须和分区数匹配,否则消息积压。
3.3 消息从 Web 服务器到 Kafka 的典型链路
假设你有一个 Web 服务器,想把访问日志实时推到 Kafka。在 0.7.1 时代,常见做法是在 Web 应用的过滤器或拦截器里,拿到请求日志后调用生产者 API 发送。代码结构大致是:初始化一个单例Producer,在doFilter里构造KeyedMessage<String, String>,然后producer.send(message)。这里有个坑:Producer是线程安全的,但如果你每次请求都 new 一个,ZooKeeper 连接会被打爆。正确做法是用静态块或 Spring 单例管理一个Producer实例,整个应用共享。另外,send在sync模式下会阻塞,如果 Kafka 集群响应慢,Web 请求线程会被拖住,所以生产环境一般用async模式,并设置queue.buffering.max.ms和queue.buffering.max.messages做缓冲。
注意:0.7.1 的
async生产者没有回调机制,消息丢了只能靠日志发现,这是和 0.8 以后版本最大的体验差距。
4. 避坑与排查:0.7.1 客户端特有的五个血泪经验
4.1 连不上 ZooKeeper 却报 Kafka 异常
现象:启动生产者后抛出kafka.common.KafkaException: Failed to connect to zookeeper,但异常栈里混着KeeperException。原因:0.7.1 的生产者把 ZooKeeper 连接失败包装成了 Kafka 异常,容易误导你去查 broker。解决:先用zkCli或telnet 127.0.0.1 2181确认 ZooKeeper 可达,再检查zk.connect地址有没有写错端口。如果 ZooKeeper 是远程的,确认防火墙放行。
4.2 snappy 压缩导致的 UnsatisfiedLinkError
现象:配置compression.codec=2后,发送消息时报java.lang.UnsatisfiedLinkError: no snappyjava in java.library.path。原因:snappy-java-1.0.4.1.jar里包含本地库,但某些 JDK 版本或操作系统下解压路径有问题。解决:把compression.codec改回0先跑通,或者手动把 jar 里的.so/.dll解压到java.library.path包含的目录。更稳妥的做法是升级 snappy-java 到 1.0.5 以上,但要注意和 Kafka 0.7.1 的兼容性。
4.3 消费者收不到消息但生产者说发送成功
现象:生产者日志显示消息已发送,消费者it.hasNext()一直返回 false。原因:消费者组 ID 和生产者主题不匹配,或者消费者启动前 offset 已经被提交到了末尾。0.7.1 的消费者默认从 ZooKeeper 里存的 offset 开始消费,如果之前有另一个消费者提交过 offset,新消费者会从那个位置继续,而不是从头。解决:换一个全新的groupid,或者手动删除 ZooKeeper 里/consumers/<groupid>节点。命令是zkCli里执行rmr /consumers/test-group。
4.4 多分区主题下消费者只消费一个分区
现象:主题设了 3 个分区,消费者只收到部分消息。原因:createMessageStreams里topicCountMap.put("test-topic", 1)只创建了一个流,0.7.1 不会自动为每个分区创建流。解决:把数字改成分区数,或者用多个消费者线程,每个线程一个流。注意流数不能超过分区数,否则多出来的流会空转。
4.5 编译通过但运行时报 NoClassDefFoundError
现象:javac编译没问题,java运行时报NoClassDefFoundError: org/apache/log4j/Logger或scala/collection/...。原因:classpath 里漏了log4j-1.2.15.jar或scala-library.jar。Kafka 0.7.1 是 Scala 写的,kafka-0.7.1.jar依赖 Scala 运行时库,但压缩包里没有显式列出scala-library.jar。解决:确认 lib 目录下是否有 Scala 库,如果没有,需要单独下载和 Kafka 0.7.1 匹配的 Scala 2.8 或 2.9 版本。常见做法是去 Kafka 0.7.1 的发行包里找scala-library-2.8.0.jar。
5. 进阶验证与参数调优:用最小集群验证消息顺序性
5.1 单机伪集群的 ZooKeeper 与 broker 配置
要验证消息顺序性,至少需要两个分区和一个消费者组。单机环境下可以起一个 ZooKeeper 和一个 Kafka broker,broker 配置里num.partitions=2。Kafka 0.7.1 的 broker 配置文件在config/server.properties,关键项:
| 参数 | 值 | 说明 |
|---|---|---|
zk.connect | 127.0.0.1:2181 | ZooKeeper 地址 |
brokerid | 0 | broker 唯一 ID |
port | 9092 | broker 监听端口 |
num.partitions | 2 | 默认分区数 |
log.dirs | /tmp/kafka-logs | 日志目录 |
启动 broker 用bin/kafka-server-start.sh config/server.properties。启动后确认 ZooKeeper 里/brokers/ids/0节点存在。
5.2 用生产者发送带 key 的消息验证分区路由
Kafka 0.7.1 的分区路由规则是:如果消息有 key,用key.hashCode() % numPartitions决定分区;如果没有 key,轮询。要验证顺序性,发送同一 key 的消息,它们会落到同一分区,消费者按分区顺序消费。代码片段:
// 发送带 key 的消息,同一 key 保证落到同一分区 for (int i = 0; i < 10; i++) { String key = "order-001"; // 同一个 key String value = "消息序号: " + i; KeyedMessage<String, String> message = new KeyedMessage<String, String>("test-topic", key, value); producer.send(message); }消费者端用两个流分别消费两个分区,打印消息时带上分区 ID。如果同一 key 的消息在同一个流里按序号递增,说明顺序性成立。注意 0.7.1 的async生产者可能因为重试导致乱序,验证时用sync模式。
5.3 从 offset 提交间隔看重复消费边界
autocommit.interval.ms设成 1000 时,消费者每秒提交一次 offset。如果你在两次提交之间 kill 消费者,重启后会从上次提交的 offset 重新消费,最多重复 1 秒的消息。要精确控制,可以关掉自动提交,手动调用connector.commitOffsets()。但 0.7.1 的手动提交 API 比较简陋,需要自己拿TopicCount和ConsumerConnector配合。我一般会先把autocommit.interval.ms设成 5000,观察重复消费的量,再决定是否值得改手动提交。从那以后我每次做 Kafka 消费者测试,都强制先跑一遍 kill -9 再重启,看重复消息的边界在哪里,这个习惯帮我省了很多线上排查时间。希望帮到你。
本文还有配套的精品资源,点击获取