public class RocketMqSpout extends Object implements IRichSpout
Constructor and Description |
---|
RocketMqSpout(Properties properties)
RocketMqSpout Constructor.
|
Modifier and Type | Method and Description |
---|---|
void |
ack(Object msgId)
Storm has determined that the tuple emitted by this spout with the msgId identifier has been fully processed.
|
void |
activate()
Called when a spout has been activated out of a deactivated mode.
|
void |
close()
Called when an ISpout is going to be shutdown.
|
void |
deactivate()
Called when a spout has been deactivated.
|
void |
declareOutputFields(OutputFieldsDeclarer declarer)
Declare the output schema for all the streams of this topology.
|
void |
fail(Object msgId)
The tuple emitted by this spout with the msgId identifier has failed to be fully processed.
|
Map<String,Object> |
getComponentConfiguration()
Declare configuration specific to this component.
|
void |
nextTuple()
When this method is called, Storm is requesting that the Spout emit tuples to the output collector.
|
void |
open(Map<String,Object> conf,
TopologyContext context,
SpoutOutputCollector collector)
Called when a task for this component is initialized within a worker on the cluster.
|
protected boolean |
process(List<org.apache.rocketmq.common.message.MessageExt> msgs)
Process pushed messages.
|
public RocketMqSpout(Properties properties)
properties
- Properties Configpublic void open(Map<String,Object> conf, TopologyContext context, SpoutOutputCollector collector)
ISpout
This includes the:
open
in interface ISpout
conf
- The Storm configuration for this spout. This is the configuration provided to the topology merged in with cluster
configuration on this machine.context
- This object can be used to get information about this task's place within the topology, including the task id and
component id of this task, input and output information, etc.collector
- The collector is used to emit tuples from this spout. Tuples can be emitted at any time, including the open and
close methods. The collector is thread-safe and should be saved as an instance variable of this spout object.protected boolean process(List<org.apache.rocketmq.common.message.MessageExt> msgs)
msgs
- messagespublic void nextTuple()
ISpout
public void ack(Object msgId)
ISpout
public void fail(Object msgId)
ISpout
public void declareOutputFields(OutputFieldsDeclarer declarer)
IComponent
declareOutputFields
in interface IComponent
declarer
- this is used to declare output stream ids, output fields, and whether or not each output stream is a direct streampublic Map<String,Object> getComponentConfiguration()
IComponent
TopologyBuilder
getComponentConfiguration
in interface IComponent
public void close()
ISpout
The one context where close is guaranteed to be called is a topology is killed when running Storm in local mode.
public void activate()
ISpout
public void deactivate()
ISpout
deactivate
in interface ISpout
Copyright © 2023 The Apache Software Foundation. All rights reserved.