public class FlumePollingInputDStream<T> extends ReceiverInputDStream<SparkFlumeEvent>
ReceiverInputDStream
that can be used to read data from several Flume agents running
SparkSink
s.Constructor and Description |
---|
FlumePollingInputDStream(StreamingContext _ssc,
scala.collection.Seq<java.net.InetSocketAddress> addresses,
int maxBatchSize,
int parallelism,
StorageLevel storageLevel,
scala.reflect.ClassTag<T> evidence$1) |
Modifier and Type | Method and Description |
---|---|
scala.collection.Seq<java.net.InetSocketAddress> |
addresses() |
Receiver<SparkFlumeEvent> |
getReceiver()
Gets the receiver object that will be sent to the worker nodes
to receive data.
|
int |
maxBatchSize() |
int |
parallelism() |
clearMetadata, compute, id, start, stop
dependencies, isTimeValid, lastValidTime, slideDuration
cache, checkpoint, checkpointData, checkpointDuration, clearCheckpointData, context, count, countByValue, countByValueAndWindow, countByWindow, creationSite, filter, flatMap, foreach, foreach, foreachRDD, foreachRDD, generatedRDDs, generateJob, getCreationSite, getOrCompute, glom, graph, initialize, isInitialized, map, mapPartitions, mustCheckpoint, parentRememberDuration, persist, persist, print, reduce, reduceByWindow, reduceByWindow, register, remember, rememberDuration, repartition, restoreCheckpointData, saveAsObjectFiles, saveAsTextFiles, setContext, setGraph, slice, slice, ssc, storageLevel, transform, transform, transformWith, transformWith, union, updateCheckpointData, validate, window, window, zeroTime
equals, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
initializeIfNecessary, initializeLogging, isTraceEnabled, log_, log, logDebug, logDebug, logError, logError, logInfo, logInfo, logName, logTrace, logTrace, logWarning, logWarning
public FlumePollingInputDStream(StreamingContext _ssc, scala.collection.Seq<java.net.InetSocketAddress> addresses, int maxBatchSize, int parallelism, StorageLevel storageLevel, scala.reflect.ClassTag<T> evidence$1)
public scala.collection.Seq<java.net.InetSocketAddress> addresses()
public int maxBatchSize()
public int parallelism()
public Receiver<SparkFlumeEvent> getReceiver()
ReceiverInputDStream
getReceiver
in class ReceiverInputDStream<SparkFlumeEvent>