T - type of objects in the collectionpublic class DStreamCollection<T> extends Object implements SparkCollection<T>
SparkCollection| Constructor and Description |
|---|
DStreamCollection(JavaSparkExecutionContext sec,
org.apache.spark.streaming.api.java.JavaDStream<T> stream) |
| Modifier and Type | Method and Description |
|---|---|
SparkCollection<RecordInfo<Object>> |
aggregate(StageSpec stageSpec,
Integer partitions,
StageStatisticsCollector collector) |
SparkCollection<T> |
cache() |
<U> SparkCollection<U> |
compute(StageSpec stageSpec,
SparkCompute<T,U> compute) |
<U> SparkCollection<U> |
flatMap(StageSpec stageSpec,
org.apache.spark.api.java.function.FlatMapFunction<T,U> function) |
<K,V> SparkPairCollection<K,V> |
flatMapToPair(org.apache.spark.api.java.function.PairFlatMapFunction<T,K,V> function) |
org.apache.spark.streaming.api.java.JavaDStream<T> |
getUnderlying() |
SparkCollection<RecordInfo<Object>> |
multiOutputTransform(StageSpec stageSpec,
StageStatisticsCollector collector) |
void |
publishAlerts(StageSpec stageSpec,
StageStatisticsCollector collector) |
void |
store(StageSpec stageSpec,
org.apache.spark.api.java.function.PairFlatMapFunction<T,Object,Object> sinkFunction) |
void |
store(StageSpec stageSpec,
SparkSink<T> sink) |
SparkCollection<RecordInfo<Object>> |
transform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<T> |
union(SparkCollection<T> other) |
SparkCollection<T> |
window(StageSpec stageSpec,
Windower windower) |
public DStreamCollection(JavaSparkExecutionContext sec, org.apache.spark.streaming.api.java.JavaDStream<T> stream)
public org.apache.spark.streaming.api.java.JavaDStream<T> getUnderlying()
getUnderlying in interface SparkCollection<T>public SparkCollection<T> cache()
cache in interface SparkCollection<T>public SparkCollection<T> union(SparkCollection<T> other)
union in interface SparkCollection<T>public SparkCollection<RecordInfo<Object>> transform(StageSpec stageSpec, StageStatisticsCollector collector)
transform in interface SparkCollection<T>public SparkCollection<RecordInfo<Object>> multiOutputTransform(StageSpec stageSpec, StageStatisticsCollector collector)
multiOutputTransform in interface SparkCollection<T>public <U> SparkCollection<U> flatMap(StageSpec stageSpec, org.apache.spark.api.java.function.FlatMapFunction<T,U> function)
flatMap in interface SparkCollection<T>public <K,V> SparkPairCollection<K,V> flatMapToPair(org.apache.spark.api.java.function.PairFlatMapFunction<T,K,V> function)
flatMapToPair in interface SparkCollection<T>public SparkCollection<RecordInfo<Object>> aggregate(StageSpec stageSpec, @Nullable Integer partitions, StageStatisticsCollector collector)
aggregate in interface SparkCollection<T>public <U> SparkCollection<U> compute(StageSpec stageSpec, SparkCompute<T,U> compute) throws Exception
compute in interface SparkCollection<T>Exceptionpublic void store(StageSpec stageSpec, org.apache.spark.api.java.function.PairFlatMapFunction<T,Object,Object> sinkFunction)
store in interface SparkCollection<T>public void store(StageSpec stageSpec, SparkSink<T> sink) throws Exception
store in interface SparkCollection<T>Exceptionpublic void publishAlerts(StageSpec stageSpec, StageStatisticsCollector collector) throws Exception
publishAlerts in interface SparkCollection<T>Exceptionpublic SparkCollection<T> window(StageSpec stageSpec, Windower windower)
window in interface SparkCollection<T>Copyright © 2018 Cask Data, Inc. Licensed under the Apache License, Version 2.0.