美文网首页
Storm自定义流划分

Storm自定义流划分

作者: EchoZhan | 来源:发表于2016-12-04 22:44 被阅读0次

    原创文章,转载请注明原作地址:www.jianshu.com/p/2db35d7bb92f

    在Storm开发过程中,经常性需要将符合某个条件的消息分发到同一个partition。官方提供了一个Stream.partitionBy("fieldName")的API,可以根据每条tuple某个字段进行流划分,划分方式是:

    field.hashCode() mod num-task
    

    虽然Clojure实际上是运行与java平台上的一种Lisp方言,但是在java里直接计算field.hashCode() % num-task,发现得到的结果和实际partitionIndex并不一致。另外,这种原生的根据字段的哈希值取模进行流划分的方式也过于单一,因此直接去写了一个Storm自定义流划分函数的方法。实现代码如下:

    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 );
        }
    } 
    

    调用方式如下,其中partitionValue是用于流划分的字段名

    Stream.partition(new HashModStreamGrouping("partitionValue"))
    

    相关文章

      网友评论

          本文标题:Storm自定义流划分

          本文链接:https://www.haomeiwen.com/subject/nqqkmttx.html