public class BatchSparkPipelineDriver extends SparkPipelineRunner implements JavaSparkMain, TxRunnable
| Constructor and Description |
|---|
BatchSparkPipelineDriver() |
| Modifier and Type | Method and Description |
|---|---|
protected SparkPairCollection<Object,Object> |
addJoinKey(StageSpec stageSpec,
String inputStageName,
SparkCollection<Object> inputCollection,
StageStatisticsCollector collector) |
protected SparkCollection<RecordInfo<Object>> |
getSource(StageSpec stageSpec,
StageStatisticsCollector collector) |
protected SparkCollection<Object> |
mergeJoinResults(StageSpec stageSpec,
SparkPairCollection<Object,List<JoinElement<Object>>> joinedInputs,
StageStatisticsCollector collector) |
void |
run(DatasetContext context) |
void |
run(JavaSparkExecutionContext sec) |
runPipelineprotected SparkCollection<RecordInfo<Object>> getSource(StageSpec stageSpec, StageStatisticsCollector collector)
getSource in class SparkPipelineRunnerprotected SparkPairCollection<Object,Object> addJoinKey(StageSpec stageSpec, String inputStageName, SparkCollection<Object> inputCollection, StageStatisticsCollector collector) throws Exception
addJoinKey in class SparkPipelineRunnerExceptionprotected SparkCollection<Object> mergeJoinResults(StageSpec stageSpec, SparkPairCollection<Object,List<JoinElement<Object>>> joinedInputs, StageStatisticsCollector collector) throws Exception
mergeJoinResults in class SparkPipelineRunnerExceptionpublic void run(JavaSparkExecutionContext sec) throws Exception
run in interface JavaSparkMainExceptionpublic void run(DatasetContext context) throws Exception
run in interface TxRunnableExceptionCopyright © 2018 Cask Data, Inc. Licensed under the Apache License, Version 2.0.