springboot入门14 – Kafka应用

简述

这几天优化了一下之前写的一个springboot kafka组件。比较起原生的spring-kafka来,我希望能够简化kafka的使用,可以更聚焦于具体的消息处理逻辑。

接下来的内容是这个组件的用法。

使用方法

添加依赖

这个组件已经提交到了maven中央仓库,可以直接通过依赖的形式引入:

Kafka java.net.SocketTimeoutException

手上有一个消费Kafka的服务,这个服务每隔一段时间使用SimpleConsumer从kafka集群获取数据。Kafka的版本是0.8.21。今天这个服务一直在报错:SocketTimeoutException。异常信息如下:

Kafka 调整partiton数目和replica factor

调整partiton

调整partition可以直接执行如下命令:

注意替换topicName、$ZK_HOST_NODE和partitionNum三个参数。

调整replica

kafka0.9 Consumer poll()方法阻塞

最近项目中用到了Kafka0.9,在使用0.9的Consumer API的时候遇到了poll()方法阻塞的问题。程序没有报任何错误,只是持续在poll()方法处阻塞。深入poll()方法可以看到是在AbstractCoordinator.ensureCoordinatorKnown()方法中出现了死循环。在循环中不停地输出如下DEBUG日志:

获取Kafka Consumer的offset

从kafka的0.8.11版本开始,它会将consumer的offset提交给ZooKeeper。然而当offset的数量(consumer数量 * partition的数量)的很多的时候,ZooKeeper的适应性就可能会出现不足。幸运的是,Kafka现在提供了一种理想的机制来存储Consumer的offset。Kafka现在是将Consumer的offset…

Kafka high level consumer

为什么选择High Level Consumer

很多时候用户只是想从kafka中读取数据,对于如何处理消息的offset则不怎么关注。在抽象了Kafka中消费事件的大部分细节后,High Level Consumer可以让用户使用起来更为简单。

首先需要知道的就是High Level Consumer保存了从zookeeper中读取某个partition最后的offset…