Class HdfsSpout

All Implemented Interfaces:
Serializable, ISpout, IComponent, IRichSpout

public class HdfsSpout extends BaseRichSpout
See Also:
  • Constructor Details

    • HdfsSpout

      public HdfsSpout()
  • Method Details

    • setHdfsUri

      public HdfsSpout setHdfsUri(String hdfsUri)
    • setReaderType

      public HdfsSpout setReaderType(String readerType)
    • setSourceDir

      public HdfsSpout setSourceDir(String sourceDir)
    • setArchiveDir

      public HdfsSpout setArchiveDir(String archiveDir)
    • setBadFilesDir

      public HdfsSpout setBadFilesDir(String badFilesDir)
    • setLockDir

      public HdfsSpout setLockDir(String lockDir)
    • setCommitFrequencyCount

      public HdfsSpout setCommitFrequencyCount(int commitFrequencyCount)
    • setCommitFrequencySec

      public HdfsSpout setCommitFrequencySec(int commitFrequencySec)
    • setMaxOutstanding

      public HdfsSpout setMaxOutstanding(int maxOutstanding)
    • setLockTimeoutSec

      public HdfsSpout setLockTimeoutSec(int lockTimeoutSec)
    • setClocksInSync

      public HdfsSpout setClocksInSync(boolean clocksInSync)
    • setIgnoreSuffix

      public HdfsSpout setIgnoreSuffix(String ignoreSuffix)
    • withOutputFields

      public HdfsSpout withOutputFields(String... fields)
      Output field names. Number of fields depends upon the reader type
    • withConfigKey

      public HdfsSpout withConfigKey(String configKey)
      set key name under which HDFS options are placed. (similar to HDFS bolt). default key name is 'hdfs.config'
    • withOutputStream

      public HdfsSpout withOutputStream(String streamName)
      Set output stream name.
    • getLockDirPath

      public org.apache.hadoop.fs.Path getLockDirPath()
    • getCollector

      public SpoutOutputCollector getCollector()
    • nextTuple

      public void nextTuple()
    • emitData

      protected void emitData(List<Object> tuple, org.apache.storm.hdfs.spout.HdfsSpout.MessageId id)
    • open

      public void open(Map<String,Object> conf, TopologyContext context, SpoutOutputCollector collector)
    • close

      public void close()
      Specified by:
      close in interface ISpout
      Overrides:
      close in class BaseRichSpout
    • ack

      public void ack(Object msgId)
      Specified by:
      ack in interface ISpout
      Overrides:
      ack in class BaseRichSpout
    • fail

      public void fail(Object msgId)
      Specified by:
      fail in interface ISpout
      Overrides:
      fail in class BaseRichSpout
    • declareOutputFields

      public void declareOutputFields(OutputFieldsDeclarer declarer)