【问题标题】:Using tick tuples with trident in storm在风暴中使用带有三叉戟的刻度元组
【发布时间】:2014-08-16 13:35:32
【问题描述】:

我可以使用标准的 spout、bolt 组合来进行流式聚合 并且在快乐的情况下工作得很好,当使用刻度元组以某个时间间隔保存数据时 使用批处理。现在我自己在做一些故障管理(跟踪未保存的元组等)。(即不是风暴中的 ootb)

但是我已经读到三叉戟给了你更高的抽象和更好的故障管理。 我不明白的是三叉戟中是否有刻度元组支持。基本上 我想在当前分钟左右在内存中批处理并保留任何聚合数据 前几分钟使用三叉戟。

这里的任何指针或设计建议都会有所帮助。

谢谢

【问题讨论】:

  • +1。似乎它没有在 api 中公开/无法通过 getSourceStreamId() 判断 TridentTuple 是否是刻度元组
  • 好问题。我刚刚为标准 spout/bolt 开发了自己的批处理工具,但我仍然不知道如何以特定频率使用 trident。

标签: apache-storm trident


【解决方案1】:

实际上微批处理是 Trident 的内置功能。您不需要任何刻度元组。当你的代码中有这样的东西时:

topology
    .newStream("myStream", spout)
    .partitionPersist(
        ElasticSearchEventState.getFactoryFor(connectionProvider),
        new Fields("field1", "field2"),
        new ElasticSearchEventUpdater()
    )

(我在这里使用我的自定义 ElasticSearch 状态/更新器,您可能会使用其他东西)

所以当你有这样的东西时,在 Trident 引擎盖下将你的流分组并执行 partitionPersist 操作,而不是对单个元组,而是对这些批次。

如果您出于任何原因仍然需要刻度元组,只需创建刻度喷口,这样的东西对我有用:

public class TickSpout implements IBatchSpout {

    public static final String TIMESTAMP_FIELD = "timestamp";
    private final long delay;

    public TickSpout(long delay) {
        this.delay = delay;
    }

    @Override
    public void open(Map conf, TopologyContext context) {
    }

    @Override
    public void emitBatch(long batchId, TridentCollector collector) {
        Utils.sleep(delay);
        collector.emit(new Values(System.currentTimeMillis()));
    }

    @Override
    public void ack(long batchId) {
    }

    @Override
    public void close() {
    }

    @Override
    public Map getComponentConfiguration() {
        return null;
    }

    @Override
    public Fields getOutputFields() {
        return new Fields(TIMESTAMP_FIELD);
    }
}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-05-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-09-27
    • 1970-01-01
    • 1970-01-01
    • 2014-11-14
    相关资源
    最近更新 更多