关于Kafka的版本号

这个内容实在是太重要了,甚至是你能否用好 Kafka 的关键。

Kafka 流行的几种 Kafka 发行版本质上都内嵌了最核心的 Apache Kafka,也就是社区版 Kafka,那今天我们就来说说 Apache Kafka 版本号的问题。

那么现在你可能会有这样的疑问:我为什么需要关心版本号的问题呢?直接使用最新版本不就好了吗?当然了,这的确是一种有效的选择版本的策略,但我想强调的是这种策略并非在任何场景下都适用。如果你不了解各个版本之间的差异和功能变化,你怎么能够准确地评判某 Kafka 版本是不是满足你的业务需求呢?因此在深入学习 Kafka 之前,花些时间搞明白版本演进,实际上是非常划算的一件事。

Kafka 版本命名

当前 Apache Kafka 已经迭代到 2.6版本。但是很多人对于 Kafka 的版本命名理解存在歧义。比如我们在官网上下载 Kafka 时,会看到这样的版本:

于是有些同学就会纳闷,难道 Kafka 版本号不是 2.11 或 2.12 吗?其实不然,前面的版本号是编译 Kafka 源代码的 Scala 编译器版本。

现在你应该知道了对于 kafka-2.11-2.1.1 的提法,真正的 Kafka 版本号实际上是 2.1.1。前面的 2 表示大版本号,即 Major Version;中间的 1 表示小版本号或次版本号,即 Minor Version;最后的 1 表示修订版本号,也就是 Patch 号。

Kafka 版本演进

Kafka 目前总共演进了 7 个大版本,分别是 0.7、0.8、0.9、0.10、0.11、1.0 和 2.0,其中的小版本和 Patch 版本很多。

  • 0.7 大版本
    • “上古”版本,提供了最基础的消息队列功能,连副本机制都没有,不推荐。
  • 0.8大版本
    • 正式引入了副本机制,至此 Kafka 成为了一个真正意义上完备的分布式高可靠消息队列解决方案。有了副本备份机制,Kafka 就能够比较好地做到消息无丢失。
    • 建议你至少要升级到 0.8.2.2 这个版本,因为该版本中老版本消费者 API 是比较稳定的。另外即使你升到了 0.8.2.2,也不要使用新版本 Producer API,此时它的 Bug 还非常多。
  • 0.9大版本
    • 增加了基础的安全认证 / 权限功能,同时使用 Java 重写了新版本消费者 API,另外还引入了 Kafka Connect 组件用于实现高性能的数据抽取。
    • 新版本 Producer API 在这个版本中算比较稳定了。但不要使用新版本 Consumer API,因为此时很不稳定。
  • 0.10.0.0大版本
    • 里程碑式的大版本,因为该版本引入了 Kafka Streams。从这个版本起,Kafka 正式升级成分布式流处理平台,虽然此时的 Kafka Streams 还基本不能线上部署使用。
    • 如果你依然在使用 0.10 大版本,建议你至少升级到 0.10.2.2 然后使用新版本 Consumer API,另外此版本的 Producer 性能也更好。
  • 0.11.0.0 大版本
    • 主流的版本之一,引入了幂等性 Producer API 以及事务(Transaction) API;重构了 Kafka 消息格式。
    • 至少将你的环境升级到 0.11.0.3,因为这个版本的消息引擎功能已经非常完善了。
  • 1.0 和 2.0大版本的区别
    • 这两个大版本主要还是 Kafka Streams 的各种改进,如果你是 Kafka Streams 的用户,至少选择 2.0.0 版本吧。

建议不论你用的是哪个版本,都请尽量保持服务器端版本和客户端版本一致,否则你将损失很多 Kafka 为你提供的性能优化收益。

小结

每个 Kafka 版本都有它恰当的使用场景和独特的优缺点,不要一味追求最新版本,选择稳妥且合适的才是最好的。

Views: 174

APACHE KAFKA 快速起步

Kafka快速起步指南。

1、获取KAFKA

下载 最新版本的Kafka并解压:

$ tar -xzf kafka_2.13-2.6.0.tgz
$ cd kafka_2.13-2.6.0

2、准备Kafka环境

注意: 本地环境需要安装 Java 8+

运行下面命令开启Kafka中内置的Zookeeper服务

# 启动Kafka中内置的ZooKeeper服务
# 注意: ZooKeeper框架未来将从Kafka框架中分离出去
$ bin/zookeeper-server-start.sh config/zookeeper.properties

打开另一个终端运行Kafka代理服务:

# 开启Kafka代理服务
$ bin/kafka-server-start.sh config/server.properties

