Storm 集群搭建和自定义调度器

编写项目并打包上传

为了方便, 这里我们把自定义调度器DirectScheduler和测试用拓扑程序DirectScheduledTopology放在一个项目中.

pom.xml

<project xmlns="http://maven.apache.org/POM/4.0.0" xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>com.niit.storm.example</groupId>
    <artifactId>storm_example</artifactId>
    <version>0.0.1-SNAPSHOT</version>
    <packaging>jar</packaging>

    <name>storm_example</name>
    <url>http://maven.apache.org</url>

    <properties>
        <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
    </properties>

    <repositories>
        <repository>
            <id>clojars.org</id>
            <url>http://clojars.org/repo</url>
        </repository>
    </repositories>

    <dependencies>
        <dependency>
            <groupId>junit</groupId>
            <artifactId>junit</artifactId>
            <version>3.8.1</version>
            <scope>test</scope>
        </dependency>
        <dependency>
            <groupId>org.apache.storm</groupId>
            <artifactId>storm-core</artifactId>
            <version>2.1.0</version>
            <scope>provided</scope>
        </dependency>
    </dependencies>
    <build>
        <plugins>
            <plugin>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>2.2-beta-5</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies
                        </descriptorRef>
                    </descriptorRefs>
                    <archive>
                        <manifest>
                            <mainClass/>
                        </manifest>
                    </archive>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-compiler-plugin</artifactId>
                <version>3.8.1</version>
                <configuration>
                    <source>8</source>
                    <target>8</target>
                </configuration>
            </plugin>
        </plugins>
    </build>
</project>

自定义Scheduler

package storm.scheduler;

import org.apache.storm.scheduler.*;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.*;

/**
 * @Author: deLucia
 * @Date: 2020/7/24
 * @Version: 1.0
 * @Description: 自定义调度器
 * DirectScheduler把划分单位缩小到组件级别,
 * 1个Spout和1个Bolt可以分别指定到某个指定的节点上运行,
 * 如果没有指定,还是按照系统自带的调度器进行调度.
 * 这个配置在Topology提交的Conf配置中可配.
 */

public class DirectScheduler implements IScheduler {
    protected static final Logger logger = LoggerFactory.getLogger(DirectScheduler.class);

    @Override
    public void prepare(Map conf) {
    }

    @Override
    public void schedule(Topologies topologies, Cluster cluster) {
        // Gets the topology which we want to schedule
        String assignedFlag; //作业是否要指定分配的标识
        Collection<TopologyDetails> topologyDetailes = topologies.getTopologies();

        for (TopologyDetails td : topologyDetailes) {
            Map<String, Object> map = td.getConf();
            assignedFlag = (String) map.get("assigned_flag"); // 标志位
            // 如果找到的拓扑逻辑的assigned_flag标记为1则代表是要指定分配到目标节点的的,否则走系统的默认调度
            if ("1".equals(assignedFlag)) {
                if (!cluster.needsScheduling(td)) { // 如果资源不紧张, 走默认调度
                    logger.info("use default - scheduler");
                    new DefaultScheduler().schedule(topologies, cluster);
                } else {
                    topologyAssign(cluster, td, map); // 否则走自定义调度
                }
            } else {
                logger.info("use default scheduler");
                new DefaultScheduler().schedule(topologies, cluster);
            }
        }
    }

    @Override
    public Map<String, Object> config() {
        return new HashMap<String, Object>();
    }

    @Override
    public void cleanup() {

    }

    /**
     * 拓扑逻辑的调度
     *
     * @param cluster  集群
     * @param topology 具体要调度的拓扑逻辑
     * @param map      map配置项
     */
    private void topologyAssign(Cluster cluster, TopologyDetails topology, Map<String, Object> map) {
        logger.info("use custom scheduler");
        if (topology == null) {
            logger.warn("topology is null");
            return;
        }

        Map<String, Object> designMap = (Map<String, Object>) map.get("design_map");

        if (designMap == null) {
            logger.warn("found no design_map in config");
            return;
        }

        // find out all the needs-scheduling components of this topology
        Map<String, List<ExecutorDetails>> componentToExecutors = cluster.getNeedsSchedulingComponentToExecutors(topology);

        // 当没有待分配线程时直接退出
        if (componentToExecutors == null || componentToExecutors.size() == 0) {
            return;
        }

        // 如果有指定待分配的节点名称则分配到指定节点
        Set<String> keys = designMap.keySet();
        for (String componentName : keys) {
            String nodeName = (String) designMap.get(componentName);
            componentAssign(cluster, topology, componentToExecutors, componentName, nodeName);
        }

    }

    /**
     * 组件调度
     *
     * @param cluster        集群的信息
     * @param topology       待调度的拓扑细节信息
     * @param totalExecutors 组件的执行器
     * @param componentName  组件的名称
     * @param supervisorName 节点的名称
     */
    private void componentAssign(Cluster cluster, TopologyDetails topology, Map<String, List<ExecutorDetails>> totalExecutors, String componentName, String supervisorName) {

        List<ExecutorDetails> executors = totalExecutors.get(componentName);

        // 由于Scheduler是轮询调用, 这里需要提前判空
        if (executors == null) {
            return;
        }

        // find out the our "special-supervisor" from the supervisor metadata
        Collection<SupervisorDetails> supervisors = cluster.getSupervisors().values();
        SupervisorDetails specialSupervisor = null;
        for (SupervisorDetails supervisor : supervisors) {
            Map meta = (Map) supervisor.getSchedulerMeta();
            if (meta != null && meta.get("name") != null) {
                if (meta.get("name").equals(supervisorName)) {
                    specialSupervisor = supervisor;
                    break;
                }
            }
        }
        // found the special supervisor
        if (specialSupervisor != null) {
            logger.info("supervisor name:" + specialSupervisor);
            List<WorkerSlot> availableSlots = cluster.getAvailableSlots(specialSupervisor);
            // 如果目标节点上已经没有空闲的slot,则进行强制释放
            if (availableSlots.isEmpty() && !executors.isEmpty()) {
                for (Integer port : cluster.getUsedPorts(specialSupervisor)) {
                    logger.info("no available work slots");
                    logger.info("will free one slot from " + specialSupervisor.getId() + ":" + port);
                    cluster.freeSlot(new WorkerSlot(specialSupervisor.getId(), port));
                }
            }
            // 重新获取可用的slot
            availableSlots = cluster.getAvailableSlots(specialSupervisor);
            // 选取节点上第一个slot,进行分配
            WorkerSlot firstSlot = availableSlots.get(0);

            for (WorkerSlot slot : availableSlots) {
                if (slot != null) {
                    logger.info("assigned executors:" + executors + " to slot: [" + firstSlot.getNodeId() + ", " + firstSlot.getPort() + "]");
                    cluster.assign(firstSlot, topology.getId(), executors);
                    return;
                }
            }

            logger.warn("failed to assigned executors:" + executors + " to slot: [" + firstSlot.getNodeId() + ", " + firstSlot.getPort() + "]");

        } else {
            logger.warn("supervisor name not existed in supervisor list!");
        }
    }
}

