【问题标题】:storm cluster mode, distributed bolt/worker load sharingStorm集群模式,分布式bolt/worker负载分担
【发布时间】:2015-05-01 05:39:36
【问题描述】:

HI:我将有一个大容量风暴分析任务。对我来说,我想在不同的节点/机器上分拆许多螺栓/工作人员来完成任务,以便每台机器都可以分担负载。我想知道如何编写螺栓/工人/拓扑,以便他们可以相互通信。在下面的代码中,我在一台机器上提交拓扑,如何在其他机器上编写bolt/worker/config,以便拓扑知道其他机器的bolt/worker。我想我无法在一台机器上提交拓扑并在其他机器上提交相同的拓扑。 关于风暴工人负载分担的任何提示?

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

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import storm.kafka.KafkaSpout;
import storm.kafka.SpoutConfig;
import storm.kafka.StringScheme;
import storm.kafka.ZkHosts;
import backtype.storm.Config;
import backtype.storm.StormSubmitter;
import backtype.storm.generated.AlreadyAliveException;
import backtype.storm.generated.InvalidTopologyException;
import backtype.storm.spout.SchemeAsMultiScheme;
import backtype.storm.task.OutputCollector;
import backtype.storm.task.TopologyContext;
import backtype.storm.topology.OutputFieldsDeclarer;
import backtype.storm.topology.TopologyBuilder;
import backtype.storm.topology.base.BaseRichBolt;
import backtype.storm.tuple.Fields;
import backtype.storm.tuple.Tuple;
import backtype.storm.tuple.Values;

public class StormClusterMain {
    private static final String SPOUTNAME="KafkaSpout"; 
    private static final String ANALYSISBOLT = "ClusterAnalysisWorker";
    private static final String CLIENTID = "ClusterStorm";
    private static final String TOPOLOGYNAME = "ClusterTopology";

    private static class AppAnalysisBolt extends BaseRichBolt {
        private static final long serialVersionUID = -6885792881303198646L;
        private static final String collectionName="clusterusers";
        private OutputCollector _collector;
        private AtomicInteger index = new AtomicInteger(0); 
        private static final Logger boltLogger = LoggerFactory.getLogger(AppAnalysisBolt.class); 

        public void prepare(Map conf, TopologyContext context, OutputCollector collector) {
            _collector = collector;
        }

        public void execute(Tuple tuple) {  
            boltLogger.error("Message received:"+tuple.getString(0));
            _collector.emit(tuple, new Values(tuple.getString(0) + "!!!"));
            _collector.ack(tuple);
        }

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


    }

   public static void main(String[] args) throws AlreadyAliveException, InvalidTopologyException{

       String zookeepers = null;
       String topicName = null;
       if(args.length == 2 ){
           zookeepers = args[0];
           topicName = args[1];
       }else{
           System.out.println("You need to have two arguments: kafka zookeeper:port and topic name");
           System.out.println("Usage :.xxx");
           System.exit(-1);
       }        

       SpoutConfig spoutConfig = new SpoutConfig(new ZkHosts(zookeepers),
            topicName,
            "",// zookeeper root path for offset storing
            CLIENTID);
       spoutConfig.scheme = new SchemeAsMultiScheme(new StringScheme());
       KafkaSpout kafkaSpout = new KafkaSpout(spoutConfig);

       TopologyBuilder builder = new TopologyBuilder();
       builder.setSpout(SPOUTNAME, kafkaSpout, 1);
       builder.setBolt(ANALYSISBOLT, new AppAnalysisBolt())
                                                 .shuffleGrouping(SPOUTNAME);

        //Configuration
        Config conf = new Config();
        conf.setDebug(false);
        //Topology run
        conf.setNumWorkers(3);
        StormSubmitter.submitTopologyWithProgressBar(TOPOLOGYNAME, conf, builder.createTopology());

【问题讨论】:

    标签: apache-storm


    【解决方案1】:

    除非出现问题,否则你已经完成了。

    当您向 Storm 提交拓扑时,Nimbus 服务会通过遍布整个集群的 Supervisor 进程查看集群上的负载。 Nimbus 然后为拓扑运行提供一定数量的资源。这些资源通常分布在集群中的各个主管节点中,它们将并行处理元组。 Nimbus 偶尔会重新审视这些决策并更改哪些节点处理什么,以试图保持集群中的负载平衡。作为用户,您永远不应该注意到该过程。

    假设您的 Storm 集群设置正确,您唯一需要做的就是提交拓扑。 Storm 会为您处理整个多节点并行处理。

    也就是说,你在其中的 AtomicInteger 会表现得很奇怪,因为 Storm 会将你的代码分割到多个服务器上,甚至是单个主机上的多个 JVM。如果您想解决单个风暴进程需要了解较大集群状态的情况,最好将其外部化到某种独立的数据存储(即 redis 或 hbase)。

    【讨论】:

    • 答案很明确。我这里只是用我自己的方式解释。它是主管节点承担负载。说 builder.setBolt(ANALYSISBOLT, new AppAnalysisBolt(),10)。 shuffleGrouping(SPOUTNAME),这是否意味着我需要为集群分配/添加/分配足够的主管,以便 10 个工人/螺栓有足够的资源(内存,cpu)来并行消耗元组。我说的对吗?
    • 不,不是真的 :) 这个整数参数,又名 parallelism_hint,指定了执行者的数量。任何给定的工作进程都支持运行多个执行器线程。因此,除了任何其他配置,您将为每个运行拓扑的工作进程生成 10 个执行器线程。有关这方面的更多信息:how to tune the parallelism hint in storm
    • 感谢您的出色回答。现在我了解了执行线程和工作进程。我目前的问题是如何增加集群的容量。我要添加哪个节点? Nimbus,主管,工作进程?如果我有一台 Nimbus 机器,一台机器主管,一台机器提交拓扑和工作人员,如果容量用完了,我该怎么办?增加工人是合理的想法。但是那么如何添加工人机器呢?如何编写worker/bolt代码,让他们知道topology并让storm找到添加的work/bolt machine?当然在topology中,我会增加worker的数量。
    • 如果您有一个正常运行的 Storm 集群,为了增加容量,您可以在集群中添加一个 Supervisor 节点。就是这样。 Storm 负责其余的工作......我将在稍后修改我的答案,并提供有关其工作原理的更多信息。
    猜你喜欢
    • 2013-08-18
    • 2015-04-20
    • 2018-10-23
    • 1970-01-01
    • 1970-01-01
    • 2018-09-29
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多