当所有服务都成功启动,一个简单的单机Kafka环境就准备好了.

3、创建一个消息主题

Kafka是一个分布式事件流平台,它允许您跨多台机器读、写、存储和处理消息(在官方文档中也称为事件 )。

示例事件有支付事务、来自移动电话的地理位置更新、发货订单、来自物联网设备或医疗设备的传感器测量等等。这些事件被组织并存储在 主题中。非常简单,主题类似于文件系统中的文件夹,事件是该文件夹中的文件。

因此,在编写第一个事件之前,必须创建一个主题。打开另一个终端会话并运行:

$ bin/kafka-topics.sh --create --topic quickstart-events --bootstrap-server localhost:9092

所有Kafka的命令行工具都有额外的选项: 运行Kafka -topic .sh命令而不带任何参数来显示使用信息。例如,它还可以显示详细信息,如新主题的分区数量:

$ bin/kafka-topics.sh --describe --topic quickstart-events --bootstrap-server localhost:9092
Topic:quickstart-events  PartitionCount:1    ReplicationFactor:1 Configs:
    Topic: quickstart-events Partition: 0    Leader: 0   Replicas: 0 Isr: 0

4、向对应主题写入消息

Kafka客户端通过网络与Kafka代理服务通信以读写消息。一旦接收到这些消息,代理将以持久和容错的方式存储这些事件,您需要多长时间就存储多长时间——甚至永远存储。

运行控制台生成器客户端,在主题中写入一些事件。默认情况下,您输入的每一行都将导致向主题写入一个单独的消息。

$ bin/kafka-console-producer.sh --topic quickstart-events --bootstrap-server localhost:9092
This is my first event
This is my second event

您可以在任何时候使用Ctrl-C停止生产者(producer)客户端

5、读取消息

打开另一个终端会话并运行控制台消费者客户端(kafka-console-consumer)来读取您刚刚创建的事件:

$ bin/kafka-console-consumer.sh --topic quickstart-events --from-beginning --bootstrap-server localhost:9092
This is my first event
This is my second event

您可以在任何时候使用Ctrl-C停止客户端。

您可以自由试验: 例如,切换回您的生产者终端(上一步)来编写额外的事件,并查看事件如何立即显示在您的消费者终端中。

因为事件长期存储在Kafka中,所以它们可以被任意多的消费者读取。您可以通过打开另一个终端会话并再次运行之前的命令来轻松地验证这一点。

6、使用KAFKA CONNECT将数据导入/导出为消息流

您可能在关系数据库或传统消息传递系统等现有系统中拥有大量数据,以及许多已经使用这些系统的应用程序。Kafka Connect允许您不断地从外部系统摄取数据到Kafka,反之亦然。因此,与Kafka集成现有系统是非常容易的。为了使这个过程更容易,有数百个这样的连接器可用。

看一看Kafka Connect section部分,了解更多关于如何持续地导入/导出数据到Kafka和导出。

7、使用KAFKA STREAMS处理消息

一旦您的数据以消息的形式存储在Kafka中,您就可以使用 Kafka Streams 客户端的Java/Scala库处理数据。它允许您实现关键任务的实时应用程序和微服务,其中输入和/或输出数据存储在Kafka主题中。

Kafka Streams将客户端编写和部署标准Java和Scala应用程序的简单性与Kafka服务器端集群技术的优势结合在一起,使这些应用程序具有高度的可伸缩性(highly scalable)、弹性(elastic)、容错性(fault-tolerant)和分布式(distributed)。这个库完全支持精确一次处理(exactly-once processing)的原语、有状态(stateful )操作和聚合(aggregation)、窗口(windowing)、连接(join)、基于事件时间的处理等等。

给你一个初步的体验,这里是如何实现流行的WordCount算法:

KStream<String, String> textLines = builder.stream("quickstart-events");

KTable<String, Long> wordCounts = textLines
            .flatMapValues(line -> Arrays.asList(line.toLowerCase().split(" ")))
            .groupBy((keyIgnored, word) -> word)
            .count();

wordCounts.toStream().to("output-topic"), Produced.with(Serdes.String(), Serdes.Long()));

Kafka 流处理演示演示和应用程序部署教程演示了如何从头到尾编写和运行这样的流处理应用程序。

8、停止KAFKA

既然你已经完成了《快速起步》,那么你可以随时停止Kafka运行——或者继续练习。

  1. 使用Ctrl-C停止生产者和消费者客户端。

  2. Ctrl-C停止Kafka代理服务器。

  3. 最后,用Ctrl-C停止ZooKeeper服务器。

如果您还想删除您本地Kafka环境的任何数据,包括您在此过程中创建的任何事件,运行命令:

