ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Kafka 0.7.1 客户端实战:从 ZooKeeper 依赖到消息顺序性验证

Kafka 0.7.1 客户端实战:从 ZooKeeper 依赖到消息顺序性验证 简介这份资源是面向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 变成可运行的 classpath2.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 下用分号分隔 classpathLinux/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 startWindows 下是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 gzip2 snappy ProducerConfig config new ProducerConfig(props); ProducerString, String producer new ProducerString, 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 以后不同 MapString, Integer topicCountMap new HashMapString, Integer(); topicCountMap.put(test-topic, new Integer(1)); MapString, ListKafkaStreambyte[], byte[] consumerMap connector.createMessageStreams(topicCountMap); KafkaStreambyte[], byte[] stream consumerMap.get(test-topic).get(0); // 消费消息 ConsumerIteratorbyte[], 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里构造KeyedMessageString, 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.codec2后发送消息时报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 brokerbroker 配置里num.partitions2。Kafka 0.7.1 的 broker 配置文件在config/server.properties关键项参数值说明zk.connect127.0.0.1:2181ZooKeeper 地址brokerid0broker 唯一 IDport9092broker 监听端口num.partitions2默认分区数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; KeyedMessageString, String message new KeyedMessageString, 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 再重启看重复消息的边界在哪里这个习惯帮我省了很多线上排查时间。希望帮到你。本文还有配套的精品资源点击获取
返回列表