| Package | Description |
|---|---|
| co.cask.cdap.etl.spark | |
| co.cask.cdap.etl.spark.batch | |
| co.cask.cdap.etl.spark.streaming |
| Modifier and Type | Method and Description |
|---|---|
SparkCollection<RecordInfo<Object>> |
SparkCollection.aggregate(StageSpec stageSpec,
Integer partitions,
StageStatisticsCollector collector) |
SparkCollection<T> |
SparkCollection.cache() |
<U> SparkCollection<U> |
SparkCollection.compute(StageSpec stageSpec,
SparkCompute<T,U> compute) |
<T> SparkCollection<T> |
SparkPairCollection.flatMap(org.apache.spark.api.java.function.FlatMapFunction<scala.Tuple2<K,V>,T> function) |
<U> SparkCollection<U> |
SparkCollection.flatMap(StageSpec stageSpec,
org.apache.spark.api.java.function.FlatMapFunction<T,U> function) |
protected abstract SparkCollection<RecordInfo<Object>> |
SparkPipelineRunner.getSource(StageSpec stageSpec,
StageStatisticsCollector collector) |
protected abstract SparkCollection<Object> |
SparkPipelineRunner.mergeJoinResults(StageSpec stageSpec,
SparkPairCollection<Object,List<JoinElement<Object>>> joinedInputs,
StageStatisticsCollector collector) |
SparkCollection<RecordInfo<Object>> |
SparkCollection.multiOutputTransform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<RecordInfo<Object>> |
SparkCollection.transform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<T> |
SparkCollection.union(SparkCollection<T> other) |
SparkCollection<T> |
SparkCollection.window(StageSpec stageSpec,
Windower windower) |
| Modifier and Type | Method and Description |
|---|---|
protected abstract SparkPairCollection<Object,Object> |
SparkPipelineRunner.addJoinKey(StageSpec stageSpec,
String inputStageName,
SparkCollection<Object> inputCollection,
StageStatisticsCollector collector) |
SparkCollection<T> |
SparkCollection.union(SparkCollection<T> other) |
| Modifier and Type | Class and Description |
|---|---|
class |
RDDCollection<T>
Implementation of
SparkCollection that is backed by a JavaRDD. |
| Modifier and Type | Method and Description |
|---|---|
SparkCollection<RecordInfo<Object>> |
RDDCollection.aggregate(StageSpec stageSpec,
Integer partitions,
StageStatisticsCollector collector) |
SparkCollection<T> |
RDDCollection.cache() |
<U> SparkCollection<U> |
RDDCollection.compute(StageSpec stageSpec,
SparkCompute<T,U> compute) |
<T> SparkCollection<T> |
PairRDDCollection.flatMap(org.apache.spark.api.java.function.FlatMapFunction<scala.Tuple2<K,V>,T> function) |
<U> SparkCollection<U> |
RDDCollection.flatMap(StageSpec stageSpec,
org.apache.spark.api.java.function.FlatMapFunction<T,U> function) |
protected SparkCollection<RecordInfo<Object>> |
BatchSparkPipelineDriver.getSource(StageSpec stageSpec,
StageStatisticsCollector collector) |
protected SparkCollection<Object> |
BatchSparkPipelineDriver.mergeJoinResults(StageSpec stageSpec,
SparkPairCollection<Object,List<JoinElement<Object>>> joinedInputs,
StageStatisticsCollector collector) |
SparkCollection<RecordInfo<Object>> |
RDDCollection.multiOutputTransform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<RecordInfo<Object>> |
RDDCollection.transform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<T> |
RDDCollection.union(SparkCollection<T> other) |
SparkCollection<T> |
RDDCollection.window(StageSpec stageSpec,
Windower windower) |
| Modifier and Type | Method and Description |
|---|---|
protected SparkPairCollection<Object,Object> |
BatchSparkPipelineDriver.addJoinKey(StageSpec stageSpec,
String inputStageName,
SparkCollection<Object> inputCollection,
StageStatisticsCollector collector) |
SparkCollection<T> |
RDDCollection.union(SparkCollection<T> other) |
| Modifier and Type | Class and Description |
|---|---|
class |
DStreamCollection<T>
JavaDStream backed
SparkCollection |
| Modifier and Type | Method and Description |
|---|---|
SparkCollection<RecordInfo<Object>> |
DStreamCollection.aggregate(StageSpec stageSpec,
Integer partitions,
StageStatisticsCollector collector) |
SparkCollection<T> |
DStreamCollection.cache() |
<U> SparkCollection<U> |
DStreamCollection.compute(StageSpec stageSpec,
SparkCompute<T,U> compute) |
<T> SparkCollection<T> |
PairDStreamCollection.flatMap(org.apache.spark.api.java.function.FlatMapFunction<scala.Tuple2<K,V>,T> function) |
<U> SparkCollection<U> |
DStreamCollection.flatMap(StageSpec stageSpec,
org.apache.spark.api.java.function.FlatMapFunction<T,U> function) |
SparkCollection<RecordInfo<Object>> |
DStreamCollection.multiOutputTransform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<RecordInfo<Object>> |
DStreamCollection.transform(StageSpec stageSpec,
StageStatisticsCollector collector) |
SparkCollection<T> |
DStreamCollection.union(SparkCollection<T> other) |
SparkCollection<T> |
DStreamCollection.window(StageSpec stageSpec,
Windower windower) |
| Modifier and Type | Method and Description |
|---|---|
SparkCollection<T> |
DStreamCollection.union(SparkCollection<T> other) |
Copyright © 2018 Cask Data, Inc. Licensed under the Apache License, Version 2.0.