$ rm -rf /tmp/kafka-logs /tmp/zookeeper

恭喜你!

您已经成功地完成了Apache Kafka快速启动。

为了解更多详情,我们建议以下步骤:

Views: 461

3-使用storm-starter测试集群

安装maven

官方示例位于storm安装文件夹下面example下的storm-starter下

安装maven3.6 (3.5也可)

下载地址

配置环境变量 vi ~/.bash_profile

#MAVEN
export M2_HOME=/home/hadoop/app/apache-maven-3.6.3
PATH=$M2_HOME/bin:$PATH

修改配置

/conf/settings.xml

<?xml version="1.0" encoding="UTF-8"?>
<settings xmlns="http://maven.apache.org/SETTINGS/1.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance" xsi:schemaLocation="http://maven.apache.org/SETTINGS/1.0.0 http://maven.apache.org/xsd/settings-1.0.0.xsd">
  <localRepository>/home/hadoop/app/apache-maven-repo/</localRepository>
  <pluginGroups>
    <pluginGroup>org.mortbay.jetty</pluginGroup>
  </pluginGroups>
  <proxies></proxies>
  <servers></servers>
  <mirrors>
    <mirror>
      <id>alimaven</id>
      <name>aliyun maven</name>
      <url>http://maven.aliyun.com/nexus/content/groups/public/</url>
      <mirrorOf>central</mirrorOf>
    </mirror>
    <mirror>
      <id>clojars</id>
      <name>clojar-maven</name>
      <!--url>http://clojars.org/repo/</url-->
      <url>https://mirrors.tuna.tsinghua.edu.cn/clojars/</url>
      <mirrorOf>clojars</mirrorOf>
    </mirror>
  </mirrors>
  <profiles>
  <profile>
    <id>jdk-1.4</id>
    <activation>
      <jdk>1.4</jdk>
    </activation>
    <repositories>
      <repository>
        <id>alimaven</id>
        <name>aliyun maven</name>
        <url>http://maven.aliyun.com/nexus/content/groups/public/</url>
        <releases>
          <enabled>true</enabled>
        </releases>
        <snapshots>
          <enabled>false</enabled>
        </snapshots>
      </repository>
      <repository>
        <id>clojars</id>
        <url>https://mirrors.tuna.tsinghua.edu.cn/clojars/</url>
      </repository>
    </repositories>
  </profile>
  </profiles>
</settings>

注意:
默认的本地仓库(localRepository)是在~/.m2/repository/
也可以根据需要自定义在合适的位置。

编译打包starter项目

编译前,如果你的虚拟机可用内存小于4G,需要首先修改一下文件的并行度,避免虚拟机资源不足.

starter项目根目录/src/jvm/org.apache.storm.starter.WordCountTopology

image-20200915002404862

在storm-starter下面执行

mvn clean

此时会启动依赖下载过程,根据网络情况可能十几分钟到半个小时...

mvn package -Dmaven.test.skip=true

一定要忽略测试过程,不然一定会报错

找到build出来的胖包(体积最大的那个)

image-20200915002448201

运行项目

首先运行集群,保证相应的进程都开启

image-20200915002535892

本地模式运行

[hadoop@hadoop00 target]$ storm local original-storm-starter-2.1.0.jar org.apache.storm.starter.WordCountTopology WCTopology

file

上传到集群运行

[hadoop@hadoop00 target]$ storm jar original-storm-starter-2.1.0.jar org.apache.storm.starter.WordCountTopology WCTopology

image-20200915002609935

过一会显示上传成功

image-20200915002623430

image-20200915002637640

查看UI界面

查看UI

file

由于开启了logviewer点击worker端口号可以查看日志(如果没有开启日志则打开链接显示无法访问)

file

如果关闭了logviewer, 则ui界面中nimbus日志不能看

日志分析

要注意的是:

分配worker进程数量, 日志中显示的slots就是worker进程数, 这里是2个

分配executor线程, 此时注意每一个线程里面都有两个数字[m,n] 代表这个线程里面会执行n-m+1个任务, 编号分别为n~m

一个是6个任务线程, 任务线程编号可以从后面的日志分析出来分别是51-56

img

Spout Executor只有1个, 任务号为11

img

Bolt Executor有10个

其中

split 开启2个线程, 一共4个任务, 因此1个线程分配有2个任务

count 开启3个线程, 一共6个任务, 因此1个线程分配有2个任务

img

比如这里是分词bolt的日志

img

统计bolt的日志

img

代码分析

package org.apache.storm.starter;

