【问题标题】:in storm, how to emit a List of primitives在风暴中,如何发出原语列表
【发布时间】:2015-08-28 19:53:30
【问题描述】:

我正在使用风暴 0.10 我有一个列表,我对其进行迭代然后发出元组

for (CustomObject o: List<CustomObject>) {
  collector.emit(STREAM_NAME, new Values(o.getFirst, o.getName, o.getAddress));
}

我不想发出多个元组,我只想发出一个元组,它是一个嵌套列表,就像这样,

我的主要问题是关于序列化。阅读 Storm 文档表明 java 序列化很昂贵,而 Storm 使用 Kryo 序列化。此外,不应该在没有 kryo 的情况下通过网络发送常规的 java POJO 类。所以我想发送一个List&lt;List&lt;Objects&gt;&gt; 如下,

List<List<Object>> valueList = new ArrayList<List<Object>>();
for (CustomObject o: List<CustomObject>) {
  v.add(new ArrayList<Object>{ 
    {
     add(o.getFirst);
     add( o.getName);
     add(o.getAddress);
    }
  });
}
collector.emit(STREAM_NAME, new Values(valueList));

所以问题是 - 这是使用 Kryo 完成的吗?

【问题讨论】:

标签: java-7 apache-storm


【解决方案1】:

你可以简单地发出你的List

collector.emit(new Values(listCustomObject));

然后当你像这样在bolt 中读回它时,

List<CustomObject> listCustomObject = (List<CustomObject>) tuple.getValue(0);

希望对你有帮助。

编辑:确实我忘了说对象必须是Serializable!为了处理列表,我会做类似的事情:

import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;


public class CustomObjectList implements Serializable {

    private static final long serialVersionUID = 6877020084704724252L;

    public class CustomObject {
        ...
    }

    private final List<CustomObject> list = new ArrayList<>();

    ...
}

EDIT2:因为我做的事情太多了,所以我通过修改storm starter ExaclamationTopology 中的ExclamationBolt 完成/测试了一个快速玩具示例

看起来像这样。

package storm.starter;

import java.io.Serializable;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;

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

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


public class ExclamationBolt extends BaseRichBolt{

    public static Logger LOG = LoggerFactory.getLogger(ExclamationBolt.class);    


    public class CustomObject {
        private String word;

        public CustomObject(String word) {
            setWord(word);
        }

        public String getWord() {
            return word;
        }

        public void setWord(String word) {
            this.word = word;
        }

    }
    public class CustomObjectList implements Serializable {

        private static final long serialVersionUID = 6877020084704724252L;

        private final List<CustomObject> list = new ArrayList<>();

        public CustomObjectList() {
        }

        public List<CustomObject> getObjects() {
            return list;
        }

        public void addObject(CustomObject o) {
            list.add(o);
        }
    }

    OutputCollector _collector;

    public ExclamationBolt() {
    }

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

    @Override
    public void execute(Tuple tuple) {
        String id = tuple.getSourceComponent();
        String word;
        if ("word".equals(id)) {
            word = tuple.getString(0);
        } else {
            CustomObjectList list = (CustomObjectList) tuple.getValue(0);
            word = list.getObjects().get(0).getWord();
        }
        word = word + "!!!";
        LOG.debug(word);
        CustomObjectList list = new CustomObjectList();
        list.addObject(new CustomObject(word));
        _collector.emit(tuple, new Values(list));
        _collector.ack(tuple);
    }

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


}

希望对你有帮助。

【讨论】:

  • 好的,我没有发出自定义对象,因为我不想处理 Kryo 序列化。 Storm 是否默认使用 kryo 序列化 List
  • 嗨。事实上,我忘了提到这个类应该是可序列化的。我正在更新答案。我希望它与我发布的代码没有问题。请告诉我。
  • 谢谢,但据我了解,根据文档,java 序列化非常昂贵且不受欢迎。所以我试图避免这种情况,问题是,通过使用“对象”列表[主要是原语和字符串],我是否可以明确地实现任何东西,并且只依赖 Storm 使用 Kryo 进行序列化,就像它对原语。
  • 好的,我明白你的意思了。我想唯一的方法是为 Storm Kyro 构建一个自定义序列化程序,但我从来没有这样做过。对不起。
【解决方案2】:

使用 KryoSerialization。

  1. 实现CustomObjectSerializer:

    class CustomObjectSerializer extends Serializer<CustomObject> {  
      @Override
      public void write(final Kryo kryo,
                        final Output output,
                        final CustomObject cObject) {
        // do serialization
      }
    
      @Override
      public CustomObject read(final Kryo kryo,
                               final Input input,
                               final Class<CustomObject> aClass) {
        // do deserialization
      }
    }
    
  2. 将 CustomObjectSerializer 注册到 Storm 配置: conf.registerSerialization(CustomObject.class,CustomObjectSerializer.class);

【讨论】:

  • 文档对此非常粗略。具体有没有任何例子, // 做序列化?关于读写?这是我试图避免的,因为我不清楚并发送嵌套 List
  • 据我所知,如果元素是可序列化的,您不必担心嵌套列表。如果 CustomObject 是可序列化的,Storm 会自动序列化 List>。如何序列化 CustomObject 取决于实现。请看github.com/EsotericSoftware/kryo
猜你喜欢
  • 1970-01-01
  • 2014-04-28
  • 2018-10-15
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2015-01-22
  • 1970-01-01
相关资源
最近更新 更多