本文由 GodPan 发表在 ScalaCool 团队博客。

看完上一篇,相信大家对消息系统以及Kafka的整体构成都有了初步了解,学习一个东西最好的办法,就是去使用它,今天就让我们一起窥探一下Kafka,并完成自己的处女作。

消息在Kafka中的历程

虽然我们掌握东西要一步一步来,但是我们在大致了解了一个东西后,会有利于我们对它的理解和学习,所以我们可以先来看一下一条消息从发出到最后被消息者接收到底经历了什么?

上图简要的说明了消息在Kafka中的整个流转过程(假设已经部署好了整个Kafka系统,并创建了相应的Topic,分区等细节后续再单独讲):

  • 1.消息生产者将消息发布到具体的Topic,根据一定算法或者随机被分发到具体的分区中;
  • 2.根据实际需求,是否需要实现处理消息逻辑;
  • 3.若需要,则实现具体逻辑后将结果发布到输出Topic;
  • 4.消费者根据需求订阅相关Topic,并消费消息;

总的来说,怎么流程还是比较清晰和简单的,下面就跟我一起来练习Kafka的基本操作,最后实现一个单词计数的小demo。

基础操作

以下代码及相应测试在以下环境测试通过:Mac OS + JDK1.8,Linux系统应该也能跑通,Windows有兴趣的同学可以去官网下载相应版本进行相应的测试练习。

下载Kafka

Mac系统同学可以使用brew安装:

brew install kafka
复制代码

Linux系统同学可以从官网下载源码解压,也可以直接执行以下命令:

cd
mkdir test-kafka && cd test-kafka
curl -o kafka_2.11-1.0.1.tgz http://mirrors.tuna.tsinghua.edu.cn/apache/kafka/1.0.1/kafka_2.11-1.0.1.tgz
tar -xzf kafka_2.11-1.0.1.tgz
cd kafka_2.11-1.0.1复制代码

启动

Kafka使用Zookeeper来维护集群信息,所以这里我们先要启动Zookeeper,Kafka与Zookeeper的相关联系跟结合后续再深入了解,毕竟不能一口吃成一个胖子。

bin/zookeeper-server-start.sh config/zookeeper.properties
复制代码

接着我们启动一个Kafka Server节点:

bin/kafka-server-start.sh config/server.properties
复制代码

这时候Kafka系统已经算是启动起来了。

创建Topic

在一切就绪之后,我们要开始做极其重要的一步,那就是创建Topic,Topic是整个系统流转的核心,另外Topic本身也包含着很多复杂的参数,比如复制因子个数,分区个数等,这里为了从简,我们将对应的参数都设为1,方便大家测试:

bin/kafka-topics.sh --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic kakfa-test
复制代码

其中参数的具体含义:

属性 功能
--create 代表创建Topic
--zookeeper zookeeper集群信息
--replication-factor 复制因子
--partitions 分区信息
--topic Topic名称

这时候我们已经创建好了一个叫kakfa-test的Topic了。

向Topic发送消息

在有了Topic后我们就可以向其发送消息:

bin/kafka-console-producer.sh --broker-list localhost:9092 --topic kakfa-test
复制代码

然后我们向控制台输入一些消息:

this is my first test kafka
so good
复制代码

这时候消息已经被发布在kakfa-test这个主题上了。

从Topic获取消息

现在Topic上已经有消息了,现在可以从中获取消息被消费:

bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic kafka-test --from-beginning
复制代码

这时候我们可以在控制台看到:

this is my first test kafka
so good
复制代码

至此我们就测试了最简单的Kafka Demo,希望大家能自己动手去试试,虽然很简单,但是这能让你对整个Kafka流程能更熟悉。

WordCount

下面我们来利用上面的一些基本操作来实现一个简单WordCount程序,它具备以下功能:

  • 1.支持词组持续输入,即生产者不断生成消息;
  • 2.程序自动从输入Topic中获取原始数据,然后经过处理,将处理结果发布在计数Topic中;
  • 3.消费者可以从计数Topic获取相应的WordCount的结果;

1.启动kafka

与上文的启动一样,按照其操作即可。