import java.util.HashMap;
import java.util.Map;
import org.apache.storm.starter.spout.RandomSentenceSpout;
import org.apache.storm.task.ShellBolt;
import org.apache.storm.topology.BasicOutputCollector;
import org.apache.storm.topology.ConfigurableTopology;
import org.apache.storm.topology.IRichBolt;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.topology.base.BaseBasicBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

/**
  * This topology demonstrates Storm's stream groupings and multilang
  * capabilities.
  */

public class WordCountTopology extends ConfigurableTopology {
    public static void main(String[] args) throws Exception {
        ConfigurableTopology.start(new WordCountTopology(), args);
    }

    @Override
    protected int run(String[] args) throws Exception {
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout("spout", new RandomSentenceSpout(), 5);
        builder.setBolt("split", new SplitSentence(), 8).shuffleGrouping("spout");
        builder.setBolt("count", new WordCount(), 12).fieldsGrouping("split", new Fields("word"));
        conf.setDebug(true);
        String topologyName = "word-count";
        conf.setNumWorkers(3);

        if (args != null && args.length > 0) {
            topologyName = args[0];
        }

        return submit(topologyName, conf, builder);

    }

    public static class SplitSentence extends ShellBolt implements IRichBolt {

        public SplitSentence() {
            super("python", "splitsentence.py");
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("word"));
        }

        @Override
        public Map<String, Object> getComponentConfiguration() {
            return null;
        }
    }

    public static class WordCount extends BaseBasicBolt {
        Map<String, Integer> counts = new HashMap<String, Integer>();

        @Override
        public void execute(Tuple tuple, BasicOutputCollector collector) {
            String word = tuple.getString(0);
            Integer count = counts.get(word);
            if (count == null) {
                count = 0;
            }
            count++;
            counts.put(word, count);
            collector.emit(new Values(word, count));
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("word", "count"));
        }
    }
}

file

问题解决

运行storm-starter项目的WordCountTopplogy实例出现如下报错

java.lang.RuntimeException: backtype.storm.multilang.NoOutputException: Pipe to subprocess seems to be broken! No output read. Serializer Exception: Traceback (most recent call last): File "splitsentence.py", line 16, in import storm ImportError: No module named storm

解决方案

下载storm.py 放在目录 /apache-storm-2.1.0/examples/storm-starter/multilang/resources/下面 , 修改WordCountTopology.java的主函数如下:

public static void main(String[] args) throws Exception {

    SplitSentence pythonSplit = new SplitSentence();
    Map env = new HashMap();
    env.put("PYTHONPATH", "/apache-storm-2.1.0/examples/storm-starter/multilang/resources/");
    pythonSplit.setEnv(env);

    TopologyBuilder builder = new TopologyBuilder();

    builder.setSpout("spout", new RandomSentenceSpout(), 5);

    builder.setBolt("split",pythonSplit, 8).shuffleGrouping("spout");
    builder.setBolt("count", new WordCount(), 12).fieldsGrouping("split", new Fields("word"));

    Config conf = new Config();
    conf.setDebug(true);

    if (args != null && args.length > 0) {
      conf.setNumWorkers(3);

      StormSubmitter.submitTopologyWithProgressBar(args[0], conf, builder.createTopology());
    }
    else {
      conf.setMaxTaskParallelism(3);

      LocalCluster cluster = new LocalCluster();
      cluster.submitTopology("word-count", conf, builder.createTopology());

      Thread.sleep(600000);

      cluster.shutdown();
    }
  }

修改完上述文件后, maven构建报错

[ERROR] Failed to execute goal org.apache.maven.plugins:maven-checkstyle-plugin:3.0.0:checdate) on project storm-starter: You have 5 Checkstyle violations. -> [Help 1]

解决方案:

修改pom.xml, 放开代码风格检查插件的限制即可

<plugin>
  <groupId>org.apache.maven.plugins</groupId>
  <artifactId>maven-checkstyle-plugin</artifactId>
  <!--Note - the version would be inherited-->
  <configuration>
      <maxAllowedViolations>100</maxAllowedViolations>
  </configuration>
</plugin>

Views: 379

消息队列介绍

Kafka分布式消息系统被认为是一种消息引擎系统,或者消息队列中间件。

队列(Queque)是一种先入先出(Fist In First Out - FIFO)的线性表数据结构, 可以使用数组或者链表实现队列, 一个队列需要维护两个指针, head指向队首, tail指向队尾, 移动队尾添加元素(入队), 移动队首指针删除元素(出队).

image-20200515093054216

实际生活中,队列的应用随处可见,比如排队、挂号、传递过程都可以用队列来描述或者实现。

什么是消息队列

