【问题标题】:Creating an Apache Storm spout which emits tuples every X seconds创建一个每 X 秒发出一次元组的 Apache Storm spout
【发布时间】:2014-10-27 19:15:07
【问题描述】:

我有一个从 MQTT 代理接收数据的拓扑,我希望 spout 的行为如下:

  1. 每 x 秒发出一批元组(或单个元组中的字符串列表)。我如何实现这一目标?我读过一些关于 Storm Trident 的文章,但它的 IBatchSpout 似乎不允许我以特定的时间间隔批量发出元组。

  2. 如果没有新数据进来,spout 应该怎么做?它不能阻塞线程,因为它是 Storm 的主线程,对吧?

【问题讨论】:

    标签: stream apache-storm mqtt trident


    【解决方案1】:

    您可以实现自己的 MQTT spout。例如,请查看MongoSpout

    重要的部分是nextTuple 方法。

    当这个方法被调用时,Storm 正在请求 Spout 发射 元组到输出收集器。 这个方法应该是非阻塞的,所以 如果 Spout 没有要发出的元组,则此方法应返回。 nextTuple、ack 和 fail 都在一个紧密循环中调用 spout 任务中的线程。当没有要发出的元组时,它是 礼貌地让 nextTuple 睡一小段时间(比如 单毫秒),以免浪费过多的 CPU。

    您不能一次等待指定的时间,但您可以实现nextTuple,以便它只偶尔发出一次元组。

    private static final EMISSION_PERIOD = 2000; // 2 seconds
    private long lastEmission;
    
    @Override
    public void nextTuple() {
        if (lastEmission == null ||
                lastEmission + EMISSION_PERIOD >= System.currentMillis()) {
            List<Object> tuple = pollMQTT();
            if (tuple != null) {
                this.collector.emit(tuple);
                return;
            }
        }
        Utils.sleep(50);
    }
    

    请注意,我找到了一个开源 MQTT spout。它看起来还没有准备好生产,但您可以将其用作起点。

    【讨论】:

      【解决方案2】:

      除了 Christian,我还为 Storm 的 MQTT 客户端找到了 this implementation。前面提到的链接还没有开发。

      【讨论】:

        猜你喜欢
        • 2018-08-14
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2018-07-10
        • 2020-05-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多