【发布时间】:2015-10-20 21:15:05
【问题描述】:
我正在尝试开始使用 Storm Trident,并使用 IOpaquePartitionedTridentSpout 设置和运行拓扑并由 OpaqueMap 支持。
但是,我很难找到让我的 spout/函数知道事务是否成功提交的方法。我没有在常规 Storm spout/bolt 界面中看到任何 ack 或 fail 方法。
我的用例是仅在同一类别的前一个被处理和持久化(或失败)时发出一个类别的元组。因为我将使用处理后的数据来更新我的下一个类别的元组。来自不同类别的元组可以并行处理。
使用partitionBy 方法按类别对流进行分区。
将max_spout_pending 设置为 1 可以消除问题,因为 Trident 一次只提交 1 个批次。但这不是可扩展的。设置为大于 1 的任何值都会使同一类别的元组(如果它们在两个连续批次中发出)在前一个事务提交之前被处理。
或者我应该为每个类别设置一个 spout 并将 max_spout_pending 设置为 1?
谢谢
【问题讨论】:
标签: java apache-storm trident