生产出美味的巧克力需要三道工序:首先将可可豆磨成可可粉,然后将可可粉加热并加入糖变成巧克力浆,最后将巧克力浆灌入模具,撒上坚果碎,冷却后就是成品巧克力了。

最开始的时候,每次研磨出一桶可可粉后,工人就会把这桶可可粉送到加工巧克力浆的工人手上,然后再回来加工下一桶可可粉。这样一来我们很快就发现,其实工人可以不用自己运送半成品,于是在每道工序之间都增加了一组传送带,研磨工人只要把研磨好的可可粉放到传送带上,就可以去加工下一桶可可粉了。 传送带解决了上下游工序之间的“通信”问题。

传送带上线后确实提高了生产效率,但也带来了新的问题:每道工序的生产速度并不相同。在巧克力浆车间,一桶可可粉传送过来时,工人可能正在加工上一批可可粉,没有时间接收。不同工序的工人们必须协调好什么时间往传送带上放置半成品,如果出现上下游工序加工速度不一致的情况,上下游工人之间必须互相等待,确保不会出现传送带上的半成品无人接收的情况。

为了解决这个问题,我们在每组传送的下游带配备了一个暂存半成品的仓库,这样上游工人就不用等待下游工人有空,任何时间都可以把加工完成的半成品丢到传送带上,无法接收的货物被暂存在仓库中,下游工人可以随时来取。传送带配备的仓库实际上起到了“通信”过程中“缓存”的作用。

img

在解决上述问题的过程中, 不知不觉我们就实现了一个消息队列.

队列能解决什么问题?

缓存数据,且保持数据存取顺序一致

为什么使用消息队列

消息异步处理

如需要快速响应的事情集中资源处理,其余的放入消息队列异步完成

很多时候,你不想也不需要立即处理消息。消息队列提供了异步处理机制,允许你把一个消息放入队列,但并不立即处理它。你想向队列中放入多少消息就放多少,然后在你想要处理的时候再去处理它们。

想一想如何设计一个秒杀系统可以达到更高的成交量。

秒杀系统需要解决的核心问题是,如何利用有限的服务器资源,尽可能多地处理短时间内的海量请求。我们知道,处理一个秒杀请求包含了很多步骤,例如:

  • 风险控制;
  • 库存锁定;
  • 生成订单;
  • 短信通知;
  • 更新统计数据。

能否决定秒杀成功,实际上只有风险控制和库存锁定这 2 个步骤。只要用户的秒杀请求通过风险控制,并在服务端完成库存锁定,就可以给用户返回秒杀结果了,对于后续的生成订单、短信通知和更新统计数据等步骤,并不一定要在秒杀请求中处理完成。

所以当服务端完成前面 2 个步骤,确定本次请求的秒杀结果后,就可以马上给用户返回响应,然后把请求的数据放入消息队列中,由消息队列异步地进行后续的操作。

file

处理一个秒杀请求,从 5 个步骤减少为 2 个步骤,这样不仅响应速度更快,并且在秒杀期间,我们可以把大量的服务器资源用来处理秒杀请求。秒杀结束后再把资源用于处理后面的步骤,充分利用有限的服务器资源处理更多的秒杀请求。

可以看到,在这个场景中,消息队列被用于实现服务的异步处理。这样做的好处是:

  • 可以更快地返回结果
  • 减少等待,自然实现了步骤之间的并发,提升系统总体的性能。

流量控制

继续说我们的秒杀系统,我们已经使用消息队列实现了部分工作的异步处理,但我们还面临一个问题:如何避免过多的请求压垮我们的秒杀系统?

加入消息队列后,整个秒杀流程变为:
file
网关在收到请求后,将请求放入请求消息队列;
后端服务从请求消息队列中获取 APP 请求,完成后续秒杀处理过程,然后返回结果。

秒杀开始后,当短时间内大量的秒杀请求到达网关时,不会直接冲击到后端的秒杀服务,而是先堆积在消息队列中,后端服务按照自己的最大处理能力,从消息队列中消费请求进行处理。

这种设计的优点是:能根据下游的处理能力自动调节流量,达到“削峰填谷”的作用。

img

但这样做同样是有代价的:

增加了系统调用链环节,导致总体的响应时延变长。
上下游系统都要将同步调用改为异步消息,增加了系统的复杂度。

服务解耦

消息队列的另外一个作用,就是实现系统应用之间的解耦。

订单是电商系统中比较核心的数据,当一个新订单创建时:

支付系统需要发起支付流程;
风控系统需要审核订单的合法性;
客服系统需要给用户发短信告知用户;
经营分析系统需要更新统计数据;
……

