public class JavaInputDStream<T> extends JavaDStream<T>
InputDStream
.Constructor and Description |
---|
JavaInputDStream(InputDStream<T> inputDStream,
scala.reflect.ClassTag<T> classTag) |
Modifier and Type | Method and Description |
---|---|
static JavaDStream<T> |
cache() |
static DStream<T> |
checkpoint(Duration interval) |
scala.reflect.ClassTag<T> |
classTag() |
static JavaRDD<T> |
compute(Time validTime) |
static StreamingContext |
context() |
static JavaDStream<Long> |
count() |
static JavaPairDStream<T,Long> |
countByValue() |
static JavaPairDStream<T,Long> |
countByValue(int numPartitions) |
static JavaPairDStream<T,Long> |
countByValueAndWindow(Duration windowDuration,
Duration slideDuration) |
static JavaPairDStream<T,Long> |
countByValueAndWindow(Duration windowDuration,
Duration slideDuration,
int numPartitions) |
static JavaDStream<Long> |
countByWindow(Duration windowDuration,
Duration slideDuration) |
static DStream<T> |
dstream() |
static JavaDStream<T> |
filter(Function<T,Boolean> f) |
static <U> JavaDStream<U> |
flatMap(FlatMapFunction<T,U> f) |
static <K2,V2> JavaPairDStream<K2,V2> |
flatMapToPair(PairFlatMapFunction<T,K2,V2> f) |
static void |
foreachRDD(VoidFunction<R> foreachFunc) |
static void |
foreachRDD(VoidFunction2<R,Time> foreachFunc) |
static <T> JavaInputDStream<T> |
fromInputDStream(InputDStream<T> inputDStream,
scala.reflect.ClassTag<T> evidence$1)
Convert a scala
InputDStream to a Java-friendly
JavaInputDStream . |
static JavaDStream<java.util.List<T>> |
glom() |
InputDStream<T> |
inputDStream() |
static <R> JavaDStream<R> |
map(Function<T,R> f) |
static <U> JavaDStream<U> |
mapPartitions(FlatMapFunction<java.util.Iterator<T>,U> f) |
static <K2,V2> JavaPairDStream<K2,V2> |
mapPartitionsToPair(PairFlatMapFunction<java.util.Iterator<T>,K2,V2> f) |
static <K2,V2> JavaPairDStream<K2,V2> |
mapToPair(PairFunction<T,K2,V2> f) |
static JavaDStream<T> |
persist() |
static JavaDStream<T> |
persist(StorageLevel storageLevel) |
static void |
print() |
static void |
print(int num) |
static JavaDStream<T> |
reduce(Function2<T,T,T> f) |
static JavaDStream<T> |
reduceByWindow(Function2<T,T,T> reduceFunc,
Duration windowDuration,
Duration slideDuration) |
static JavaDStream<T> |
reduceByWindow(Function2<T,T,T> reduceFunc,
Function2<T,T,T> invReduceFunc,
Duration windowDuration,
Duration slideDuration) |
static JavaDStream<T> |
repartition(int numPartitions) |
static JavaDStream<Long> |
scalaIntToJavaLong(DStream<Object> in) |
static java.util.List<R> |
slice(Time fromTime,
Time toTime) |
static <U> JavaDStream<U> |
transform(Function<R,JavaRDD<U>> transformFunc) |
static <U> JavaDStream<U> |
transform(Function2<R,Time,JavaRDD<U>> transformFunc) |
static <K2,V2> JavaPairDStream<K2,V2> |
transformToPair(Function<R,JavaPairRDD<K2,V2>> transformFunc) |
static <K2,V2> JavaPairDStream<K2,V2> |
transformToPair(Function2<R,Time,JavaPairRDD<K2,V2>> transformFunc) |
static <U,W> JavaDStream<W> |
transformWith(JavaDStream<U> other,
Function3<R,JavaRDD<U>,Time,JavaRDD<W>> transformFunc) |
static <K2,V2,W> JavaDStream<W> |
transformWith(JavaPairDStream<K2,V2> other,
Function3<R,JavaPairRDD<K2,V2>,Time,JavaRDD<W>> transformFunc) |
static <U,K2,V2> JavaPairDStream<K2,V2> |
transformWithToPair(JavaDStream<U> other,
Function3<R,JavaRDD<U>,Time,JavaPairRDD<K2,V2>> transformFunc) |
static <K2,V2,K3,V3> |
transformWithToPair(JavaPairDStream<K2,V2> other,
Function3<R,JavaPairRDD<K2,V2>,Time,JavaPairRDD<K3,V3>> transformFunc) |
static JavaDStream<T> |
union(JavaDStream<T> that) |
static JavaDStream<T> |
window(Duration windowDuration) |
static JavaDStream<T> |
window(Duration windowDuration,
Duration slideDuration) |
static JavaRDD<T> |
wrapRDD(RDD<T> rdd) |
cache, compute, dstream, filter, fromDStream, persist, persist, repartition, union, window, window, wrapRDD
equals, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
checkpoint, context, count, countByValue, countByValue, countByValueAndWindow, countByValueAndWindow, countByWindow, flatMap, flatMapToPair, foreachRDD, foreachRDD, glom, map, mapPartitions, mapPartitionsToPair, mapToPair, print, print, reduce, reduceByWindow, reduceByWindow, scalaIntToJavaLong, slice, transform, transform, transformToPair, transformToPair, transformWith, transformWith, transformWithToPair, transformWithToPair
public JavaInputDStream(InputDStream<T> inputDStream, scala.reflect.ClassTag<T> classTag)
public static <T> JavaInputDStream<T> fromInputDStream(InputDStream<T> inputDStream, scala.reflect.ClassTag<T> evidence$1)
InputDStream
to a Java-friendly
JavaInputDStream
.inputDStream
- (undocumented)evidence$1
- (undocumented)public static JavaDStream<Long> scalaIntToJavaLong(DStream<Object> in)
public static void print()
public static void print(int num)
public static JavaDStream<Long> count()
public static JavaPairDStream<T,Long> countByValue()
public static JavaPairDStream<T,Long> countByValue(int numPartitions)
public static JavaDStream<Long> countByWindow(Duration windowDuration, Duration slideDuration)
public static JavaPairDStream<T,Long> countByValueAndWindow(Duration windowDuration, Duration slideDuration)
public static JavaPairDStream<T,Long> countByValueAndWindow(Duration windowDuration, Duration slideDuration, int numPartitions)
public static JavaDStream<java.util.List<T>> glom()
public static StreamingContext context()
public static <R> JavaDStream<R> map(Function<T,R> f)
public static <K2,V2> JavaPairDStream<K2,V2> mapToPair(PairFunction<T,K2,V2> f)
public static <U> JavaDStream<U> flatMap(FlatMapFunction<T,U> f)
public static <K2,V2> JavaPairDStream<K2,V2> flatMapToPair(PairFlatMapFunction<T,K2,V2> f)
public static <U> JavaDStream<U> mapPartitions(FlatMapFunction<java.util.Iterator<T>,U> f)
public static <K2,V2> JavaPairDStream<K2,V2> mapPartitionsToPair(PairFlatMapFunction<java.util.Iterator<T>,K2,V2> f)
public static JavaDStream<T> reduce(Function2<T,T,T> f)
public static JavaDStream<T> reduceByWindow(Function2<T,T,T> reduceFunc, Duration windowDuration, Duration slideDuration)
public static JavaDStream<T> reduceByWindow(Function2<T,T,T> reduceFunc, Function2<T,T,T> invReduceFunc, Duration windowDuration, Duration slideDuration)
public static void foreachRDD(VoidFunction<R> foreachFunc)
public static void foreachRDD(VoidFunction2<R,Time> foreachFunc)
public static <U> JavaDStream<U> transform(Function<R,JavaRDD<U>> transformFunc)
public static <U> JavaDStream<U> transform(Function2<R,Time,JavaRDD<U>> transformFunc)
public static <K2,V2> JavaPairDStream<K2,V2> transformToPair(Function<R,JavaPairRDD<K2,V2>> transformFunc)
public static <K2,V2> JavaPairDStream<K2,V2> transformToPair(Function2<R,Time,JavaPairRDD<K2,V2>> transformFunc)
public static <U,W> JavaDStream<W> transformWith(JavaDStream<U> other, Function3<R,JavaRDD<U>,Time,JavaRDD<W>> transformFunc)
public static <U,K2,V2> JavaPairDStream<K2,V2> transformWithToPair(JavaDStream<U> other, Function3<R,JavaRDD<U>,Time,JavaPairRDD<K2,V2>> transformFunc)
public static <K2,V2,W> JavaDStream<W> transformWith(JavaPairDStream<K2,V2> other, Function3<R,JavaPairRDD<K2,V2>,Time,JavaRDD<W>> transformFunc)
public static <K2,V2,K3,V3> JavaPairDStream<K3,V3> transformWithToPair(JavaPairDStream<K2,V2> other, Function3<R,JavaPairRDD<K2,V2>,Time,JavaPairRDD<K3,V3>> transformFunc)
public static DStream<T> dstream()
public static JavaDStream<T> filter(Function<T,Boolean> f)
public static JavaDStream<T> cache()
public static JavaDStream<T> persist()
public static JavaDStream<T> persist(StorageLevel storageLevel)
public static JavaDStream<T> window(Duration windowDuration)
public static JavaDStream<T> window(Duration windowDuration, Duration slideDuration)
public static JavaDStream<T> union(JavaDStream<T> that)
public static JavaDStream<T> repartition(int numPartitions)
public InputDStream<T> inputDStream()
public scala.reflect.ClassTag<T> classTag()
classTag
in interface JavaDStreamLike<T,JavaDStream<T>,JavaRDD<T>>
classTag
in class JavaDStream<T>