public class MockSource extends StreamingSource<co.cask.cdap.api.data.format.StructuredRecord>
| Modifier and Type | Class and Description |
|---|---|
static class |
MockSource.Conf
Config for mock source.
|
| Modifier and Type | Field and Description |
|---|---|
static PluginClass |
PLUGIN_CLASS |
PLUGIN_TYPE| Constructor and Description |
|---|
MockSource(MockSource.Conf conf) |
| Modifier and Type | Method and Description |
|---|---|
void |
configurePipeline(PipelineConfigurer pipelineConfigurer) |
static ETLPlugin |
getPlugin(co.cask.cdap.api.data.schema.Schema schema,
List<co.cask.cdap.api.data.format.StructuredRecord> records) |
static ETLPlugin |
getPlugin(co.cask.cdap.api.data.schema.Schema schema,
List<co.cask.cdap.api.data.format.StructuredRecord> records,
Long intervalMillis) |
org.apache.spark.streaming.api.java.JavaDStream<co.cask.cdap.api.data.format.StructuredRecord> |
getStream(StreamingContext context) |
getRequiredExecutorspublic static final PluginClass PLUGIN_CLASS
public MockSource(MockSource.Conf conf)
public void configurePipeline(PipelineConfigurer pipelineConfigurer)
configurePipeline in interface PipelineConfigurableconfigurePipeline in class StreamingSource<co.cask.cdap.api.data.format.StructuredRecord>public org.apache.spark.streaming.api.java.JavaDStream<co.cask.cdap.api.data.format.StructuredRecord> getStream(StreamingContext context) throws Exception
getStream in class StreamingSource<co.cask.cdap.api.data.format.StructuredRecord>Exceptionpublic static ETLPlugin getPlugin(co.cask.cdap.api.data.schema.Schema schema, List<co.cask.cdap.api.data.format.StructuredRecord> records) throws IOException
IOExceptionpublic static ETLPlugin getPlugin(co.cask.cdap.api.data.schema.Schema schema, List<co.cask.cdap.api.data.format.StructuredRecord> records, Long intervalMillis) throws IOException
IOExceptionCopyright © 2017 Cask Data, Inc. Licensed under the Apache License, Version 2.0.