测试用拓扑

package storm.scheduler;

import org.apache.storm.Config;
import org.apache.storm.StormSubmitter;
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.topology.base.BaseRichSpout;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;
import java.util.HashMap;
import java.util.Map;

/**
 * 使用Storm实现积累求和的操作 - 使用自定义调度器
 * 自定义DiectScheduler可以为某个组件可以去指定不同节点运行, 开发者可以根据实际集群的资源情况选择适合的节点, 工作中有时候是很有用的.
 * 我们的目的是调度spout到supervisor002节点   调度bolt到supervisor003节点
 * 首先上传拓扑
 * bin/storm jar storm_example.jar storm.scheduler.SumTopology  sum supervisor002 supervisor003
 */
public class DirectScheduledTopology {

    /**
     * Spout需要继承BaseRichSpout
     * 数据源需要产生数据并发射
     */
    public static class DataSourceSpout extends BaseRichSpout {

        private SpoutOutputCollector collector;

        /**
         * 初始化方法,只会被调用一次
         *
         * @param conf      配置参数
         * @param context   上下文
         * @param collector 数据发射器
         */
        @Override
        public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) {
            this.collector = collector;
        }

        int number = 0;

        /**
         * 会产生数据,在生产上肯定是从消息队列中获取数据
         * 这个方法是一个死循环,会一直不停的执行
         */
        @Override
        public void nextTuple() {
            this.collector.emit(new Values(++number));

            System.out.println("Spout: " + number);

            // 防止数据产生太快
            Utils.sleep(1000);

        }

        /**
         * 声明输出字段
         * @param declarer
         */
        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {
            declarer.declare(new Fields("num"));
        }
    }

    /**
     * 数据的累积求和Bolt:接收数据并处理
     */
    public static class SumBolt extends BaseRichBolt {
        private OutputCollector out;

        /**
         * 初始化方法,会被执行一次
         *
         * @param stormConf
         * @param context
         * @param collector
         */
        @Override
        public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
            out = collector;
        }

        int sum = 0;

        /**
         * 其实是一个死循环调用,职责:获取Spout发送过来的数据
         * @param input
         */
        @Override
        public void execute(Tuple input) {

            // Bolt中获取值可以根据index获取
            // 也可以根据上一个环节中定义的field的名称获取(建议)
            Integer value = input.getIntegerByField("num");
            sum += value;
            System.out.println("sum: " + sum);
            out.emit(new Values(sum));
        }

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

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

        String topologyName = "sum-topology";
        String spoutSupervisorName = "supervisor002";
        String boltSupervisorName = "supervisor003";

        if (args != null && args.length > 2) {
            topologyName = args[0];
            spoutSupervisorName = args[1];
            boltSupervisorName = args[2];
        } else {
            System.out.println("Usage:");
            System.out.println("storm jar path-to-jar-file main-class topology-name spout-supervisor-name bolt supervisor-name");
            System.exit(1);
        }

        String spoutId = "num-spout";
        String boltId = "sum-bolt";

        Config config = new Config();
        config.setDebug(true);
        config.setNumWorkers(2);

        Map<String, String> component2Node = new HashMap<>();
        component2Node.put(spoutId, spoutSupervisorName);
        component2Node.put(boltId, boltSupervisorName);
        // 此标识代表topology需要被 
        // 我们的自定义调度器 DirectScheduler  调度
        config.put("assigned_flag", "1");
        // 具体的组件节点对信息
        config.put("design_map", component2Node);

        // TopologyBuilder根据Spout和Bolt来构建出Topology
        // Storm中任何一个作业都是通过Topology的方式进行提交的
        // Topology中需要指定Spout和Bolt的执行顺序
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout("num-spout", new DataSourceSpout(), 2);
        builder.setBolt("sum-bolt", new SumBolt(), 2).shuffleGrouping("num-spout");

        // 创建一个本地Storm集群:本地模式运行,不需要搭建Storm集群
        // new LocalCluster().submitTopology("sum-topology", config, builder.createTopology());

        StormSubmitter.submitTopology(topologyName, config, builder.createTopology());
    }

}

打成jar包, 上传至nimbus所在节点的storm下的lib目录下

集群环境搭建(3节点)

准备三台机器,hadoop001, hadoop002, hadoop003.

我们希望hadoop001启动: nimbus, ui, logviewer, supervisor

hadoop002,hadoop002启动: supervisor, logviewer

各节点的分工如下:

hostname nimbus supervisor zookeeper
hadoop001 √(supervisor001) server.1
hadoop002 √(supervisor002) server.2
hadoop003 √(supervisor003) server.3

虚拟机准备

可以先配置一台机器, 假设主机名为hadoop001, 静态IP为192.168.186.100. 接下来分别安装jdk环境, 解压zookeeper和storm的安装文件到指定目录~/app下面, 配置环境变量如下:

~/.bash_profile

...

# JAVA ENV
export JAVA_HOME=/home/hadoop/app/jdk1.8.0_211
PATH=$JAVA_HOME/bin:$PATH

# ZOOKEEPER HOME
export ZOOKEEPER_HOME=/home/hadoop/app/zookeeper-3.4.11
PATH=$ZOOKEEPER_HOME/bin:$PATH

# STORM HOME
export STORM_HOME=/home/hadoop/app/storm-2.1.0
PATH=$STORM_HOME/bin:$PATH