2.创建输入Topic

bin/kafka-topics.sh --zookeeper localhost:2181 --create --topic kafka-word-count-input --partitions 1 --replication-factor 1
复制代码

3.向Topic输入消息

bin/kafka-console-producer.sh --broker-list localhost:9092 --topic kafka-word-count-input
复制代码

4.流处理逻辑

这部分内容是整个例子的核心,这部分代码有Java 8+和Scala版本,个人认为流处理用函数式语法表达的更加简洁清晰,推荐大家用函数式的思维去尝试写以下,发现自己再也不想写Java匿名内部类这种语法了。

我们先来看一个Java 8的版本:

public class WordCount {public static void main(String[] args) throws Exception {Properties props = new Properties();props.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-word-count");props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());final StreamsBuilder builder = new StreamsBuilder();KStream<String, String> source = builder.<String, String>stream("kafka-word-count-input");Pattern pattern = Pattern.compile("\\W+");source.flatMapValues(value -> Arrays.asList(pattern.split(value.toLowerCase(Locale.getDefault())))).groupBy((key, value) -> value).count(Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("counts-store")).mapValues(value->Long.toString(value)).toStream().to("kafka-word-count-output");final KafkaStreams streams = new KafkaStreams(builder.build(), props);streams.start();}
}
复制代码

是不是很惊讶,用java也能写出如此简洁的代码,所以说如果有适用场景,推荐大家尝试的用函数式的思维去写写java代码。

我们再来看看Scala版本的:


object WordCount {def main(args: Array[String]) {val props: Properties = {val p = new Properties()p.put(StreamsConfig.APPLICATION_ID_CONFIG, "kafka-word-count")p.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092")p.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String.getClass)p.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String.getClass)p}val builder: StreamsBuilder = new StreamsBuilder()val source: KStream[String, String] = builder.stream("kafka-word-count-input")source.flatMapValues(textLine => textLine.toLowerCase.split("\\W+").toIterable.asJava).groupBy((_, word) => word).count(Materialized.as[String, Long, KeyValueStore[Bytes, Array[Byte]]]("counts-store")).toStream.to("kafka-word-count-output")val streams: KafkaStreams = new KafkaStreams(builder.build(), props)streams.start()}
}
复制代码

可以发现使用Java 8函数式风格编写的代码已经跟Scala很相似了。

5.启动处理逻辑

很多同学电脑上并没有装sbt,所以这里演示的利用Maven构建的Java版本,具体执行步骤请参考戳这里kafka-word-count上的说明。

6.启动消费者进程

最后我们启动消费者进程,并在生产者中输入一些单词,比如:

最后我们可以在消费者进程中看到以下输出:

bin/kafka-console-consumer.sh --topic kafka-word-count-output --from-beginning --bootstrap-server localhost:9092  --property print.key=true
复制代码

总结

本篇文章主要是讲解了Kafka的基本运行过程和一些基础操作,但这是我们学习一个东西必不可少的一步,只有把基础扎实好,才能更深入的去了解它,理解它为什么这么设计,我在这个过程中也遇到很多麻烦,所以还是希望大家能够自己动手去实践一下,最终能收获更多。

Kafka 学习笔记(二) :初探 Kafka相关推荐

  1. kafka学习(二)kafka工作流程分析

    本文借鉴:再过半小时,你就能明白kafka的工作原理了(特此感谢!) 一.发送数据 PS:Producer在写入数据的时候永远的找leader,不会直接将数据写入follower 1.follower ...

  2. Kafka学习笔记(一):Kafka基本概念理解

    ActiveMQ.RabbitMQ是用的比较多得消息队列,但是随着时间的推移,大数据的应运而生,这两种消息队列使用的也是越来越少了,Kafka渐渐进入开发人员的视线,再加上Kafka天生的集群运行.大 ...

  3. 大数据 -- kafka学习笔记:知识点整理(部分转载)

    一 为什么需要消息系统 1.解耦 允许你独立的扩展或修改两边的处理过程,只要确保它们遵守同样的接口约束. 2.冗余 消息队列把数据进行持久化直到它们已经被完全处理,通过这一方式规避了数据丢失风险.许多 ...

