Skip navigation links
A B C D E F G I J L M N O P R S T U W 

A

addJoinKey(StageSpec, String, SparkCollection<Object>, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 
addJoinKey(StageSpec, String, SparkCollection<Object>, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.SparkPipelineRunner
 
addOutput(String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
addOutput(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
addOutput(String, OutputFormatProvider) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
addOutput(Output) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
aggregate(StageSpec, Integer, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
aggregate(StageSpec, Integer, StageStatisticsCollector) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
aggregate(StageSpec, Integer, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
AggregatorAggregateFunction<GROUP_KEY,GROUP_VAL,OUT> - Class in co.cask.cdap.etl.spark.function
Function that uses a BatchAggregator to perform the aggregate part of the aggregator.
AggregatorAggregateFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.AggregatorAggregateFunction
 
AggregatorGroupByFunction<GROUP_KEY,GROUP_VAL> - Class in co.cask.cdap.etl.spark.function
Function that uses a BatchAggregator to perform the groupBy part of the aggregator.
AggregatorGroupByFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.AggregatorGroupByFunction
 
AlertPassFilter - Class in co.cask.cdap.etl.spark.function
Filters a SparkCollection containing both output and errors to one that just contains alerts.
AlertPassFilter() - Constructor for class co.cask.cdap.etl.spark.function.AlertPassFilter
 

B

BasicSparkExecutionPluginContext - Class in co.cask.cdap.etl.spark.batch
Implementation of SparkExecutionPluginContext by delegating to JavaSparkExecutionContext.
BasicSparkExecutionPluginContext(JavaSparkExecutionContext, JavaSparkContext, DatasetContext, PipelineRuntime, StageSpec) - Constructor for class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
BasicSparkPluginContext - Class in co.cask.cdap.etl.spark.batch
Implementation of SparkPluginContext that delegates to a SparkContext.
BasicSparkPluginContext(SparkClientContext, PipelineRuntime, StageSpec, DatasetContext, Admin) - Constructor for class co.cask.cdap.etl.spark.batch.BasicSparkPluginContext
 
BatchSinkFunction<IN,OUT_KEY,OUT_VAL> - Class in co.cask.cdap.etl.spark.function
Function that uses a BatchSink to transform one object into a pair.
BatchSinkFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.BatchSinkFunction
 
BatchSourceFunction - Class in co.cask.cdap.etl.spark.function
Function that uses a BatchSource to transform a pair of objects into a single object.
BatchSourceFunction(PluginFunctionContext, int) - Constructor for class co.cask.cdap.etl.spark.function.BatchSourceFunction
 
BatchSparkPipelineDriver - Class in co.cask.cdap.etl.spark.batch
Batch Spark pipeline driver.
BatchSparkPipelineDriver() - Constructor for class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 

C

cache() - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
cache() - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
cache() - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
call(Tuple2<GROUP_KEY, Iterable<GROUP_VAL>>) - Method in class co.cask.cdap.etl.spark.function.AggregatorAggregateFunction
 
call(GROUP_VAL) - Method in class co.cask.cdap.etl.spark.function.AggregatorGroupByFunction
 
call(RecordInfo<Object>) - Method in class co.cask.cdap.etl.spark.function.AlertPassFilter
 
call(IN) - Method in class co.cask.cdap.etl.spark.function.BatchSinkFunction
 
call(Tuple2<Object, Object>) - Method in class co.cask.cdap.etl.spark.function.BatchSourceFunction
 
call(T) - Method in class co.cask.cdap.etl.spark.function.CountingFunction
 
call(RecordInfo<Object>) - Method in class co.cask.cdap.etl.spark.function.ErrorPassFilter
 
call(ErrorRecord<T>) - Method in class co.cask.cdap.etl.spark.function.ErrorTransformFunction
 
call(T) - Method in interface co.cask.cdap.etl.spark.function.FlatMapFunc
 
call(T) - Method in class co.cask.cdap.etl.spark.function.InitialJoinFunction
 
call(Tuple2<List<JoinElement<T>>, T>) - Method in class co.cask.cdap.etl.spark.function.JoinFlattenFunction
 
call(Tuple2<JOIN_KEY, List<JoinElement<INPUT_RECORD>>>) - Method in class co.cask.cdap.etl.spark.function.JoinMergeFunction
 
call(INPUT_RECORD) - Method in class co.cask.cdap.etl.spark.function.JoinOnFunction
 
call(Tuple2<List<JoinElement<T>>, Optional<T>>) - Method in class co.cask.cdap.etl.spark.function.LeftJoinFlattenFunction
 
call(T) - Method in class co.cask.cdap.etl.spark.function.MultiOutputTransformFunction
 
call(Tuple2<Optional<List<JoinElement<T>>>, Optional<T>>) - Method in class co.cask.cdap.etl.spark.function.OuterJoinFlattenFunction
 
call(RecordInfo<Object>) - Method in class co.cask.cdap.etl.spark.function.OutputPassFilter
 
call(T) - Method in interface co.cask.cdap.etl.spark.function.PairFlatMapFunc
 
call(T) - Method in class co.cask.cdap.etl.spark.function.TransformFunction
 
call(JavaRDD<T>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.ComputeTransformFunction
 
call(JavaRDD<T>) - Method in class co.cask.cdap.etl.spark.streaming.function.CountingTransformFunction
 
call(JavaPairRDD<GROUP_KEY, Iterable<GROUP_VAL>>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicAggregatorAggregate
 
call(JavaRDD<GROUP_VAL>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicAggregatorGroupBy
 
call(JavaPairRDD<JOIN_KEY, List<JoinElement<INPUT_RECORD>>>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicJoinMerge
 
call(JavaRDD<T>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicJoinOn
 
call(JavaRDD<T>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicTransform
 
call(JavaRDD<T>) - Method in class co.cask.cdap.etl.spark.streaming.function.preview.LimitingFunction
 
call(JavaRDD<Alert>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.StreamingAlertPublishFunction
 
call(JavaRDD<T>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.StreamingBatchSinkFunction
 
call(JavaRDD<T>, Time) - Method in class co.cask.cdap.etl.spark.streaming.function.StreamingSparkSinkFunction
 
call(T) - Method in class co.cask.cdap.etl.spark.streaming.function.WrapOutputTransformFunction
 
co.cask.cdap.etl.spark - package co.cask.cdap.etl.spark
 
co.cask.cdap.etl.spark.batch - package co.cask.cdap.etl.spark.batch
 
co.cask.cdap.etl.spark.function - package co.cask.cdap.etl.spark.function
 
co.cask.cdap.etl.spark.plugin - package co.cask.cdap.etl.spark.plugin
 
co.cask.cdap.etl.spark.streaming - package co.cask.cdap.etl.spark.streaming
 
co.cask.cdap.etl.spark.streaming.function - package co.cask.cdap.etl.spark.streaming.function
 
co.cask.cdap.etl.spark.streaming.function.preview - package co.cask.cdap.etl.spark.streaming.function.preview
 
CombinedEmitter<T> - Class in co.cask.cdap.etl.spark
An emitter used in Spark to collect all output and errors emitted.
CombinedEmitter(String) - Constructor for class co.cask.cdap.etl.spark.CombinedEmitter
 
Compat - Class in co.cask.cdap.etl.spark
Utility class to handle incompatibilities between Spark1 and Spark2.
compute(StageSpec, SparkCompute<T, U>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
compute(StageSpec, SparkCompute<T, U>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
compute(StageSpec, SparkCompute<T, U>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
ComputeTransformFunction<T,U> - Class in co.cask.cdap.etl.spark.streaming.function
Function used to implement a SparkCompute stage in a DStream.
ComputeTransformFunction(JavaSparkExecutionContext, StageSpec, SparkCompute<T, U>) - Constructor for class co.cask.cdap.etl.spark.streaming.function.ComputeTransformFunction
 
configure() - Method in class co.cask.cdap.etl.spark.batch.ETLSpark
 
configurePipeline(PipelineConfigurer) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkCompute
 
configurePipeline(PipelineConfigurer) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkSink
 
configurePipeline(PipelineConfigurer) - Method in class co.cask.cdap.etl.spark.plugin.WrappedStreamingSource
 
configurePipeline(PipelineConfigurer) - Method in class co.cask.cdap.etl.spark.plugin.WrappedWindower
 
convert(FlatMapFunc<T, R>) - Static method in class co.cask.cdap.etl.spark.Compat
 
convert(PairFlatMapFunc<T, K, V>) - Static method in class co.cask.cdap.etl.spark.Compat
 
CountingFunction<T> - Class in co.cask.cdap.etl.spark.function
Function that doesn't transform anything, but just emits counts for the number of records from that stage.
CountingFunction(String, Metrics, String, DataTracer) - Constructor for class co.cask.cdap.etl.spark.function.CountingFunction
 
CountingTransformFunction<T> - Class in co.cask.cdap.etl.spark.streaming.function
Function used to emit a metric for every item in an RDD.
CountingTransformFunction(String, Metrics, String, DataTracer) - Constructor for class co.cask.cdap.etl.spark.streaming.function.CountingTransformFunction
 
createBatchRuntimeContext() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
createPlugin() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
createSparkInterpreter() - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
createSparkInterpreter() - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
createStageMetrics() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
createTopic(String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
createTopic(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 

D

DatasetInfoTypeAdapter - Class in co.cask.cdap.etl.spark.batch
Type adapter for DatasetInfo
DatasetInfoTypeAdapter() - Constructor for class co.cask.cdap.etl.spark.batch.DatasetInfoTypeAdapter
 
DefaultStreamingContext - Class in co.cask.cdap.etl.spark.streaming
Default implementation of StreamingContext for Spark.
DefaultStreamingContext(StageSpec, JavaSparkExecutionContext, JavaStreamingContext) - Constructor for class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 
deleteTopic(String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
deserialize(JsonElement, Type, JsonDeserializationContext) - Method in class co.cask.cdap.etl.spark.batch.DatasetInfoTypeAdapter
 
deserialize(JsonElement, Type, JsonDeserializationContext) - Method in class co.cask.cdap.etl.spark.batch.InputFormatProviderTypeAdapter
 
deserialize(JsonElement, Type, JsonDeserializationContext) - Method in class co.cask.cdap.etl.spark.batch.OutputFormatProviderTypeAdapter
 
destroy() - Method in class co.cask.cdap.etl.spark.batch.ETLSpark
 
discardDataset(Dataset) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
discardDataset(Dataset) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
discardDataset(Dataset) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
DStreamCollection<T> - Class in co.cask.cdap.etl.spark.streaming
JavaDStream backed SparkCollection
DStreamCollection(JavaSparkExecutionContext, JavaDStream<T>) - Constructor for class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
DynamicAggregatorAggregate<GROUP_KEY,GROUP_VAL,OUT> - Class in co.cask.cdap.etl.spark.streaming.function
Serializable function that can be used to perform the aggregate part of an Aggregator.
DynamicAggregatorAggregate(DynamicDriverContext) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicAggregatorAggregate
 
DynamicAggregatorGroupBy<GROUP_KEY,GROUP_VAL> - Class in co.cask.cdap.etl.spark.streaming.function
Serializable function that can be used to perform the group by part of an Aggregator.
DynamicAggregatorGroupBy(DynamicDriverContext) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicAggregatorGroupBy
 
DynamicDriverContext - Class in co.cask.cdap.etl.spark.streaming
Serializable context that can be used to dynamically instantiate plugins from a driver context.
DynamicDriverContext() - Constructor for class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
DynamicDriverContext(StageSpec, JavaSparkExecutionContext, StageStatisticsCollector) - Constructor for class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
DynamicJoinMerge<JOIN_KEY,INPUT_RECORD,OUT> - Class in co.cask.cdap.etl.spark.streaming.function
Serializable function that can be used to perform the merge part of a Joiner.
DynamicJoinMerge(DynamicDriverContext) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicJoinMerge
 
DynamicJoinOn<JOIN_KEY,T> - Class in co.cask.cdap.etl.spark.streaming.function
Serializable function that can be used to perform the joinOn part of a Joiner.
DynamicJoinOn(DynamicDriverContext, String) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicJoinOn
 
DynamicSparkCompute<T,U> - Class in co.cask.cdap.etl.spark.streaming.function
This class is required to make sure that macro substitution occurs each time a pipeline is run instead of just the first time the pipeline is run.
DynamicSparkCompute(DynamicDriverContext, SparkCompute<T, U>) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicSparkCompute
 
DynamicTransform<T> - Class in co.cask.cdap.etl.spark.streaming.function
Serializable function that can be used to perform a flat map on a DStream.
DynamicTransform(DynamicDriverContext, boolean) - Constructor for class co.cask.cdap.etl.spark.streaming.function.DynamicTransform
 

E

emit(T) - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
emit(String, Object) - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
emitAlert(Map<String, String>) - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
emitError(InvalidEntry<T>) - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
ErrorPassFilter<T> - Class in co.cask.cdap.etl.spark.function
Filters a SparkCollection containing both output and errors to one that just contains errors.
ErrorPassFilter() - Constructor for class co.cask.cdap.etl.spark.function.ErrorPassFilter
 
ErrorTransformFunction<T,U> - Class in co.cask.cdap.etl.spark.function
Function that uses an ErrorTransform to perform a flatmap.
ErrorTransformFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.ErrorTransformFunction
 
ETLSpark - Class in co.cask.cdap.etl.spark.batch
Configures and sets up runs of BatchSparkPipelineDriver.
ETLSpark(BatchPhaseSpec) - Constructor for class co.cask.cdap.etl.spark.batch.ETLSpark
 
execute(TxRunnable) - Method in class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 
execute(int, TxRunnable) - Method in class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 

F

flatMap(FlatMapFunction<Tuple2<K, V>, T>) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
flatMap(StageSpec, FlatMapFunction<T, U>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
flatMap(StageSpec, FlatMapFunction<T, U>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
flatMap(FlatMapFunction<Tuple2<K, V>, T>) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
flatMap(StageSpec, FlatMapFunction<T, U>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
flatMap(FlatMapFunction<Tuple2<K, V>, T>) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
FlatMapFunc<T,R> - Interface in co.cask.cdap.etl.spark.function
A function that returns zero or more output records from each input record.
flatMapToPair(PairFlatMapFunction<T, K, V>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
flatMapToPair(PairFlatMapFunction<T, K, V>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
flatMapToPair(PairFlatMapFunction<T, K, V>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
foreachRDD(JavaDStream<T>, Function2<JavaRDD<T>, Time, Void>) - Static method in class co.cask.cdap.etl.spark.Compat
 
foreachRDD(JavaDStream<T>, Function2<JavaRDD<T>, Time, Void>) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 
fromDataset(String) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromDataset(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromDataset(String, Map<String, String>, Iterable<? extends Split>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromDataset(String) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromDataset(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromDataset(String, Map<String, String>, Iterable<? extends Split>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromStream(String) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromStream(String, long, long) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromStream(String, Class<V>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromStream(String, long, long, Class<V>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromStream(String, long, long, Class<? extends StreamEventDecoder<K, V>>, Class<K>, Class<V>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
fromStream(String) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromStream(String, long, long) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromStream(String, Class<V>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromStream(String, long, long, Class<V>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fromStream(String, long, long, Class<? extends StreamEventDecoder<K, V>>, Class<K>, Class<V>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
fullOuterJoin(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
fullOuterJoin(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
fullOuterJoin(JavaPairRDD<K, V1>, JavaPairRDD<K, V2>) - Static method in class co.cask.cdap.etl.spark.Compat
 
fullOuterJoin(JavaPairRDD<K, V1>, JavaPairRDD<K, V2>, int) - Static method in class co.cask.cdap.etl.spark.Compat
 
fullOuterJoin(SparkPairCollection<K, T>) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
fullOuterJoin(SparkPairCollection<K, T>, int) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
fullOuterJoin(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
fullOuterJoin(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
fullOuterJoin(JavaPairDStream<K, V1>, JavaPairDStream<K, V2>) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 
fullOuterJoin(JavaPairDStream<K, V1>, JavaPairDStream<K, V2>, int) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 

G

getDataset(String) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getDataset(String, String) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getDataset(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getDataset(String, String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getDataset(String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
getDataset(String, String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
getDataset(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
getDataset(String, String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
getDataset(String) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getDataset(String, String) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getDataset(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getDataset(String, String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getDataTracer() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
getDirectMessagePublisher() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
getEmitted() - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
getErrorRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
getInputRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
getLogicalStartTime() - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getLogicalStartTime() - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getMessageFetcher() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
getMessagePublisher() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
getOrCreate(String, Function0<JavaStreamingContext>) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 
getOutputRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
getPluginContext() - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getPluginContext() - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getPluginFunctionContext() - Method in class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
getRuntimeArguments() - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getRuntimeArguments() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
getRuntimeArguments() - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getSlideInterval() - Method in class co.cask.cdap.etl.spark.plugin.WrappedWindower
 
getSource(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 
getSource(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.SparkPipelineRunner
 
getSparkBatchSinkFactory() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceSinkFactoryInfo
 
getSparkBatchSourceFactory() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceSinkFactoryInfo
 
getSparkContext() - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
getSparkContext() - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
getSparkExecutionContext() - Method in class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 
getSparkExecutionContext() - Method in class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
getSparkStreamingContext() - Method in class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 
getStageName() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
getStagePartitions() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceSinkFactoryInfo
 
getStageSpec() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
getStageStatisticsCollector() - Method in class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
getStream(StreamingContext) - Method in class co.cask.cdap.etl.spark.plugin.WrappedStreamingSource
 
getTopicProperties(String) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
getUnderlying() - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
getUnderlying() - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
getUnderlying() - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
getUnderlying() - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
getUnderlying() - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
getUnderlying() - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
getWidth() - Method in class co.cask.cdap.etl.spark.plugin.WrappedWindower
 

I

incrementErrorRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
incrementInputRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
incrementOutputRecordCount() - Method in class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
initialize() - Method in class co.cask.cdap.etl.spark.batch.ETLSpark
 
initialize(SparkExecutionPluginContext) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkCompute
 
initialize(SparkExecutionPluginContext) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicSparkCompute
 
InitialJoinFunction<T> - Class in co.cask.cdap.etl.spark.function
Transforms an Object into a singleton list containing the JoinElement of that object.
InitialJoinFunction(String) - Constructor for class co.cask.cdap.etl.spark.function.InitialJoinFunction
 
InputFormatProviderTypeAdapter - Class in co.cask.cdap.etl.spark.batch
Type adapter for InputFormatProvider
InputFormatProviderTypeAdapter() - Constructor for class co.cask.cdap.etl.spark.batch.InputFormatProviderTypeAdapter
 
INSTANCE - Static variable in class co.cask.cdap.etl.spark.NoLookupProvider
 
isPreviewEnabled() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
isPreviewEnabled() - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 

J

join(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
join(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
join(SparkPairCollection<K, T>) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
join(SparkPairCollection<K, T>, int) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
join(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
join(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
JoinFlattenFunction<T> - Class in co.cask.cdap.etl.spark.function
Flattens the Tuple2 of list and object returned by a join into a single list.
JoinFlattenFunction(String) - Constructor for class co.cask.cdap.etl.spark.function.JoinFlattenFunction
 
JoinMergeFunction<JOIN_KEY,INPUT_RECORD,OUT> - Class in co.cask.cdap.etl.spark.function
Function that merges a join result using a BatchJoiner.
JoinMergeFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.JoinMergeFunction
 
JoinOnFunction<JOIN_KEY,INPUT_RECORD> - Class in co.cask.cdap.etl.spark.function
Function that uses a BatchJoiner to perform the joinOn part of the join.
JoinOnFunction(PluginFunctionContext, String) - Constructor for class co.cask.cdap.etl.spark.function.JoinOnFunction
 

L

LeftJoinFlattenFunction<T> - Class in co.cask.cdap.etl.spark.function
Flattens the Tuple2 of list and optional object returned by a left outer join into a single list.
LeftJoinFlattenFunction(String) - Constructor for class co.cask.cdap.etl.spark.function.LeftJoinFlattenFunction
 
leftOuterJoin(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
leftOuterJoin(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
leftOuterJoin(JavaPairRDD<K, V1>, JavaPairRDD<K, V2>) - Static method in class co.cask.cdap.etl.spark.Compat
 
leftOuterJoin(JavaPairRDD<K, V1>, JavaPairRDD<K, V2>, int) - Static method in class co.cask.cdap.etl.spark.Compat
 
leftOuterJoin(SparkPairCollection<K, T>) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
leftOuterJoin(SparkPairCollection<K, T>, int) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
leftOuterJoin(SparkPairCollection<K, T>) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
leftOuterJoin(SparkPairCollection<K, T>, int) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
leftOuterJoin(JavaPairDStream<K, V1>, JavaPairDStream<K, V2>) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 
leftOuterJoin(JavaPairDStream<K, V1>, JavaPairDStream<K, V2>, int) - Static method in class co.cask.cdap.etl.spark.StreamingCompat
 
LimitingFunction<T> - Class in co.cask.cdap.etl.spark.streaming.function.preview
Function used to limit the number of records
LimitingFunction(int) - Constructor for class co.cask.cdap.etl.spark.streaming.function.preview.LimitingFunction
 

M

mapValues(Function<V, T>) - Method in class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
mapValues(Function<V, T>) - Method in interface co.cask.cdap.etl.spark.SparkPairCollection
 
mapValues(Function<V, T>) - Method in class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
mergeJoinResults(StageSpec, SparkPairCollection<Object, List<JoinElement<Object>>>, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 
mergeJoinResults(StageSpec, SparkPairCollection<Object, List<JoinElement<Object>>>, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.SparkPipelineRunner
 
multiOutputTransform(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
multiOutputTransform(StageSpec, StageStatisticsCollector) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
multiOutputTransform(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
MultiOutputTransformFunction<T> - Class in co.cask.cdap.etl.spark.function
Function that uses a MultiOutputTransform to perform a flatmap.
MultiOutputTransformFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.MultiOutputTransformFunction
 

N

NoLookupProvider - Class in co.cask.cdap.etl.spark
A LookupProvider that doesn't work because lookups don't work in Spark.
NoLookupProvider() - Constructor for class co.cask.cdap.etl.spark.NoLookupProvider
 

O

onRunFinish(boolean, SparkPluginContext) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkSink
 
OuterJoinFlattenFunction<T> - Class in co.cask.cdap.etl.spark.function
Flattens the Tuple2 of optional list and optional object returned by a full outer join into a single list.
OuterJoinFlattenFunction(String) - Constructor for class co.cask.cdap.etl.spark.function.OuterJoinFlattenFunction
 
OutputFormatProviderTypeAdapter - Class in co.cask.cdap.etl.spark.batch
Type adapter for OutputFormatProvider
OutputFormatProviderTypeAdapter() - Constructor for class co.cask.cdap.etl.spark.batch.OutputFormatProviderTypeAdapter
 
OutputPassFilter<T> - Class in co.cask.cdap.etl.spark.function
Filters a SparkCollection containing both output and errors to one that just contains output from a specific port.
OutputPassFilter() - Constructor for class co.cask.cdap.etl.spark.function.OutputPassFilter
 
OutputPassFilter(String) - Constructor for class co.cask.cdap.etl.spark.function.OutputPassFilter
 

P

PairDStreamCollection<K,V> - Class in co.cask.cdap.etl.spark.streaming
JavaPairDStream backed SparkPairCollection
PairDStreamCollection(JavaSparkExecutionContext, JavaPairDStream<K, V>) - Constructor for class co.cask.cdap.etl.spark.streaming.PairDStreamCollection
 
PairFlatMapFunc<T,K,V> - Interface in co.cask.cdap.etl.spark.function
A function that returns zero or more key-value pair records from each input record.
PairRDDCollection<K,V> - Class in co.cask.cdap.etl.spark.batch
Implementation of SparkCollection that is backed by a JavaPairRDD.
PairRDDCollection(JavaSparkExecutionContext, JavaSparkContext, DatasetContext, SparkBatchSinkFactory, JavaPairRDD<K, V>) - Constructor for class co.cask.cdap.etl.spark.batch.PairRDDCollection
 
PluginFunctionContext - Class in co.cask.cdap.etl.spark.function
Serializable collection of objects that can be used in Spark closures to instantiate plugins.
PluginFunctionContext(StageSpec, JavaSparkExecutionContext, StageStatisticsCollector) - Constructor for class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
PluginFunctionContext(StageSpec, JavaSparkExecutionContext, Map<String, String>, long, StageStatisticsCollector) - Constructor for class co.cask.cdap.etl.spark.function.PluginFunctionContext
 
prepareRun(SparkPluginContext) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkSink
 
provide(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
provide(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.NoLookupProvider
 
provide(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
publishAlerts(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
publishAlerts(StageSpec, StageStatisticsCollector) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
publishAlerts(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 

R

RDDCollection<T> - Class in co.cask.cdap.etl.spark.batch
Implementation of SparkCollection that is backed by a JavaRDD.
RDDCollection(JavaSparkExecutionContext, JavaSparkContext, DatasetContext, SparkBatchSinkFactory, JavaRDD<T>) - Constructor for class co.cask.cdap.etl.spark.batch.RDDCollection
 
readExternal(ObjectInput) - Method in class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
registerLineage(String) - Method in class co.cask.cdap.etl.spark.streaming.DefaultStreamingContext
 
releaseDataset(Dataset) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
releaseDataset(Dataset) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
releaseDataset(Dataset) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
reset() - Method in class co.cask.cdap.etl.spark.CombinedEmitter
 
run(JavaSparkExecutionContext) - Method in class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 
run(DatasetContext) - Method in class co.cask.cdap.etl.spark.batch.BatchSparkPipelineDriver
 
run(JavaSparkExecutionContext) - Method in class co.cask.cdap.etl.spark.plugin.WrappedJavaSparkMain
 
run(SparkExecutionPluginContext, JavaRDD<IN>) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkSink
 
runPipeline(PipelinePhase, String, JavaSparkExecutionContext, Map<String, Integer>, PluginContext, Map<String, StageStatisticsCollector>) - Method in class co.cask.cdap.etl.spark.SparkPipelineRunner
 

S

saveAsDataset(JavaPairRDD<K, V>, String) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
saveAsDataset(JavaPairRDD<K, V>, String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkExecutionPluginContext
 
saveAsDataset(JavaPairRDD<K, V>, String) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
saveAsDataset(JavaPairRDD<K, V>, String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
serialize(DatasetInfo, Type, JsonSerializationContext) - Method in class co.cask.cdap.etl.spark.batch.DatasetInfoTypeAdapter
 
serialize(InputFormatProvider, Type, JsonSerializationContext) - Method in class co.cask.cdap.etl.spark.batch.InputFormatProviderTypeAdapter
 
serialize(OutputFormatProvider, Type, JsonSerializationContext) - Method in class co.cask.cdap.etl.spark.batch.OutputFormatProviderTypeAdapter
 
setInput(Input) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
setSparkConf(SparkConf) - Method in class co.cask.cdap.etl.spark.batch.BasicSparkPluginContext
 
SPARK_COMPAT - Static variable in class co.cask.cdap.etl.spark.Compat
 
SparkBatchRuntimeContext - Class in co.cask.cdap.etl.spark.batch
Default implementation of BatchRuntimeContext for spark contexts.
SparkBatchRuntimeContext(PipelineRuntime, StageSpec) - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchRuntimeContext
 
SparkBatchSinkContext - Class in co.cask.cdap.etl.spark.batch
Default implementation of BatchSinkContext for spark contexts.
SparkBatchSinkContext(SparkBatchSinkFactory, SparkClientContext, PipelineRuntime, DatasetContext, StageSpec) - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
SparkBatchSinkContext(SparkBatchSinkFactory, JavaSparkExecutionContext, DatasetContext, PipelineRuntime, StageSpec) - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchSinkContext
 
SparkBatchSinkFactory - Class in co.cask.cdap.etl.spark.batch
Handles writes to batch sinks.
SparkBatchSinkFactory() - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchSinkFactory
 
SparkBatchSourceContext - Class in co.cask.cdap.etl.spark.batch
Default implementation of BatchSourceContext for spark contexts.
SparkBatchSourceContext(SparkBatchSourceFactory, SparkClientContext, PipelineRuntime, DatasetContext, StageSpec) - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 
SparkBatchSourceSinkFactoryInfo - Class in co.cask.cdap.etl.spark.batch
Stores all the information of SparkBatchSinkFactory and stagePartitions
SparkBatchSourceSinkFactoryInfo(SparkBatchSourceFactory, SparkBatchSinkFactory, Map<String, Integer>) - Constructor for class co.cask.cdap.etl.spark.batch.SparkBatchSourceSinkFactoryInfo
 
SparkCollection<T> - Interface in co.cask.cdap.etl.spark
Abstraction over different types of spark collections with common shared operations on those collections.
SparkPairCollection<K,V> - Interface in co.cask.cdap.etl.spark
Abstraction over different types of spark pair collections with common shared operations on those collections.
SparkPipelinePluginContext - Class in co.cask.cdap.etl.spark.plugin
Wraps spark specific plugin types.
SparkPipelinePluginContext(PluginContext, Metrics, boolean, boolean) - Constructor for class co.cask.cdap.etl.spark.plugin.SparkPipelinePluginContext
 
SparkPipelineRunner - Class in co.cask.cdap.etl.spark
Base Spark program to run a Hydrator pipeline.
SparkPipelineRunner() - Constructor for class co.cask.cdap.etl.spark.SparkPipelineRunner
 
SparkPipelineRuntime - Class in co.cask.cdap.etl.spark
PipelineRuntime that can be created using Spark contexts.
SparkPipelineRuntime(SparkClientContext) - Constructor for class co.cask.cdap.etl.spark.SparkPipelineRuntime
 
SparkPipelineRuntime(JavaSparkExecutionContext) - Constructor for class co.cask.cdap.etl.spark.SparkPipelineRuntime
 
SparkPipelineRuntime(JavaSparkExecutionContext, long) - Constructor for class co.cask.cdap.etl.spark.SparkPipelineRuntime
 
SparkStageStatisticsCollector - Class in co.cask.cdap.etl.spark
Implementation of StageStatisticsCollector for batch spark pipelines.
SparkStageStatisticsCollector(JavaSparkContext) - Constructor for class co.cask.cdap.etl.spark.SparkStageStatisticsCollector
 
SparkStreamingExecutionContext - Class in co.cask.cdap.etl.spark.streaming
Implementation of SparkExecutionPluginContext by delegating to JavaSparkExecutionContext.
SparkStreamingExecutionContext(JavaSparkExecutionContext, JavaSparkContext, long, StageSpec) - Constructor for class co.cask.cdap.etl.spark.streaming.SparkStreamingExecutionContext
 
store(StageSpec, PairFlatMapFunction<T, Object, Object>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
store(StageSpec, SparkSink<T>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
store(StageSpec, PairFlatMapFunction<T, Object, Object>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
store(StageSpec, SparkSink<T>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
store(StageSpec, PairFlatMapFunction<T, Object, Object>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
store(StageSpec, SparkSink<T>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
StreamingAlertPublishFunction - Class in co.cask.cdap.etl.spark.streaming.function
Function used to publish alerts with a JavaDStream.
StreamingAlertPublishFunction(JavaSparkExecutionContext, StageSpec) - Constructor for class co.cask.cdap.etl.spark.streaming.function.StreamingAlertPublishFunction
 
StreamingBatchSinkFunction<T> - Class in co.cask.cdap.etl.spark.streaming.function
Function used to write a batch of data to a batch sink for use with a JavaDStream.
StreamingBatchSinkFunction(PairFlatMapFunction<T, Object, Object>, JavaSparkExecutionContext, StageSpec) - Constructor for class co.cask.cdap.etl.spark.streaming.function.StreamingBatchSinkFunction
 
StreamingCompat - Class in co.cask.cdap.etl.spark
Utility class to handle incompatibilities between Spark1 and Spark2 streaming.
StreamingSparkSinkFunction<T> - Class in co.cask.cdap.etl.spark.streaming.function
Function used to write a batch of data to a SparkSink for use with a JavaDStream.
StreamingSparkSinkFunction(JavaSparkExecutionContext, StageSpec) - Constructor for class co.cask.cdap.etl.spark.streaming.function.StreamingSparkSinkFunction
 

T

transform(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
transform(SparkExecutionPluginContext, JavaRDD<IN>) - Method in class co.cask.cdap.etl.spark.plugin.WrappedSparkCompute
 
transform(StageSpec, StageStatisticsCollector) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
transform(StageSpec, StageStatisticsCollector) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
transform(SparkExecutionPluginContext, JavaRDD<T>) - Method in class co.cask.cdap.etl.spark.streaming.function.DynamicSparkCompute
 
TransformFunction<T> - Class in co.cask.cdap.etl.spark.function
Function that uses a Transform to perform a flatmap.
TransformFunction(PluginFunctionContext) - Constructor for class co.cask.cdap.etl.spark.function.TransformFunction
 

U

union(SparkCollection<T>) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
union(SparkCollection<T>) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
union(SparkCollection<T>) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
updateTopic(String, Map<String, String>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSourceContext
 

W

window(StageSpec, Windower) - Method in class co.cask.cdap.etl.spark.batch.RDDCollection
 
window(StageSpec, Windower) - Method in interface co.cask.cdap.etl.spark.SparkCollection
 
window(StageSpec, Windower) - Method in class co.cask.cdap.etl.spark.streaming.DStreamCollection
 
WrapOutputTransformFunction<T> - Class in co.cask.cdap.etl.spark.streaming.function
Simply wraps all elements into a non-error, non-port RecordInfo.
WrapOutputTransformFunction(String) - Constructor for class co.cask.cdap.etl.spark.streaming.function.WrapOutputTransformFunction
 
WrappedJavaSparkMain - Class in co.cask.cdap.etl.spark.plugin
Wrapper around a JavaSparkMain that makes sure logging, classloading, and other pipeline capabilities are setup correctly.
WrappedJavaSparkMain(JavaSparkMain, Caller) - Constructor for class co.cask.cdap.etl.spark.plugin.WrappedJavaSparkMain
 
WrappedSparkCompute<IN,OUT> - Class in co.cask.cdap.etl.spark.plugin
Wrapper around a SparkCompute that makes sure logging, classloading, and other pipeline capabilities are setup correctly.
WrappedSparkCompute(SparkCompute<IN, OUT>, Caller) - Constructor for class co.cask.cdap.etl.spark.plugin.WrappedSparkCompute
 
WrappedSparkSink<IN> - Class in co.cask.cdap.etl.spark.plugin
Wrapper around a SparkCompute that makes sure logging, classloading, and other pipeline capabilities are setup correctly.
WrappedSparkSink(SparkSink<IN>, Caller) - Constructor for class co.cask.cdap.etl.spark.plugin.WrappedSparkSink
 
WrappedStreamingSource<T> - Class in co.cask.cdap.etl.spark.plugin
Wrapper around a Windower that makes sure logging, classloading, and other pipeline capabilities are setup correctly.
WrappedStreamingSource(StreamingSource<T>, Caller) - Constructor for class co.cask.cdap.etl.spark.plugin.WrappedStreamingSource
 
WrappedWindower - Class in co.cask.cdap.etl.spark.plugin
Wrapper around a Windower that makes sure logging, classloading, and other pipeline capabilities are setup correctly.
WrappedWindower(Windower, Caller) - Constructor for class co.cask.cdap.etl.spark.plugin.WrappedWindower
 
wrapUnknownPlugin(String, Object, Caller) - Method in class co.cask.cdap.etl.spark.plugin.SparkPipelinePluginContext
 
writeExternal(ObjectOutput) - Method in class co.cask.cdap.etl.spark.streaming.DynamicDriverContext
 
writeFromRDD(JavaPairRDD<K, V>, JavaSparkExecutionContext, String, Class<K>, Class<V>) - Method in class co.cask.cdap.etl.spark.batch.SparkBatchSinkFactory
 
A B C D E F G I J L M N O P R S T U W 
Skip navigation links

Copyright © 2018 Cask Data, Inc. Licensed under the Apache License, Version 2.0.