...

之后再克隆出另外两台虚拟机, 修改IP分别为192.168.186.102和192.168.186.103.修改主机名为hadoop002和hadoop003.

克隆方式如下:

image-20211014104419615

三台机器/etc/hosts文件修改如下:

192.168.186.101 hadoop001
192.168.186.102 hadoop002
192.168.186.103 hadoop003

配置SSH免密

为了便于集群中个节点互访, 每台机器都需要通过ssh-keygen命令后输入三次回车生成一对密钥(公钥和私钥), 可以在~/.ssh下查看, 并将自己的公钥通过ssh-copy-id命令分发给其他机器的~/.ssh/authorized_keys(必须是600权限)文件中:

[hadoop@hadoop001 .ssh]$ ssh-keygen -t rsa
[hadoop@hadoop002 .ssh]$ ssh-keygen -t rsa
[hadoop@hadoop003 .ssh]$ ssh-keygen -t rsa

[hadoop@hadoop001 .ssh]$ ssh-copy-id hadoop@hadoop002
[hadoop@hadoop001 .ssh]$ ssh-copy-id hadoop@hadoop003

[hadoop@hadoop002 .ssh]$ ssh-copy-id hadoop@hadoop001
[hadoop@hadoop002 .ssh]$ ssh-copy-id hadoop@hadoop003

[hadoop@hadoop003 .ssh]$ ssh-copy-id hadoop@hadoop001
[hadoop@hadoop003 .ssh]$ ssh-copy-id hadoop@hadoop002

测试: (以从hadoop001连接hadoop002为例)


[hadoop@hadoop001 .ssh]$ ssh hadoop002
Last login: Sun Oct 17 01:15:17 2021 from hadoop001
[hadoop@hadoop002 ~]$ exit
登出
Connection to hadoop002 closed.

搭建zookeeper集群

跟单机基本一致, 配置文件相同

# The number of milliseconds of each tick
tickTime=2000
# The number of ticks that the initial
# synchronization phase can take
initLimit=10
# The number of ticks that can pass between
# sending a request and getting an acknowledgement
syncLimit=5
# the directory where the snapshot is stored.
# do not use /tmp for storage, /tmp here is just
# example sakes.
dataDir=/home/hadoop/app/zookeeper-3.4.11/dataDir
dataLogDir=/home/hadoop/app/zookeeper-3.4.11/dataLogDir
# the port at which the clients will connect
clientPort=2181
# the maximum number of client connections.
# increase this if you need to handle more clients
#maxClientCnxns=60
#
# Be sure to read the maintenance section of the
# administrator guide before turning on autopurge.
#
# http://zookeeper.apache.org/doc/current/zookeeperAdmin.html#sc_maintenance
#
# The number of snapshots to retain in dataDir
#autopurge.snapRetainCount=3
# Purge task interval in hours
# Set to "0" to disable auto purge feature
#autopurge.purgeInterval=1

# 将所有命令添加到白名单中
4lw.commands.whitelist=*

# 配置zookeeper服务器的 [主机名或主机ID]:[同步端口(唯一)]:[选举端口(唯一)]
server.1=hadoop001:2888:3888
server.2=hadoop002:2888:3888
server.3=hadoop003:2888:3888

但是$dataDir下的myid文件保存的数字应该和对应server的id一致

然后分别启动三个机器的zookeeper就可以了

使用status查看应该有2个follower和1个leader

搭建Storm集群

除了storm.zookeeper.serversnimbus.seeds的配置不同之外, 和单机的storm配置基本相同.

hadoop001的配置:

 "pology.eventlogger.executors": 1
 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

 # 自定义调度器, 需要提前将jar包放在nimbus节点storm安装目录的lib目录里
 storm.scheduler: "storm.scheduler.DirectScheduler"
 # 自定义属性, 配合自定义调度器使用
 supervisor.scheduler.meta:
   name: "supervisor001"

hadoop002的配置:


 "pology.eventlogger.executors": 1
 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

 supervisor.scheduler.meta:
   name: "supervisor002"

hadoop003的配置:


 "pology.eventlogger.executors": 1
 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

 supervisor.scheduler.meta:
   name: "supervisor003"

然后分别启动以下进程

[hadoop@hadoop001 storm-2.1.0]$ jps
2964 Jps
2310 Nimbus
2871 LogviewerServer
2186 QuorumPeerMain
2555 UIServer
2412 Supervisor

[hadoop@hadoop002 storm-2.1.0]$ jps
2309 LogviewerServer
2217 Supervisor
2558 Jps
2095 QuorumPeerMain

[hadoop@hadoop003 storm-2.1.0]$ jps
2165 Supervisor
2234 LogviewerServer
2333 Jps
2095 QuorumPeerMain

集群启停脚本编写

start-zookeeper-cluster.sh

#/bin/bash

echo "========== zk start =========="
for node in hadoop001 hadoop002 hadoop003
do
  ssh $node "source $HOME/.bash_profile;zkServer.sh start $ZOOKEEPER_HOME/conf/zoo.cfg"
done

echo "wait 5 secs .... ...."
sleep 5

echo "========== zk status =========="
for node in hadoop001 hadoop002 hadoop003
do
  ssh $node "source $HOME/.bash_profile;zkServer.sh status $ZOOKEEPER_HOME/conf/zoo.cfg"
done

echo "wait 5 secs .... ...."
sleep 5

start-storm-cluster.sh

#/bin/bash

echo "========== storm start =========="
ssh hadoop001 "source $HOME/.bash_profile;nohup storm nimbus > /dev/null  2>&1 &"
ssh hadoop001 "source $HOME/.bash_profile;nohup storm ui > /dev/null  2>&1 &"
ssh hadoop001 "source $HOME/.bash_profile;nohup storm supervisor > /dev/null  2>&1 &"
ssh hadoop001 "source $HOME/.bash_profile;nohup storm logviewer > /dev/null  2>&1 &"

ssh hadoop002 "source $HOME/.bash_profile;nohup storm supervisor > /dev/null  2>&1 &"
ssh hadoop002 "source $HOME/.bash_profile;nohup storm logviewer > /dev/null  2>&1 &"