这些订单下游的系统都需要实时获得订单数据。随着业务不断发展,这些订单下游系统不断的增加,不断变化

img

引入消息队列后,订单服务在订单变化时发送一条消息到消息队列的一个主题 Order 中,所有下游系统都订阅主题 Order,这样每个下游系统都可以获得一份实时完整的订单数据。

img

无论增加、减少下游系统或是下游系统需求如何变化,订单服务都无需做任何更改,实现了订单服务与下游服务的解耦。

总结

以上就是消息队列最常被使用的三种场景:异步处理、流量控制和服务解耦。当然,消息队列的适用范围不仅仅局限于这些场景.简单的说,我们在单体应用里面需要用队列解决的问题,在分布式系统中大多都可以用消息队列来解决。

消息队列的选型

  • 开源
  • 流行, 活跃
  • 兼容
  • 确保消息的可靠传递
  • 支持集群, 高可用
  • 性能要求

可供选择的消息队列

1. Rabbit MQ (AMQP协议)
  • "messaging that just works"
  • 老牌产品
2. Rocket MQ
  • Alibaba 2012年开源
  • 2017年成为Apache顶级项目
  • 经过多次双十一考验, 稳定可靠
3. Kafka
  • 也是Apache顶级项目
  • 最初设计目的是处理海量日志
    • 不保证消息可靠性(可能丢消息)
    • 不支持集群
  • 当下的Kafka
    • 非常成熟的消息队列产品
    • 稳定可靠, 功能特性可以满足绝大多数场景要求
    • 兼容性好,尤其是大数据和流计算领域优先支持Kafka!!!
    • 异步性能最好(同步性能反而较差).
4. 第二梯队
  1. Active MQ(JMS)
  2. Zero MQ
  3. 雅虎Pulsar

Views: 240

2-Storm集群安装(伪分布式)

基础环境:百度网盘,提取码:NIIT
--来自百度网盘超级会员V4的分享)

安装ZooKeeper集群

Storm 使用 Zookeeper 协调管理集群. Zookeeper 并不是 用于消息传递, 所以 Storm 对Zookeeper造成的负载压力非常低. 单节点Zookeeper集群在大多数情况下应该是足够的,但是如果您想要故障转移或部署大型Storm集群,则可能需要较大的Zookeeper集群. 部署Zookeeper的说明是这里.

这里我们安装一个单机伪分布式的Zookeeper集群:

conf/zoo_1.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk1/data
dataLogDir=/home/hadoop/tmp/zk1/dataLog
clientPort=2181
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

conf/zoo_2.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk2/data
dataLogDir=/home/hadoop/tmp/zk2/dataLog
clientPort=2182
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

conf/zoo_3.cfg

tickTime=2000
initLimit=10
syncLimit=5
dataDir=/home/hadoop/tmp/zk3/data
dataLogDir=/home/hadoop/tmp/zk3/dataLog
clientPort=2183
4lw.commands.whitelist=*
server.1=localhost:2891:3891
server.2=localhost:2892:3892
server.3=localhost:2893:3893

在dataDir中创建文件myid内容是1或2或3

echo 1 > /home/hadoop/tmp/zk1/data/myid
echo 2 > /home/hadoop/tmp/zk2/data/myid
echo 3 > /home/hadoop/tmp/zk3/data/myid

启动服务

bin/zkServer.sh start conf/zoo_1.cfg
bin/zkServer.sh start conf/zoo_2.cfg
bin/zkServer.sh start conf/zoo_3.cfg

关于Zookeeper部署的几点注意事项:

1.在监督(supervision)下运行Zookeeper至关重要,因为Zookeeper是快速失败的,如果遇到任何错误的情况都将退出进程. 有关详细信息,请参阅这里 .

2.建立一个cron定时任务来压缩Zookeeper的数据和事务日志至关重要. Zookeeper守护进程本身不会这样做,如果没有设置cron,Zookeeper将很快耗尽磁盘空间. 有关详细信息,请参阅这里.

Nimbus和worker节点的安装环境

接下来你需要准备Nimbus 和 worker 节点的安装环境:

  1. Java 8+ (Apache Storm 2.x is tested through travisci against a java 8 JDK)
  2. Python 2.6.6 (Python 3.x should work too,but is not tested as part of our CI enviornment)

这些依赖版本是Storm已经测试过的. Storm 在不同的Java 或Python版本上也许会存在问题.

下载解压 Storm

接下来,下载一个Storm版本,并解压zip文件到Nimbus和每个worker机器上的某个目录下. Storm版本可以从这里下载.

当前最新的版本是2.2.0,这里使用2.1.0的版本。

我们需要下载:

