【发布时间】:2015-11-27 02:17:26
【问题描述】:
我正在尝试实现有保证的消息处理,但没有调用 Spout 上的 ack 或 fail 方法。
我正在通过 spout 传递消息 ID 对象。 我用每个螺栓传递元组并在每个螺栓中调用collector.ack(tuple)。
问题 ack 或 fail 没有被调用,我不知道为什么?
这是一个简短的代码示例。
使用 BaseRichSpout 的 Spout 代码
public void nextTuple() {
for( String usage : usageData ) {
.... further code ....
String msgID = UUID.randomUUID().toString()
+ System.currentTimeMillis();
Values value = new Values(splitUsage[0], splitUsage[1],
splitUsage[2], msgID);
outputCollector.emit(value, msgID);
}
}
@Override
public void ack(Object msgId) {
this.pendingTuples.remove(msgId);
LOG.info("Ack " + msgId);
}
@Override
public void fail(Object msgId) {
// Re-emit the tuple
LOG.info("Fail " + msgId);
this.outputCollector.emit(this.pendingTuples.get(msgId), msgId);
}
使用 BaseRichBolt 的螺栓代码
@Override
public void execute(Tuple inputTuple) {
this.outputCollector.emit(inputTuple, new Values(serverData, msgId));
this.outputCollector.ack(inputTuple);
}
最终螺栓
@Override
public void execute(Tuple inputTuple) {
..... Simply reports does not emit .....
this.outputCollector.ack(inputTuple);
}
【问题讨论】:
标签: java apache-storm