  4. Kafka学习笔记(一):什么是消息队列?什么是Kafka?

    目录 一.消息队列的概述 (一)前置知识点 1.集群和分布式 2.队列(Queue)的含义 3.同步与异步的含义 (二)消息队列的含义与特点 二.Kafka (一) 概述 (二) 常用名词含义 导航栏 ...

  5. Kafka学习笔记(3)----Kafka的数据复制(Replica)与Failover

    1. CAP理论 1.1 Cosistency(一致性) 通过某个节点的写操作结果对后面通过其他节点的读操作可见. 如果更新数据后,并发访问的情况下可立即感知该更新,称为强一致性 如果允许之后部分或全 ...

  6. qml学习笔记(二):可视化元素基类Item详解(上半场anchors等等)

    原博主博客地址:http://blog.csdn.net/qq21497936 本文章博客地址:http://blog.csdn.net/qq21497936/article/details/7851 ...

  7. [转载]dorado学习笔记(二)

    原文地址:dorado学习笔记(二)作者:傻掛 ·isFirst, isLast在什么情况下使用?在遍历dataset的时候会用到 ·dorado执行的顺序,首先由jsp发送请求,调用相关的ViewM ...

  8. PyTorch学习笔记(二)——回归

    PyTorch学习笔记(二)--回归 本文主要是用PyTorch来实现一个简单的回归任务. 编辑器:spyder 1.引入相应的包及生成伪数据 import torch import torch.nn ...

  9. tensorflow学习笔记二——建立一个简单的神经网络拟合二次函数

    tensorflow学习笔记二--建立一个简单的神经网络 2016-09-23 16:04 2973人阅读 评论(2) 收藏 举报  分类: tensorflow(4)  目录(?)[+] 本笔记目的 ...

  10. Scapy学习笔记二

    Scapy学习笔记二 Scapy Sniffer的用法: http://blog.csdn.net/qwertyupoiuytr/article/details/54670489 Scapy Snif ...

最新文章

  1. Java如何校验两个文件内容是相同的?
  2. 拼图游戏 复制粘贴一个叫lemene的人的,这个人是c++博客的用户,我不是,怕以后找不到这篇文章,所以复制粘贴了。文中最后给出了原文链接连接...
  3. [数论]Gcd/ExGcd欧几里得学习笔记
  4. gesturedetector.java_android使用gesturedetector手势识别示例分享
  5. spark 用户画像挖掘分析_如何基于Spark进行用户画像?
  6. 程序员面试金典 - 面试题 16.17. 连续数列(DP/分治)
  7. Linux 多用户和多用户边界
  8. [ActionScript 3.0] 通过as3操作web内容
  9. day24-XSS过滤及单实例
  10. mysql单表查询怎么做_mysql单表查询
  11. 树莓派3连接ps4无线手柄
  12. updating homebrew
  13. mysql - rank函数的使用
  14. 金山云CDN:国内最佳付费CDN
  15. RuoYi-Vue——关于登录后不同角色跳不同页面
  16. qq修改群名服务器失败,qq建群失败什么原因 q群一直创建失败 - 云骑士一键重装系统...
  17. 高速串行总线设计基础(七)揭秘SERDES高速面纱之时钟校正与通道绑定技术
  18. 计算机系统使用寿命,笔记本电脑的使用寿命是多久,你了解吗?
  19. 信息系统面临的安全威胁
  20. SpringCloud动态获取yml文件中的配置

热门文章

  1. php自己遇到的一些问题
  2. HTML5-Video(视屏播放)
  3. Android的数据库(SQLite)学习
  4. iview 级联选择组件_iView Cascader级联选择器
  5. 定时重启_SpringBoot基于数据库的定时任务实现方法
  6. ThinkPhp 使用 PHP_XLSXWriter 代替 PHPExcel 百万级数据单次导出
  7. Go新手上路(时不时更新)
  8. phpspider 爬取汉谜网
  9. PHP可变变量($$)
  10. Linux用php上传表单文件,文件太大提示[413 Request Entity Too Large]