天天看點

流式處理架構storm淺析(下篇)

本文來自網易雲社群

作者:汪建偉

  • 舉個栗子

1 實作的目标  

設計一個系統,來實作對一個文本裡面的單詞出現的頻率進行統計。

2 設計Topology結構:  

這是一個簡單的例子,topology也非常簡單。整個topology如下:

流式處理架構storm淺析(下篇)

整個topology分為三個部分:

WordReader:資料源,負責發送sentence

WordNormalizer:負責将sentence切分

Wordcounter:負責對單詞的頻率進行累加

3 代碼實作

1. 建構maven環境,添加storm依賴

<repositories>  
	      <!-- Repository where we can found the storm dependencies  -->  
	      <repository>  
	          <id>clojars.org</id>  
	          <url>http://clojars.org/repo</url>  
	      </repository>  
	</repositories>  
	<dependencies>  
	      <dependency>   
	        <groupId>storm</groupId>  
	        <artifactId>storm</artifactId>  
	        <version>0.7.1</version>  
	     </dependency>  
	</dependencies>      

    2. 定義Topology

public class TopologyMain {  
	    public static void main(String[] args) throws InterruptedException {  
	           
	        //Topology definition  
	        TopologyBuilder builder = new TopologyBuilder();  
	        builder.setSpout("word-reader",new WordReader());  
	        builder.setBolt("word-normalizer", new WordNormalizer())  
	            .shuffleGrouping("word-reader");  
	        builder.setBolt("word-counter", new WordCounter(),1)  
	            .fieldsGrouping("word-normalizer", new Fields("word"));  
	          
	        //Configuration  
	        Config conf = new Config();  
	        conf.put("wordsFile", args[0]);  
	        conf.setDebug(false);  
	        //Topology run  
	        conf.put(Config.TOPOLOGY_MAX_SPOUT_PENDING, 1);  
	        LocalCluster cluster = new LocalCluster();  
	        cluster.submitTopology("Getting-Started-Toplogie", conf, builder.createTopology());  
	        Thread.sleep(1000);  
	        cluster.shutdown();  
	    }  
	}      

    3. 實作WordReader Spout

public class WordReader extends BaseRichSpout {  
	  
	    private SpoutOutputCollector collector;  
	    private FileReader fileReader;  
	    private boolean completed = false;  
	    public void ack(Object msgId) {  
	        System.out.println("OK:"+msgId);  
	    }  
	    public void close() {}  
	    public void fail(Object msgId) {  
	        System.out.println("FAIL:"+msgId);  
	    }  
	  
	    public void nextTuple() {  
	        if(completed){  
	            try {  
	                Thread.sleep(1000);  
	            } catch (InterruptedException e) {  
	            }  
	            return;  
	        }  
	        String str;  
	        BufferedReader reader = new BufferedReader(fileReader);  
	        try{  
	            while((str = reader.readLine()) != null){  
	                this.collector.emit(new Values(str),str);  
	            }  
	        }catch(Exception e){  
	            throw new RuntimeException("Error reading tuple",e);  
	        }finally{  
	            completed = true;  
	        }  
	    }  
	  
	    public void open(Map conf, TopologyContext context,  
	                     SpoutOutputCollector collector) {  
	        try {  
	            this.fileReader = new FileReader(conf.get("wordsFile").toString());  
	        } catch (FileNotFoundException e) {  
	            throw new RuntimeException("Error reading file ["+conf.get("wordFile")+"]");  
	        }  
	        this.collector = collector;  
	    }  
	  
	    public void declareOutputFields(OutputFieldsDeclarer declarer) {  
	        declarer.declare(new Fields("line"));  
	    }  
	}      

第一個被調用的spout方法都是public void open(Map conf, TopologyContext context, SpoutOutputCollector collector)。它接收如下參數:配置對象,在定義topology對象是建立;TopologyContext對象,包含所有拓撲資料;還有SpoutOutputCollector對象,它能讓我們釋出交給bolts處理的資料。

    4. 實作WordNormalizer bolt

public class WordNormalizer extends BaseBasicBolt {  
	  
	    public void cleanup() {}  
	  
	    public void execute(Tuple input, BasicOutputCollector collector) {  
	        String sentence = input.getString(0);  
	        String[] words = sentence.split(" ");  
	        for(String word : words){  
	            word = word.trim();  
	            if(!word.isEmpty()){  
	                word = word.toLowerCase();  
	                collector.emit(new Values(word));  
	            }  
	        }  
	    }  
	      
	    public void declareOutputFields(OutputFieldsDeclarer declarer) {  
	        declarer.declare(new Fields("word"));  
	    }  
	}      

bolt最重要的方法是void execute(Tuple input),每次接收到元組時都會被調用一次,還會再釋出若幹個元組。

    5. 實作WordCounter bolt  

public class WordCounter extends BaseBasicBolt {  
	  
	    Integer id;  
	    String name;  
	    Map counters;  
	  
	    @Override  
	    public void cleanup() {  
	        System.out.println("-- Word Counter ["+name+"-"+id+"] --");  
	        for(Map.Entry entry : counters.entrySet()){  
	            System.out.println(entry.getKey()+": "+entry.getValue());  
	        }  
	    }  
	  
