public class BasicSparkExecutionPluginContext extends AbstractTransformContext implements SparkExecutionPluginContext
SparkExecutionPluginContext by delegating to JavaSparkExecutionContext.arguments| Constructor and Description |
|---|
BasicSparkExecutionPluginContext(JavaSparkExecutionContext sec,
org.apache.spark.api.java.JavaSparkContext jsc,
DatasetContext datasetContext,
PipelineRuntime pipelineRuntime,
StageSpec stageSpec) |
| Modifier and Type | Method and Description |
|---|---|
SparkInterpreter |
createSparkInterpreter() |
void |
discardDataset(Dataset dataset) |
<K,V> org.apache.spark.api.java.JavaPairRDD<K,V> |
fromDataset(String datasetName) |
<K,V> org.apache.spark.api.java.JavaPairRDD<K,V> |
fromDataset(String datasetName,
Map<String,String> arguments) |
<K,V> org.apache.spark.api.java.JavaPairRDD<K,V> |
fromDataset(String datasetName,
Map<String,String> arguments,
Iterable<? extends Split> splits) |
org.apache.spark.api.java.JavaRDD<StreamEvent> |
fromStream(String streamName) |
<V> org.apache.spark.api.java.JavaPairRDD<Long,V> |
fromStream(String streamName,
Class<V> valueType) |
org.apache.spark.api.java.JavaRDD<StreamEvent> |
fromStream(String streamName,
long startTime,
long endTime) |
<K,V> org.apache.spark.api.java.JavaPairRDD<K,V> |
fromStream(String streamName,
long startTime,
long endTime,
Class<? extends co.cask.cdap.api.stream.StreamEventDecoder<K,V>> decoderClass,
Class<K> keyType,
Class<V> valueType) |
<V> org.apache.spark.api.java.JavaPairRDD<Long,V> |
fromStream(String streamName,
long startTime,
long endTime,
Class<V> valueType) |
<T extends Dataset> |
getDataset(String name) |
<T extends Dataset> |
getDataset(String name,
Map<String,String> arguments) |
<T extends Dataset> |
getDataset(String namespace,
String name) |
<T extends Dataset> |
getDataset(String namespace,
String name,
Map<String,String> arguments) |
long |
getLogicalStartTime() |
PluginContext |
getPluginContext() |
Map<String,String> |
getRuntimeArguments() |
org.apache.spark.api.java.JavaSparkContext |
getSparkContext() |
<T> Lookup<T> |
provide(String table,
Map<String,String> arguments) |
void |
releaseDataset(Dataset dataset) |
<K,V> void |
saveAsDataset(org.apache.spark.api.java.JavaPairRDD<K,V> rdd,
String datasetName) |
<K,V> void |
saveAsDataset(org.apache.spark.api.java.JavaPairRDD<K,V> rdd,
String datasetName,
Map<String,String> arguments) |
getArguments, getInputSchema, getInputSchemas, getMetrics, getNamespace, getOutputPortSchemas, getOutputSchema, getPipelineName, getPluginProperties, getPluginProperties, getServiceURL, getServiceURL, getStageName, loadPluginClass, newPluginInstanceclone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitgetArguments, getInputSchema, getInputSchemas, getMetrics, getNamespace, getOutputPortSchemas, getOutputSchema, getPipelineName, getPluginProperties, getPluginProperties, getStageName, loadPluginClass, newPluginInstancegetServiceURL, getServiceURLpublic BasicSparkExecutionPluginContext(JavaSparkExecutionContext sec, org.apache.spark.api.java.JavaSparkContext jsc, DatasetContext datasetContext, PipelineRuntime pipelineRuntime, StageSpec stageSpec)
public long getLogicalStartTime()
getLogicalStartTime in interface SparkExecutionPluginContextgetLogicalStartTime in interface StageContextgetLogicalStartTime in class AbstractStageContextpublic Map<String,String> getRuntimeArguments()
getRuntimeArguments in interface SparkExecutionPluginContextpublic <K,V> org.apache.spark.api.java.JavaPairRDD<K,V> fromDataset(String datasetName)
fromDataset in interface SparkExecutionPluginContextpublic <K,V> org.apache.spark.api.java.JavaPairRDD<K,V> fromDataset(String datasetName, Map<String,String> arguments)
fromDataset in interface SparkExecutionPluginContextpublic <K,V> org.apache.spark.api.java.JavaPairRDD<K,V> fromDataset(String datasetName, Map<String,String> arguments, @Nullable Iterable<? extends Split> splits)
fromDataset in interface SparkExecutionPluginContextpublic org.apache.spark.api.java.JavaRDD<StreamEvent> fromStream(String streamName)
fromStream in interface SparkExecutionPluginContextpublic org.apache.spark.api.java.JavaRDD<StreamEvent> fromStream(String streamName, long startTime, long endTime)
fromStream in interface SparkExecutionPluginContextpublic <V> org.apache.spark.api.java.JavaPairRDD<Long,V> fromStream(String streamName, Class<V> valueType)
fromStream in interface SparkExecutionPluginContextpublic <V> org.apache.spark.api.java.JavaPairRDD<Long,V> fromStream(String streamName, long startTime, long endTime, Class<V> valueType)
fromStream in interface SparkExecutionPluginContextpublic <K,V> org.apache.spark.api.java.JavaPairRDD<K,V> fromStream(String streamName, long startTime, long endTime, Class<? extends co.cask.cdap.api.stream.StreamEventDecoder<K,V>> decoderClass, Class<K> keyType, Class<V> valueType)
fromStream in interface SparkExecutionPluginContextpublic <K,V> void saveAsDataset(org.apache.spark.api.java.JavaPairRDD<K,V> rdd,
String datasetName)
saveAsDataset in interface SparkExecutionPluginContextpublic <K,V> void saveAsDataset(org.apache.spark.api.java.JavaPairRDD<K,V> rdd,
String datasetName,
Map<String,String> arguments)
saveAsDataset in interface SparkExecutionPluginContextpublic org.apache.spark.api.java.JavaSparkContext getSparkContext()
getSparkContext in interface SparkExecutionPluginContextpublic PluginContext getPluginContext()
getPluginContext in interface SparkExecutionPluginContextpublic SparkInterpreter createSparkInterpreter() throws IOException
createSparkInterpreter in interface SparkExecutionPluginContextIOExceptionpublic <T extends Dataset> T getDataset(String name) throws DatasetInstantiationException
getDataset in interface DatasetContextDatasetInstantiationExceptionpublic <T extends Dataset> T getDataset(String namespace, String name) throws DatasetInstantiationException
getDataset in interface DatasetContextDatasetInstantiationExceptionpublic <T extends Dataset> T getDataset(String name, Map<String,String> arguments) throws DatasetInstantiationException
getDataset in interface DatasetContextDatasetInstantiationExceptionpublic <T extends Dataset> T getDataset(String namespace, String name, Map<String,String> arguments) throws DatasetInstantiationException
getDataset in interface DatasetContextDatasetInstantiationExceptionpublic void releaseDataset(Dataset dataset)
releaseDataset in interface DatasetContextpublic void discardDataset(Dataset dataset)
discardDataset in interface DatasetContextpublic <T> Lookup<T> provide(String table, Map<String,String> arguments)
provide in interface LookupProviderprovide in class AbstractTransformContextCopyright © 2018 Cask Data, Inc. Licensed under the Apache License, Version 2.0.