ssh hadoop003 "source $HOME/.bash_profile;nohup storm supervisor > /dev/null  2>&1 &"
ssh hadoop003 "source $HOME/.bash_profile;nohup storm logviewer > /dev/null  2>&1 &"

需要保证三台机器都设置正确的环境变量, 如果必要可以将hadoop001的环境变量通过scp命令覆盖到其他机器上去如:

[hadoop@hadoop001]$ scp ~/.bash_profile hadoop@hadoop002:~/.bash_profile
[hadoop@hadoop001]$ scp ~/.bash_profile hadoop@hadoop003:~/.bash_profile

下面这个xRunCommand.sh脚本用来在集群中执行相同的命令并返回结果

#/bin/bash

if (( $# == 0 ));then
  echo "Usage /path/to/xRunCommand.sh \"<COMMAND>\""
  exit 0
fi

for node in hadoop001 hadoop002 hadoop003
do
  echo "======== $node ========"
  echo "ssh $node $1"
  ssh $node "source ~/.bash_profile;$1"
done

比如可以xRunCommand.sh jps来查看集群中三台机器中的所有java进程

[hadoop@hadoop001 scripts]$ ./xRunCommand.sh jps
======== hadoop001 ========
ssh hadoop001 jps
8656 Jps
7220 QuorumPeerMain
7976 UIServer
8026 LogviewerServer
7933 Nimbus
======== hadoop002 ========
ssh hadoop002 jps
6119 Supervisor
7031 Jps
5753 QuorumPeerMain
6235 Worker
6223 LogWriter
======== hadoop003 ========
ssh hadoop003 jps
6896 Jps
6002 Supervisor
6119 Worker
5641 QuorumPeerMain
6107 LogWriter

提交拓扑到集群

在hadoop001及nimbus所在节点,使用以下命令提交拓扑

storm jar $STORM_HOME/lib/storm_example-0.0.1-SNAPSHOT.jar storm.scheduler.DirectScheduledTopology sum-topo supervisor002 supervisor003

其中:

  1. 第一个参数sum-topo 为拓扑名称

  2. 第二个参数supervisor002和第三个参数supervisor003表示调度的目的地

    即调度spout组件到supervisor002节点,调度bolt组件到supervisor003节点

    具体逻辑请查看DirectScheduledTopology类中的源代码

提交后, 即可在UI界面查看结果:

image-20211017023208132

可见我们的自定义调度器确实起作用了.

也可以通过$STORM_HOME/logs/nimbus.log查看到DirectScheduler类输出的相关日志, 例如:

2021-10-17 01:20:40.659 s.s.DirectScheduler timer [INFO] use custom scheduler
2021-10-17 01:20:40.659 s.s.DirectScheduler timer [INFO] supervisor name:SupervisorDetails ID: a0e622ed-0193-4acc-a85f-6f61aa55d2f4-192.168.186.102 HOST: hadoop002 META: null SCHED_META: {name=supervisor002} PORTS: [6700, 6701, 6702, 6703]
2021-10-17 01:20:40.660 s.s.DirectScheduler timer [INFO] assigned executors:[[4, 4], [3, 3]] to slot: [a0e622ed-0193-4acc-a85f-6f61aa55d2f4-192.168.186.102, 6700]
2021-10-17 01:20:40.673 s.s.DirectScheduler timer [INFO] supervisor name:SupervisorDetails ID: 76e19329-990f-441e-8633-623019800d73-192.168.186.103 HOST: hadoop003 META: null SCHED_META: {name=supervisor003} PORTS: [6700, 6701, 6702, 6703]
2021-10-17 01:20:40.673 s.s.DirectScheduler timer [INFO] assigned executors:[[6, 6], [5, 5]] to slot: [76e19329-990f-441e-8633-623019800d73-192.168.186.103, 6700]

Nimbus失败异常处理

问题描述

启动Storm集群后发现错误如下:

org.apache.storm.utils.NimbusLeaderNotFoundException: Could not find leader nimbus from seed hosts ["hadoop001"]. Did you specify a valid list of nimbus hosts for config nimbus.seeds?

问题原因

可能是Zookeeper所保存的数据和现在集群的状态不一致, 也可能是nimbus进程异常退出(具体可在nimbus节点查看nimbus日志$STORM_HOME/logs/nimbus.log进行错误排查)

解决办法

不管是何种原因, 都可以使用下面的方法简单粗暴的解决:

注意, 以下处理方法会导致Storm集群中的所有运行的拓扑任务丢失.

  1. 先停止storm集群再停止zookeeper集群

    # 停止storm集群只要使用kill -9 杀死相关进程即可
    kill -9 ... ... ... ...
    ~/scripts/stop-zookeeper-cluster.sh
  2. 删除所有zookeeper下dataDir除了myid之外的所有文件.

    [hadoop@hadoop001]$ rm -rf ~/app/zookeeper-3.4.11/dataDir/version-2/
    [hadoop@hadoop002]$ rm -rf ~/app/zookeeper-3.4.11/dataDir/version-2/
    [hadoop@hadoop003]$ rm -rf ~/app/zookeeper-3.4.11/dataDir/version-2/
  3. 启动zookeeper集群

    ~/scripts/start-zookeeper-cluster.sh
  4. 在任意zookeeper集群的节点使用以下命令删除Storm保存在zookeeper中的数据

    $ zkCli.sh
    > rmr /storm
  5. 最后启动storm集群即可.

    $ ~/scripts/start-storm-cluster.sh

谢谢!

Views: 322

Storm计算网站PV和UV(实现可靠处理)

需求分析

编写Storm拓扑实现可靠计算网站当日PV和UV

重点:

  1. 去重计算模式
  2. 实现可靠处理

电商常用指标之PV、UV、VV、独立IP

  • PV(访问量):Page View, 即页面浏览量或点击量,用户每次访问即被计算一次。
  • UV(独立访客):Unique Visitor, 访问您网站的一台电脑客户端为一个访客。00:00-24:00内相同的客户端只会被计算一次。
  • VV即Visit View,访客访问的次数,用以记录所有访客一天内访问量多少次网站。
  • IP(独立IP):指独立IP数。00:00-24:00内相同IP地址之被计算一次。

如果我们要统计一个网站当日有多少实际用户访问,使用UV也就是独立用户作为统计量有什么好处?它比独立IP的统计更加准确吗?

比如你早上8点使用PC访问了某商城的两个页面,下午2点又使用手机访问了这个商城网站的3个页面, 那么对于这个商城网站, 当日的PV、UV、VV、IP各项指标该如何计算呢?

PV为5, PV指浏览量,因此PV指等于上午浏览的2个页面和下午浏览的3个页面之和;

UV为1, UV指独立访客数,因此一天内同一访客的多次访问只计为1个UV;

VV为2, VV指访客的访问次数,上午和下午分别有一次访问行为,因此VV为2

IP为2, IP为独立IP数,由于PC和手机访问时的IP不同,因此独立IP数为2.

可见IP是一个反映网络虚拟地址对象的概念,UV是一个反映实际使用者的概念,每个UV相对于每个IP更加准确地对应一个实际的浏览者。

方案设计

综上所述:使用UV作为统计量,可以更加准确的了解单位时间内实际上有多少个访问者来到了相应的页面。

用Cookie携带的SessionID分析UV值:当客户端第一次访问某个网站服务器的时候,网站服务器会给这个客户端的电脑发出一个Cookie,通常放在这个客户端电脑的C盘当中。在这个Cookie中会分配一个独一无二的编号,这其中会记录一些访问服务器的信息,如访问时间,访问了哪些页面等等。当你下次再访问这个服务器的时候,服务器就可以直接从你的电脑中找到上一次放进去的Cookie文件,并且对其进行一些更新,但那个独一无二的编号是不会变的。

如果把sessionID 放入Set集合实现自动去重,再通过Set.size() 获得UV。该方案在单线程和单JVM下没有问题。高并发情况的多线程情况不适用。

推荐的方案:编写Storm拓扑, Spout获取外部服务器日志然后再通过shuffleGrouping发给多个下级bolt进行数据清洗, 清洗后的数据通过fieldGrouping 进行多线程局部汇总得到,下级blot进行单线程保存sessionID和count数到Map,下一级blot3进行Map遍历,可以得到:Pv、UV、访问深度(每个session_id 的浏览数)

file

创建IDEA项目(基于Maven)

pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>cn.delucia</groupId>
    <artifactId>uv-with-reliability</artifactId>
    <version>1.0-SNAPSHOT</version>

    <dependencies>
        <dependency>
            <groupId>org.apache.storm</groupId>
            <artifactId>storm-core</artifactId>
            <version>2.1.0</version>
            <scope>compile</scope>
        </dependency>
    </dependencies>

</project>

项目结构

file

代码编写

创建工具类cn.delucia.storm.utils.SocketClientUtil

package cn.delucia.storm.utils;

import java.io.OutputStream;
import java.lang.management.ManagementFactory;
import java.net.InetAddress;
import java.net.Socket;
import java.net.UnknownHostException;

/**
 * @Author: deLucia
 * @Date: 2021/9/29
 * @Version: 1.0
 * @Description: 工具类
 */
public class SocketClientUtil {
    /**
     * 是否启用send方法向Socket Server发送信息
     */
    public static boolean ENABLE = true;
    public static final String HOST = "localhost";
    public static final int PORT = 9876;

    private SocketClientUtil() {
    }

    private static String getHostname() {
        try {
            return InetAddress.getLocalHost().getHostName();
        } catch (UnknownHostException e) {
            e.printStackTrace();
        }
        return null;
    }

    /**
     * 返回进程pid
     */
    private static String getPID() {
        String info = ManagementFactory.getRuntimeMXBean().getName();
        return info.split("@")[0];
    }

    /**
     * 线程信息
     */
    private static String getTID() {
        return Thread.currentThread().getName();
    }

    //对象信息
    private static String getOID(Object obj) {
        String cname = obj.getClass().getSimpleName();
        int hash = obj.hashCode();
        return cname + "@" + hash;
    }

    public static String info(Object obj, String msg) {
        return getHostname() + ", " + getPID() + ", " + getTID() + ", " + getOID(obj) + ", " + msg;
    }

    /**
     * 向远端发送sock消息
     * 远端需要先开启nc进行监听:
     * windows: nc -l localhost -L -p 9876
     * linux: nc -l localhost -k -p 9876
     */
    public static void send(Object obj, String msg) {
        if (!ENABLE) {
            return;
        }
        String info = info(obj, msg);
        try (Socket sock = new Socket(HOST, PORT); OutputStream os = sock.getOutputStream()) {
            os.write((info + "\r\n").getBytes());
            os.flush();
        } catch (Exception e) {
            e.printStackTrace();
        }
    }
}

创建类cn.delucia.storm.log.LogTopology

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import com.google.common.collect.ImmutableList;
import org.apache.storm.Config;
import org.apache.storm.LocalCluster;
import org.apache.storm.StormSubmitter;
import org.apache.storm.topology.TopologyBuilder;
import org.apache.storm.tuple.Fields;

/**
 * @author Delucia
 */
public class LogTopology {

    public static final String HOST = "hadoop000";
    public static final Integer ZK_SERVER_PORT = 2181;
    public static final String SPOUT1_ID = LogMockSpout.class.getSimpleName();
    public static final String BOLT1_ID = LogCheckBolt.class.getSimpleName();
    public static final String BOLT2_ID = LogCountBolt.class.getSimpleName();
    public static final String BOLT3_ID = LogGlobalSumBolt.class.getSimpleName();

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

        // 构造拓扑 DAG(有向无环图)
        TopologyBuilder builder = new TopologyBuilder();
        builder.setSpout(SPOUT1_ID, new LogMockSpout(), 1);
        builder.setBolt(BOLT1_ID, new LogCheckBolt(), 4).shuffleGrouping(SPOUT1_ID);
        builder.setBolt(BOLT2_ID, new LogCountBolt(), 4).fieldsGrouping(BOLT1_ID,
                new Fields("date", "sid"));
        // 单线程全局汇总
        builder.setBolt(BOLT3_ID, new LogGlobalSumBolt(), 1).globalGrouping(BOLT2_ID);
        // Storm相关配置
        Config conf = new Config();
        // TOPOLOGY_DEBUG: ON 日志会详细记录所有发射的数据
        conf.setDebug(false);
        // Executor number for event loggers
        conf.setNumEventLoggers(0);
        // Executor number for ackers
        conf.setNumAckers(1);
        // number of workers
        conf.setNumWorkers(1);
        // transfer queue between workers
        conf.put(Config.TOPOLOGY_TRANSFER_BUFFER_SIZE, 64);
        // reveicer queue for each executor
        conf.put(Config.TOPOLOGY_EXECUTOR_RECEIVE_BUFFER_SIZE, 16384);

        if (args.length > 0 ) {
            SocketClientUtil.ENABLE = false;
            conf.put(Config.STORM_ZOOKEEPER_SERVERS, ImmutableList.of(HOST));
            conf.put(Config.STORM_ZOOKEEPER_PORT, ZK_SERVER_PORT);
            conf.put(Config.NIMBUS_SEEDS, ImmutableList.of(HOST));
            conf.put(Config.NIMBUS_THRIFT_PORT, 16627);
            String jarLocation = ".\\target\\parallelisim-1.0-SNAPSHOT.jar";
            System.setProperty("storm.jar", jarLocation);

            // 使用 StormSubmitter 提交本地构建的jar包
            StormSubmitter.submitTopology(args[0], conf, builder.createTopology());

        } else {
            LocalCluster cluster = new LocalCluster();
            cluster.submitTopology(LogTopology.class.getSimpleName(), conf, builder.createTopology());
        }
    }
}

创建类cn.delucia.storm.log.LogMockSpout

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import org.apache.storm.spout.SpoutOutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichSpout;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Values;
import org.apache.storm.utils.Utils;

import java.util.Map;
import java.util.Random;
import java.util.concurrent.*;
import java.util.concurrent.atomic.AtomicInteger;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 统计当日日志中的PV和UV指标
 * PV: Page View 每次访问技术一次
 * UV: Unique View 相同用户无论访问多少次都只计数1次
 * 如何判断是否为相同用户:
 * - IP 用户可能更换设备,或者使用VPN工具, 因此根据IP的统计可能不准确
 * - Cookie(SessionID) 一般使用cookie判断
 */
public class LogMockSpout extends BaseRichSpout {
    /**
     * Mock Data
     * 注意这里伪造的日期需要有当天的日期,否则最终统计的PV和UV结果都为 0
     * 假设今天是 2021年10月2日。
     * 2022年的日期因为晚于今天, 因此是无效数据, 需要重试3次
     */
    private final String[] dates = {"2021-10-02 08:40:50", "2022-10-02 18:40:50", "2022-10-02 19:40:50", "2021-10-02 20:40:50"};
    private final String[] sids = {"user100", "user101", "user102", "user103"};
    private final String[] urls = {"https://cn.delucia/a.html", "https://cn.delucia/b.html", "https://cn.delucia/c.html", "http://cn.delucia/d.html"};
    /**
     * Spout发射的tuple数量
     */
    private static final Integer EMMIT_TUPLE_NUMBER = 30;
    /**
     * tuple处理失败后进行重试的次数
     */
    private Integer retryTimes = 0;
    /**
     * 允许的最大重试次数
     */
    private static final Integer MAX_RETRY_TIMES = 3;
    private Integer count = 0;
    /**
     * 缓存所有发送的消息, 消息发射后移除
     */
    private ConcurrentHashMap<Object, String> cachedMessages;
    /**
     * 用于发射tuple到下游
     */
    private transient SpoutOutputCollector collector;
    /**
     * 用于生产随机数
     */
    private Random rand;
    /**
     * 自定义线程池 - 用于延迟发送需要重试的数据
     */
    private transient ThreadPoolExecutor executor;

    /**
     * 每个线程初始化时调用
     *
     * @param conf
     * @param context
     * @param collector
     */
    @Override
    public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) {
        this.collector = collector;
        this.cachedMessages = new ConcurrentHashMap<>();
        this.rand = new Random();

        // 自定义线程池 - 用于执行重试任务
        int corePoolSize = 2;
        int maximumPoolSize = 4;
        long keepAliveTime = 10;
        TimeUnit unit = TimeUnit.SECONDS;
        BlockingQueue<Runnable> workQueue = new ArrayBlockingQueue<>(3);
        ThreadFactory threadFactory = new NameThreadFactory();
        RejectedExecutionHandler handler = new IgnorePolicy();

        executor = new ThreadPoolExecutor(corePoolSize, maximumPoolSize, keepAliveTime, unit, workQueue, threadFactory, handler);
        executor.prestartAllCoreThreads(); // 预启动所有核心线程
    }

    @Override
    public void nextTuple() {
        if (count >= EMMIT_TUPLE_NUMBER) {
            return;
        }
        int k = rand.nextInt(4);
        int l = rand.nextInt(4);
        int m = rand.nextInt(4);
        String line = dates[k] + "\t" + sids[l] + "\t" + urls[m];
        String msg = ">> Emit Tuple[\"line\"]: " + line;
        System.err.println(msg);
        SocketClientUtil.send(this, msg);
        Long time = System.currentTimeMillis();
        this.collector.emit(new Values(line), time);
        cachedMessages.put(time, line);
        Utils.sleep(500);
        count++;
    }

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

    @Override
    public void ack(Object msgId) {
        cachedMessages.remove(msgId);
    }

    @Override
    public void fail(Object msgId) {
        cachedMessages.forEach((id, message) -> {
            if (msgId.equals(id)) {
                cachedMessages.put(id, message);
                retryTimes += 1;
                // 重试
                if (retryTimes <= MAX_RETRY_TIMES) {
                    // 同步改异步
                    // collector.emit(new Values(message), msgId);
                    RetryTask retryTask = new RetryTask(collector, message, msgId, retryTimes);
                    executor.execute(retryTask);
                } else {
                    // 超过最大重试次数, 放弃
                    retryTimes = 0;
                    cachedMessages.remove(msgId);
                }
            }
        });
    }

    class NameThreadFactory implements ThreadFactory {

        private final AtomicInteger mThreadNum = new AtomicInteger(1);

        @Override
        public Thread newThread(Runnable r) {
            Thread t = new Thread(r, "my-thread-" + mThreadNum.getAndIncrement());
            // String msg = "Start Thread[" + t.getName() + "] will retry in 30 secs, max retry times is " + MAX_RETRY_TIMES + ".";
            // System.err.println(msg);
            // SocketClientUtil.send(this, msg);
            return t;
        }
    }

    static class RetryTask implements Runnable {
        private final SpoutOutputCollector collector;
        private final String messageBody;
        private final Object messageId;
        private final Integer retryTimes;
        private final Random rand;

        public RetryTask(SpoutOutputCollector collector, String name, Object msgId, Integer retryTimes) {
            this.collector = collector;
            this.messageBody = name;
            this.messageId = msgId;
            this.retryTimes = retryTimes;
            this.rand = new Random();
        }

        @Override
        public void run() {
            // 30秒内进行重试
            Utils.sleep(rand.nextInt(30000));
            collector.emit(new Values(messageBody), messageId);
            String msg = ">> Replay: " + retryTimes + " times: " + messageBody;
            System.err.println(msg);
            SocketClientUtil.send(this, msg);
        }

        @Override
        public String toString() {
            return "Task: " + messageBody;
        }
    }

    public static class IgnorePolicy implements RejectedExecutionHandler {

        @Override
        public void rejectedExecution(Runnable r, ThreadPoolExecutor e) {
            log(r, e);
        }

        private void log(Runnable r, ThreadPoolExecutor e) {
            String msg = r.toString() + " rejected";
            System.err.println(msg);
            SocketClientUtil.send(this, msg);
        }
    }
}

