【问题标题】:Creating an Apache Storm spout which emits tuples every X seconds创建一个每 X 秒发出一次元组的 Apache Storm spout
【发布时间】:2014-10-27 19:15:07
【问题描述】:
我有一个从 MQTT 代理接收数据的拓扑,我希望 spout 的行为如下:
每 x 秒发出一批元组(或单个元组中的字符串列表)。我如何实现这一目标?我读过一些关于 Storm Trident 的文章,但它的 IBatchSpout 似乎不允许我以特定的时间间隔批量发出元组。
如果没有新数据进来,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。它看起来还没有准备好生产,但您可以将其用作起点。