Apache Storm:如何在Bolt中转换并创建新的Tuple?
在Apache Storm里,Tuple是不可变对象,所以你没法直接修改原有Tuple的values,必须创建一个新的Tuple实例来传递给doExecute。这里有个符合Storm惯用风格的实现方式,不用自己手动复制修改,而是利用Storm提供的内部Tuple实现类来构建:
首先,你需要在Bolt的prepare方法里保存好TopologyContext的引用,因为创建新Tuple需要用到它:
private lateinit var context: TopologyContext private lateinit var collector: OutputCollector override fun prepare(config: MutableMap<Any?, Any>?, context: TopologyContext?, collector: OutputCollector?) { super.prepare(config, context, collector) this.context = context!! this.collector = collector!! }
然后在execute方法里,用TupleImpl(Storm的标准Tuple实现类)来构建新的Tuple,同时保留原Tuple的源元数据(组件名、任务ID、流ID),这样能保证Tuple的上下文完整性:
override fun execute(input: Tuple?) { val filteredValues = input?.values?.filterIsInstance(MyClass::class.java) ?: emptyList() // 构建新Tuple,复用原Tuple的源信息 val newTuple = TupleImpl( context, filteredValues, input?.sourceComponent, input?.sourceTaskId, input?.sourceStreamId ) if (doExecute(newTuple)) { collector.ack(input) } else { collector.fail(input) } }
为什么这么做?
- Storm的Tuple设计为不可变,所以直接修改原有Tuple的values是不可能的,必须创建新实例
TupleImpl是Storm官方提供的Tuple实现,用它来构建新Tuple是最符合Storm生态的方式- 保留原Tuple的源元数据,能避免一些潜在的问题(比如后续如果有依赖Tuple来源的逻辑,不会出问题)
你之前猜测的“复制输入Tuple并修改values”的思路其实不可行,因为Tuple的values属性是只读的,没有提供修改的方法。上面的实现就是Storm场景下处理这类需求的惯用方案。
内容的提问来源于stack exchange,提问作者breezymri
相关产品推荐
相关产品推荐