创建类cn.delucia.storm.log.LogCheckBolt

package cn.delucia.storm.log;

import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 检查日志的格式是否符合要求
 */
public class LogCheckBolt extends BaseRichBolt {
    private static final long serialVersionUID = 1L;

    OutputCollector collector;

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {

        String line = input.getStringByField("line");
        if (line == null || line.isEmpty()) {
            this.collector.fail(input);
            return;
        }
        String date = line.split("\t")[0];
        String sid = line.split("\t")[1];

        SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
        try {
            Date d = sdf.parse(date);
            if (d.after(new Date())) {
                this.collector.fail(input);
                return;
            }
            date = sdf.format(d);
            this.collector.emit(input, new Values(date, sid));
            this.collector.ack(input);
        } catch (ParseException e) {
            e.printStackTrace();
            // 日期格式不正确
            this.collector.fail(input);
        }
    }

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

创建类cn.delucia.storm.log.LogCountBolt

package cn.delucia.storm.log;

import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Fields;
import org.apache.storm.tuple.Tuple;
import org.apache.storm.tuple.Values;

import java.util.HashMap;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 对相同的sessionID进行聚合计数
 */
public class LogCountBolt extends BaseRichBolt {

    OutputCollector collector;
    // <session_id>, <count>
    Map<String, Long> map = new HashMap<>();

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {

        String date = input.getStringByField("date");
        String ip = input.getStringByField("sid");
        String key = date + "_" + ip;
        Long count = 0L;
        try {
            count = map.get(key);
            if (count == null) {
                count = 0L;
            }
            count++;
            map.put(key, count);
            this.collector.emit(input, new Values(date + "_" + ip, count));
            this.collector.ack(input);
        } catch (Exception e) {
            e.printStackTrace();
            System.err.println(LogCountBolt.class.getSimpleName() + " failed, date: " + date + ", sid: " + ip + ", count: " + count);
            this.collector.fail(input);
        }
    }

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

创建类cn.delucia.storm.log.LogGlobalSumBolt

package cn.delucia.storm.log;

import cn.delucia.storm.utils.SocketClientUtil;
import org.apache.storm.task.OutputCollector;
import org.apache.storm.task.TopologyContext;
import org.apache.storm.topology.OutputFieldsDeclarer;
import org.apache.storm.topology.base.BaseRichBolt;
import org.apache.storm.tuple.Tuple;

import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.util.Date;
import java.util.HashMap;
import java.util.Map;

/**
 * @Author: deLucia
 * @Date: 2021/9/28
 * @Version: 1.0
 * @Description: 全局聚合 - 只统计当日的 pv 和 uv
 */
public class LogGlobalSumBolt extends BaseRichBolt {
    SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd");
    private static final long serialVersionUID = 1L;
    transient OutputCollector collector;
    Map<String, Long> map = new HashMap<>();

