Storm自定義流劃分

原創(chuàng)文章,轉(zhuǎn)載請注明原作地址:www.itdecent.cn/p/2db35d7bb92f

在Storm開發(fā)過程中,經(jīng)常性需要將符合某個條件的消息分發(fā)到同一個partition。官方提供了一個Stream.partitionBy("fieldName")的API,可以根據(jù)每條tuple某個字段進(jìn)行流劃分,劃分方式是:

field.hashCode() mod num-task

雖然Clojure實際上是運行與java平臺上的一種Lisp方言,但是在java里直接計算field.hashCode() % num-task,發(fā)現(xiàn)得到的結(jié)果和實際partitionIndex并不一致。另外,這種原生的根據(jù)字段的哈希值取模進(jìn)行流劃分的方式也過于單一,因此直接去寫了一個Storm自定義流劃分函數(shù)的方法。實現(xiàn)代碼如下:

import java.util.Arrays;
import java.util.List;

import org.apache.storm.generated.GlobalStreamId;
import org.apache.storm.grouping.CustomStreamGrouping;
import org.apache.storm.task.WorkerTopologyContext;
import org.apache.storm.tuple.Fields;

public class HashModStreamGrouping implements CustomStreamGrouping {
    
    private List<Integer> targetTasks;
    private String partitionKeyName;
    private int partitionKeyIndex;

    public HashModStreamGrouping(String partitionKeyName) {
        this.partitionKeyName = partitionKeyName;
    }
    
    @Override
    public List<Integer> chooseTasks(int taskId, List<Object> values) {
        String partitionStr = String.valueOf(values.get(this.partitionKeyIndex));
        int partitionVal = getTaskIndex(partitionStr, this.targetTasks.size());
        return Arrays.asList(this.targetTasks.get(partitionVal));
    }

    @Override
    public void prepare(WorkerTopologyContext context, GlobalStreamId streamId,
            List<Integer> targetTasks) {
        this.targetTasks = targetTasks;
        
        Fields outputFields = context.getComponentOutputFields(streamId);
        this.partitionKeyIndex = outputFields.fieldIndex(this.partitionKeyName);
    }
    
    public static int getTaskIndex(String partitionStr, int targetSize) {
        return  (Math.abs(partitionStr.hashCode()) % targetSize );
    }
} 

調(diào)用方式如下,其中partitionValue是用于流劃分的字段名

Stream.partition(new HashModStreamGrouping("partitionValue"))
最后編輯于
?著作權(quán)歸作者所有,轉(zhuǎn)載或內(nèi)容合作請聯(lián)系作者
【社區(qū)內(nèi)容提示】社區(qū)部分內(nèi)容疑似由AI輔助生成,瀏覽時請結(jié)合常識與多方信息審慎甄別。
平臺聲明:文章內(nèi)容(如有圖片或視頻亦包括在內(nèi))由作者上傳并發(fā)布,文章內(nèi)容僅代表作者本人觀點,簡書系信息發(fā)布平臺,僅提供信息存儲服務(wù)。

相關(guān)閱讀更多精彩內(nèi)容

友情鏈接更多精彩內(nèi)容