	    @Override  
	    public void prepare(Map stormConf, TopologyContext context) {  
	        this.counters = new HashMap();  
	        this.name = context.getThisComponentId();  
	        this.id = context.getThisTaskId();  
	    }  
	  
	    @Override  
	    public void declareOutputFields(OutputFieldsDeclarer declarer) {}  
	  
	    @Override  
	    public void execute(Tuple input, BasicOutputCollector collector) {  
	        String str = input.getString(0);  
	        if(!counters.containsKey(str)){  
	            counters.put(str, 1);  
	        }else{  
	            Integer c = counters.get(str) + 1;  
	            counters.put(str, c);  
	        }  
	    }  
	}      

    6. 使用本地模式運作Topology

在這個目錄下面建立一個檔案,/src/main/resources/words.txt,一個單詞一行,然後用下面的指令運作這個拓撲:mvn exec:java -Dexec.main -Dexec.args=”src/main/resources/words.txt”。

如果你的words.txt檔案有如下内容: Storm test are great is an Storm simple application but very powerful really Storm is great 你應該會在日志中看到類似下面的内容: is: 2 application: 1 but: 1 great: 1 test: 1 simple: 1 Storm: 3 really: 1 are: 1 great: 1 an: 1 powerful: 1 very: 1 在這個例子中,每類節點隻有一個執行個體。

  • 附-Storm記錄級容錯的基本原理

首先來看一下什麼叫做記錄級容錯?storm允許使用者在spout中發射一個新的源tuple時為其指定一個message id, 這個message id可以是任意的object對象。多個源tuple可以共用一個message id,表示這多個源 tuple對使用者來說是同一個消息單元。storm中記錄級容錯的意思是說,storm會告知使用者每一個消息單元是否在指定時間内被完全處理了。那什麼叫做完全處理呢,就是該message id綁定的源tuple及由該源tuple後續生成的tuple經過了topology中每一個應該到達的bolt的處理。舉個例子。在圖4-1中,在spout由message 1綁定的tuple1和tuple2經過了bolt1和bolt2的處理生成兩個新的tuple,并最終都流向了bolt3。當這個過程完成處理完時,稱message 1被完全處理了。

流式處理架構storm淺析(下篇)

在storm的topology中有一個系統級元件,叫做acker。這個acker的任務就是追蹤從spout中流出來的每一個message id綁定的若幹tuple的處理路徑,如果在使用者設定的最大逾時時間内這些tuple沒有被完全處理,那麼acker就會告知spout該消息處理失敗了,相反則會告知spout該消息處理成功了。在剛才的描述中,我們提到了”記錄tuple的處理路徑”,如果曾經嘗試過這麼做的同學可以仔細地思考一下這件事的複雜程度。但是storm中卻是使用了一種非常巧妙的方法做到了。在說明這個方法之前,我們來複習一個數學定理。

A xor A = 0.

A xor B…xor B xor A = 0,其中每一個操作數出現且僅出現兩次。

storm中使用的巧妙方法就是基于這個定理。具體過程是這樣的:在spout中系統會為使用者指定的message id生成一個對應的64位整數,作為一個root id。root id會傳遞給acker及後續的bolt作為該消息單元的唯一辨別。同時無論是spout還是bolt每次新生成一個tuple的時候,都會賦予該tuple一個64位的整數的id。Spout發射完某個message id對應的源tuple之後,會告知acker自己發射的root id及生成的那些源tuple的id。而bolt呢,每次接受到一個輸入tuple處理完之後,也會告知acker自己處理的輸入tuple的id及新生成的那些tuple的id。Acker隻需要對這些id做一個簡單的異或運算,就能判斷出該root id對應的消息單元是否處理完成了。下面通過一個圖示來說明這個過程。

流式處理架構storm淺析(下篇)

上圖 spout中綁定message 1生成了兩個源tuple,id分别是0010和1011.

流式處理架構storm淺析(下篇)

上圖 bolt1處理tuple 0010時生成了一個新的tuple,id為0110.

流式處理架構storm淺析(下篇)

上圖 bolt2處理tuple 1011時生成了一個新的tuple,id為0111.

流式處理架構storm淺析(下篇)

上圖 bolt3中接收到tuple 0110和tuple 0111,沒有生成新的tuple.

容錯過程存在一個可能出錯的地方,那就是,如果生成的tuple id并不是完全各異的,acker可能會在消息單元完全處理完成之前就錯誤的計算為0。這個錯誤在理論上的确是存在的,但是在實際中其機率是極低極低的,完全可以忽略。

相關閱讀:流式處理架構storm淺析(上篇)

網易雲免費體驗館,0成本體驗20+款雲産品! 

更多網易研發、産品、營運經驗分享請通路網易雲社群。

相關文章:

【推薦】 一個隻有十行的精簡MVVM架構(下篇)

【推薦】 責任鍊模式的使用-Netty ChannelPipeline和Mina IoFilterChain分析

繼續閱讀