    @Override
    public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) {
        this.collector = collector;
    }

    @Override
    public void execute(Tuple input) {
        int pv = 0;
        int uv = 0;

        String dateIp = input.getStringByField("date_sid");
        Long count = input.getLongByField("count");
        String currDate = sdf.format(new Date());
        collector.ack(input);
        try {
            Date accessDate = sdf.parse(dateIp.split("_")[0]);
            // 访问日期不能晚于今天
            map.put(dateIp, count);
            String msg = ">> date_sid: " + dateIp + ", pv: " + count;
            System.err.println(msg);
            SocketClientUtil.send(this, msg);

            // 只计算当日的pv和uv
            if (!dateIp.split("_")[0].startsWith(currDate)) {
                return;
            }
            if (!map.isEmpty()) {
                for (Map.Entry<String, Long> e : map.entrySet()) {
                    uv++;
                    pv += e.getValue();
                }
            }
            String result = ">> " + currDate + " pv: " + pv + ", uv: " + uv;
            System.err.println(result);
            SocketClientUtil.send(this, result);
        } catch (ParseException e1) {
            e1.printStackTrace();
        } finally {
            // 到此整棵tuple树处理完成
            collector.ack(input);
        }

    }

    @Override
    public void declareOutputFields(OutputFieldsDeclarer declarer) {

    }

}

本地运行

在本机安装netcat,使用一下命令开启Socket Server:

>nc -l localhost -L -p 9876

然后本地运行项目(不要带参数),等待Socket Server接收到输出结果,如图所示:

file

本地提交打好的jar包提交到远程集群运行

首先在pom.xml中将依赖storm-core的scope改为provided
再使用maven package打包,生成的文件在项目根目录下的target文件夹内
修改LogTopology.java中的代码,指定jar包的本地提交路径:

String jarLocation = ".\\target\\parallelisim-1.0-SNAPSHOT.jar";
System.setProperty("storm.jar", jarLocation);

然后直接运行项目,即可通过Thrift协议协议将jar包远程提交到集群中的nimbus节点运行。
最后可以在Storm UI界面查看结果。

Views: 351

Storm拓扑之Stream Grouping

Stream的分组策略

Stream Grouping - 定义了一个流在Bolt任务间该如何被切分,谁来处理哪些数据流,按照什么规则来分配.

随机分组

Shuffle Grouping- 随机分组, 随机派发stream里面的tuple,保证每个bolt接收到的tuple数目大致相同。

字段分组

Fields grouping - 根据指定字段的值进行分组。比如说,一个数据流根据'word'字段进行分组,所有具有相同的'word'字段值的tuple会路由到同一个bolt的task中。

广播分组

All grouping- 全复制分组(广播分组), 将所有的tuple复制后分发给所有的bolt task。每个订阅数据流的task都会接收到tuple的拷贝。

全局分组

Global Grouping- 全局分组:这种分组方式将所有的tuple路由到唯一一个task上。Storm按照最小的taskID来选取接收数据的task。注意!!当使用全局分组方式时,设置bolt的task并发度是没有意义的,因为所有tuple都转发到同一个task上了。使用全局分组的时候需要注意,因为所有tuple都转发到一个JVM实例上,可能会引起Storm集群中某个JVM或者服务器出现性能瓶颈或奔溃。

无分组

None Grouping - 不分组,这个分组的意思是说stream不关心到底谁会收到它的tuple。目前这种分组和Shuffle grouping是一样的效果, 有一点不同的是storm会把这个bolt放到这个bolt的订阅者同一个线程里面去执行。

直接分组

Direct Grouping — 直接分组, 这是一种比较特别的分组方法,用这种分组意味着消息的发送者指定由消息接收者的哪个task处理这个消息。 只有被声明为Direct Stream的消息流可以声明这种分组方法。而且这种消息tuple必须使用emitDirect方法来发射。消息处理者可以通过TopologyContext来获取处理它的消息的task的id (OutputCollector.emit方法也会返回task的id)。

本地或随机分组

Local or shuffle grouping - 本地或随机分组:和随机分组类似,但是,会将tuple分发给同一个worker内的bolt task(如果worker内有接收数据的bolt task)。其他情况下,采用随机分组的方式。取决于topology的并发度,本地或随机分组可以减少网络传输,从而提高topology性能。

作业:

自己探索并实现其中任意一种分组方式,并添加说明性注释和运行截图

代码下载:
1494-2021-09-18-storm_Chap2Ex

Views: 274

Storm 累加拓扑示例

创建Spout发送递增数字数列

private static class NumSpout extends BaseRichSpout {
    private SpoutOutputCollector collector;
    private int num = 1;