也可以使用wget下载,如

wget https://dlcdn.apache.org/storm/apache-storm-2.1.0/apache-storm-2.1.0.tar.gz

在storm.yaml中设置必要的配置

Storm 发布包中在目录conf/storm.yaml 下包含一个默认的配置文件. 你可以在[这里](http://github.com/apache/storm/blob/master /conf/defaults.yaml)查看默认值. storm.yaml 中的存在的配置项会覆盖掉 defaults.yaml中相应的配置项. 下面一些配置是集群运行时所必要的:

  1. storm.zookeeper.servers: 这是一个Storm集群所依赖 Zookeeper 集群的hosts列表. 类似于:
storm.zookeeper.servers:
  - "111.222.333.444"
  - "555.666.777.888"

如果配置的Zookeeper集群不是默认的端口, 你应该设置 storm.zookeeper.port 选项.

  1. storm.local.dir: Nimbus 和 Supervisor 守护进程需要配置一个本地目录来存储少量状态信息(例如jars包,配置文件等等). 您应该在每个机器上创建该目录,给予适当的权限,然后使用此配置填写目录位置. 例如:
storm.local.dir: "/mnt/storm"

如果您在windows下运行Strom,应该如下: yaml storm.local.dir: "C:\\storm-local" 如果您使用相对路径,那么路径是相对于(STORM_HOME). 您也可以使用默认值 $STORM_HOME/storm-local

  1. nimbus.seeds: worker节点需要知道哪些机器是主机的候选者,以便下载 topology jar和confs(nimbus.host 在1.0之后已经废弃,这里实现了HA). 例如:
nimbus.seeds: ["111.222.333.44"]

鼓励您填写机器的FQDN (Fully Qualified Domain Name,全域名)列表. 如果要设置Nimbus HA,则必须解决运行nimbus的所有机器的FQDN.当您只想设置“伪分布式”集群时您可能希望将其保留为默认值,仍然鼓励您填写FQDN.

  1. supervisor.slots.ports: 对于每个worker节点,您可以使用此配置设置在该计算机上运行的worker数量. 每个worker使用单个端口接收消息,并且此设置定义哪些端口打开以供使用. 如果您在此定义五个端口,那么Storm将分配最多五个worker在本机上运行. 如果您定义了三个端口,Storm将只能运行三个worker. 默认情况下,此设置被配置为在端口6700,6701,6702和6703上运行4个worker:
supervisor.slots.ports:
    - 6700
    - 6701
    - 6702
    - 6703 

以上是一些配置的介绍,下面我们的具体配置如下:

 storm.zookeeper.servers:
    - "hadoop000"
 storm.local.dir: "/home/hadoop/app/storm-2.1.0/data"
 nimbus.seeds: ["hadoop000"]
 storm.zookeeper.port: 2181
 supervisor.slots.ports:
    - 6700
    - 6701
    - 6702
    - 6703
 ui.port: 8082

注意:配置storm.zookeeper.servers前面有空格,等等细节很严格。另外ui服务进程的端口8080由于很容易和其他服务冲突,这里改成了8082

启动Storm进程

  • 主节点启动nimbus服务

    启动Nimbus。Run the command bin/storm nimbus under supervision on the master machine.

    storm nimbus &
  • 主节点启动UI服务:

    Run the Storm UI (a site you can access from the browser that gives diagnostics on the cluster and topologies) by running the command “bin/storm ui” under supervision. The UI can be accessed by navigating your web browser to http://{ui host}:8080. by default

    但是由于之前在配置文件中已经将UI端口修改成8082(避免和tomcat的8080冲突), 因此现在的访问方式是http://{ui host}:8082

    bin/storm ui &
  • 在从节点(工作节点)上启动supervisor服务(只不过目前是配置的伪分布式,所有服务都是在一个节点上,没有备用节点):

    storm supervisor &
  • 用jps判断是否启动成功(如果失败,则检查日志)

    img

  • 避免打印信息或者报错信息出现在屏幕上

    为了避免打印信息或者报错信息出现的屏幕上,也可以这样

    bin/storm ui  >/dev/null 2>&1 &  

    其中,2>&1 是将标准出错重定向到标准输出,但是最好的方式是采用nohup方式运行。

    nohup命令:如果你正在运行一个进程,而且你觉得在退出帐户时该进程还不会结束,那么可以使用nohup命令。该命令可以在你退出帐户/关闭终端之后继续运行相应的进程。nohup就是不挂起的意思( no hang up), 像这样:

    nohup ./bin/storm nimbus  > /dev/null  2>&1 & 
    nohup ./bin/storm supervisor > /dev/null 2>&1 & 
    nohup ./bin/storm ui > /dev/null  2>&1 &

查看UI界面

访问http://hadoop000:8082/ (端口是在配置文件中配置过的)

img

img

img

Numbus的配置超级繁杂, 有14页之多

查看log

log日志在storm包的logs文件夹中

如果想要在UI界面点击主机名或者端口号查看日志,则对应的机器需要开启logviewer服务

storm logviewer &

配置集群

假设当前机器hadoop001作为主节点, 如果要增加从节点机器hadoop002和hadoop003`,只需要修改如下:

storm.zookeeper.servers:
     - "hadoop001"
     - "hadoop002"
     - "hadoop003"
 storm.local.dir: "/home/hadoop/app/storm-2.1.0/data"
 nimbus.seeds: ["hadoop001"]
 storm.zookeeper.port: 2181
 supervisor.slots.ports:
   - 6700
   - 6701
   - 6702
   - 6703
 ui.port: 8082  

给bin目录的文件添加可执行权限:

bin>   chmod u+x *

开启nimbus和ui后台进程,把整个安装文件夹使用scp命令复制到其他两台机器(需要提前配置好通信机器之间的ssh免密访问),清空logsdata文件夹里面的内容,给bin目录的文件添加可执行权限,开启supervisor进程,这样就准备好集群环境了。

打开UI界面查看 可以看到Nimbus是hadoop001, 而supervisor是hadoop002和hadoop003

Storm UI

Cluster Summary

Version Supervisors Used slots Free slots Total slots Executors Tasks
2.1.0 2 0 8 8 0 0

Nimbus Summary

Search:

Host Port Status Version Uptime
hadoop001 6627 Leader 2.1.0 9m 16s

Showing 1 to 1 of 1 entries

Owner Summary

Search:

Owner Total Topologies Total Executors Total Workers Memory Usage (MB)
No data available in table

Showing 0 to 0 of 0 entries

Topology Summary

Search:

Name Owner Status Uptime Num workers Num executors Num tasks Replication count Assigned Mem (MB) Scheduler Info Topology Version Storm Version
No data available in table

Showing 0 to 0 of 0 entries

Supervisor Summary

Search:

Host Id Uptime Slots Used slots Avail slots Used Mem (MB) Version
hadoop002 (log) d8412147-0dbb-4933-a2d3-1bb28c806dc2-192.168.186.102 10m 33s 4 0 4 0 2.1.0
hadoop003 (log) 57798559-c338-4844-9e29-eb9aa99fa2fe-192.168.186.103 1m 7s 4 0 4 0 2.1.0

Showing 1 to 2 of 2 entries

Nimbus Configuration

Show entries

Search:

Key Value
blacklist.scheduler.reporter "org.apache.storm.scheduler.blacklist.reporters.LogReporter"
blacklist.scheduler.resume.time.secs 1800
blacklist.scheduler.strategy "org.apache.storm.scheduler.blacklist.strategies.DefaultBlacklistStrategy"
blacklist.scheduler.tolerance.count 3
blacklist.scheduler.tolerance.time.secs 300
client.blobstore.class "org.apache.storm.blobstore.NimbusBlobStore"
dev.zookeeper.path "/tmp/dev-storm-zookeeper"
drpc.authorizer.acl.filename "drpc-auth-acl.yaml"
drpc.authorizer.acl.strict false
drpc.childopts "-Xmx768m"
drpc.disable.http.binding true
drpc.http.creds.plugin "org.apache.storm.security.auth.DefaultHttpCredentialsPlugin"
drpc.http.port 3774
drpc.https.keystore.password ""
drpc.https.keystore.type "JKS"
drpc.https.port -1
drpc.invocations.port 3773
drpc.invocations.threads 64
drpc.max_buffer_size 1048576
drpc.port 3772

Showing 1 to 20 of 266 entries

为nimbus配置高可用

最后, 尽管前面我们使用了三台机器搭建了集群,由于nimbus服务没有配置高可用, 实际上还是算伪分布式。为nimbus配置高可用也非常简单, 比如说让hadoop001hadoop002其中一个作为主节点,另一个作为备用主节点,则配置如下:

nimbus.seeds: ["hadoop001", "hadoop002"]

Nimbus Summary

Search:

Host Port Status Version Uptime
hadoop001 6627 Leader 2.1.0 9m 16s
hadoop002 6627 Follower 2.1.0 ...

作业

使用Centos7系统搭建Storm集群, 将搭建过程形成笔记,并附上截图。

练习

在互联网上查询Twitter是如何使用storm进行情感分析的,了解storm的重要性以及用途。

Views: 895