简述
这几天优化了一下之前写的一个springboot kafka组件。比较起原生的spring-kafka来,我希望能够简化kafka的使用,可以更聚焦于具体的消息处理逻辑。
接下来的内容是这个组件的用法。
使用方法
添加依赖
这个组件已经提交到了maven中央仓库,可以直接通过依赖的形式引入:
|
1 2 3 4 5 |
<dependency> <group |
这几天优化了一下之前写的一个springboot kafka组件。比较起原生的spring-kafka来,我希望能够简化kafka的使用,可以更聚焦于具体的消息处理逻辑。
接下来的内容是这个组件的用法。
这个组件已经提交到了maven中央仓库,可以直接通过依赖的形式引入:
|
1 2 3 4 5 |
<dependency> <group |
手上有一个消费Kafka的服务,这个服务每隔一段时间使用SimpleConsumer从kafka集群获取数据。Kafka的版本是0.8.21。今天这个服务一直在报错:SocketTimeoutException。异常信息如下:
|
1 2 3 4 5 6 7 8 9 10 11 12 13 14 |
2018-12-22 15:18:32 INFO SimpleConsumer |
很多时候用户只是想从kafka中读取数据,对于如何处理消息的offset则不怎么关注。在抽象了Kafka中消费事件的大部分细节后,High Level Consumer可以让用户使用起来更为简单。
首先需要知道的就是High Level Consumer保存了从zookeeper中读取某个partition最后的offset…
Kafka的Simple Consumer并不简单。相反的,它是直接在对Kafka的partition和offset进行操作,是一种比较接近底层的方案。
使用SimpleConsumer最主要的原因是想在消费消息时获取更大的权限。比如说要做下面这些事情:
这次是要写一个Producer示例程序。使用的kafka版本是0.8.2。开发语言是java。
程序主体是一个Producer类。这个Producer类主要是用来为指定的topic创建消息。
首先需要引入一些支持类:
|
1 2 3 |
import kafka.javaapi.producer.Producer; import kafka.producer.KeyedMes |
kafka是一个分布式的、可分区的、可复制的日志提交服务。它提供了消息传递的功能,但是有着独特的设计。
首先,先了解一些基础概念:
从…