    @Override
    public void open(Map<String, Object> conf, TopologyContext context, SpoutOutputCollector collector) {

        this.collector = collector;
    }

    @Override
    public void nextTuple() {

        collector.emit(new Values(num));
        System.err.println("+ " + num);
        Utils.sleep(1000);
        num++;
    }

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

创建Bolt负责计算累加结果

    private static class SumBolt extends BaseRichBolt {

        // private OutputCollector collector;

        private int sum = 0;
        @Override
        public void prepare(Map<String, Object> topoConf, TopologyContext context, OutputCollector collector) {

            // this.collector = collector;
        }

        @Override
        public void execute(Tuple input) {
            //int number = input.getInteger(0);
            int number = input.getIntegerByField("num");
            sum += number;
            // collector.emit(new Values(sum));
            System.err.println("= " + sum);
        }

        @Override
        public void declareOutputFields(OutputFieldsDeclarer declarer) {

            // declarer.declare(new Fields("sum"));
        }
    }

本地运行

本地运行完整代码

public class MySumTopology {

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

        Config config = new Config();
        config.setNumWorkers(1);
        config.setDebug(true);

        TopologyBuilder topologyBuilder = new TopologyBuilder();
        topologyBuilder.setSpout("num-spout", new NumSpout());
        topologyBuilder.setBolt("sum-bolt", new SumBolt()).shuffleGrouping("num-spout");

        LocalCluster localCluster = new LocalCluster();
        localCluster.submitTopology("sum-topo", config, topologyBuilder.createTopology());
        // 运行1分钟停止
        Utils.sleep(60000);
        localCluster.shutdown();
    }
}

保证数据可靠处理

  1. Spout在使用nextTuple()方法发送数据时需要传入消息ID
    @Override
    public void nextTuple() {
        collector.emit(new Values(num), num);
        ...
    }
  1. Bolt中execute()方法中标记tuple是否处理成功
  • 处理成功 collector.ack(input)

  • 处理失败collector.fail(input)

    注意: ackfail方法需要锚定到发射过来的tuple上.

    private static class SumBolt extends BaseRichBolt {

        private OutputCollector collector;
        private Random rand;

        private int sum = 0;

        @Override
        public void execute(Tuple input) {

            if (rand.nextInt(10) < 2) { 
                collector.fail(input); // 20%几率 处理失败 
                System.err.println("!搞错了");
                return;
            }

            int number = input.getIntegerByField("num");
            sum += number;
            // collector.emit(new Values(sum));
            System.err.println("= " + sum);
            collector.ack(input); // 处理成功: Acknowledge
        }
    ...
  1. Spout中对处理失败的元组触发回调

    这里把处理失败的数字重新发射出去

    @Override
    public void fail(Object msgId) {
        int oldNum = num - 1;
        collector.emit(new Values(oldNum), msgId);
        System.out.println("重来 + " + oldNum);
        Utils.sleep(1000);
    }

远程运行模式

如果是需要提交到集群, 需要使用StormSubmitter替换new LocalCluster()来调用submitTopology(..)方法

// local mode运行 - LocalCluster
LocalCluster cluster = new LocalCluster();
cluster.submitTopology("sum-topo", config, builder.createTopology());

// remote mode - run on a production cluster
// StormSubmitter.submitTopology("sum-topo",config,builder.createTopology());

然后打成jar包上传到集群运行

Views: 222

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: 376

Index