public class HadoopFlow extends BaseFlow<JobConf>
Flow
.
HadoopFlow must be created through a FlowConnector
sub-class instance.
If classpath paths are provided on the FlowDef
, the Hadoop distributed cache mechanism will be used
to augment the remote classpath.
Any path elements that are relative will be uploaded to HDFS, and the HDFS URI will be used on the JobConf. Note
all paths are added as "files" to the JobConf, not archives, so they aren't needlessly uncompressed cluster side.FlowConnector
BaseFlow.FlowHolder
flowStats, sinks, sources, stop, stopJobsOnExit, thread
CASCADING_FLOW_ID
NULL
Modifier | Constructor and Description |
---|---|
protected |
HadoopFlow() |
|
HadoopFlow(PlatformInfo platformInfo,
Map<Object,Object> properties,
JobConf jobConf,
FlowDef flowDef) |
protected |
HadoopFlow(PlatformInfo platformInfo,
Map<Object,Object> properties,
JobConf jobConf,
String name,
Map<String,String> flowDescriptor) |
Modifier and Type | Method and Description |
---|---|
JobConf |
getConfig() |
Map<Object,Object> |
getConfigAsProperties() |
JobConf |
getConfigCopy() |
FlowProcess<JobConf> |
getFlowProcess() |
protected int |
getMaxNumParallelSteps() |
String |
getProperty(String key)
Method getProperty returns the value associated with the given key from the underlying properties system.
|
protected long |
getTotalSliceCPUMilliSeconds() |
protected void |
initConfig(Map<Object,Object> properties,
JobConf parentConfig) |
protected void |
initFromProperties(Map<Object,Object> properties) |
protected void |
internalClean(boolean stop) |
protected void |
internalShutdown() |
protected void |
internalStart() |
boolean |
isPreserveTemporaryFiles()
Method isPreserveTemporaryFiles returns false if temporary files will be cleaned when this Flow completes.
|
protected JobConf |
newConfig(JobConf defaultConfig) |
protected void |
setConfigProperty(JobConf config,
Object key,
Object value) |
boolean |
stepsAreLocal() |
addListener, addPlannerProperties, addStepListener, areSinksStale, areSourcesNewer, cleanup, complete, createConfig, createFlowThread, deleteCheckpointsIfNotUpdate, deleteCheckpointsIfReplace, deleteSinks, deleteSinksIfNotUpdate, deleteSinksIfReplace, deleteTrapsIfNotUpdate, deleteTrapsIfReplace, fireOnCompleted, fireOnStarting, fireOnStopping, fireOnThrowable, fireOnThrowable, getCascadeID, getCascadingServices, getCheckpointNames, getCheckpoints, getCheckpointsCollection, getClassPath, getFieldsFor, getFlowDescriptor, getFlowSession, getFlowSkipStrategy, getFlowStats, getFlowSteps, getFlowStepStrategy, getHolder, getID, getName, getPlannerInfo, getPlatformInfo, getRunID, getSink, getSink, getSinkModified, getSinkNames, getSinks, getSinksCollection, getSource, getSourceNames, getSources, getSourcesCollection, getSpawnStrategy, getStats, getSubmitPriority, getTags, getTrapNames, getTraps, getTrapsCollection, handleExecutorShutdown, hasListeners, hasStepListeners, initialize, initializeNewJobsMap, initSteps, internalStopAllJobs, isDebugEnabled, isInfoEnabled, isSkipFlow, isStopJobsOnExit, logDebug, logError, logError, logInfo, logWarn, logWarn, logWarn, openSink, openSink, openSource, openSource, openTapForRead, openTapForWrite, openTrap, openTrap, prepare, presentSinkFields, presentSourceFields, registerShutdownHook, removeListener, removeStepListener, resourceExists, retrieveSinkFields, retrieveSourceFields, setCascade, setCheckpoints, setFlowSkipStrategy, setFlowStepGraph, setFlowStepStrategy, setName, setPlannerInfo, setSinks, setSources, setSpawnStrategy, setSubmitPriority, setTraps, start, stop, toString, updateSchemes, writeDOT, writeStepsDOT
protected HadoopFlow()
protected HadoopFlow(PlatformInfo platformInfo, Map<Object,Object> properties, JobConf jobConf, String name, Map<String,String> flowDescriptor)
public HadoopFlow(PlatformInfo platformInfo, Map<Object,Object> properties, JobConf jobConf, FlowDef flowDef)
protected void initFromProperties(Map<Object,Object> properties)
initFromProperties
in class BaseFlow<JobConf>
protected void initConfig(Map<Object,Object> properties, JobConf parentConfig)
initConfig
in class BaseFlow<JobConf>
protected void setConfigProperty(JobConf config, Object key, Object value)
setConfigProperty
in class BaseFlow<JobConf>
public JobConf getConfigCopy()
public Map<Object,Object> getConfigAsProperties()
public String getProperty(String key)
key
- of type Stringpublic FlowProcess<JobConf> getFlowProcess()
public boolean isPreserveTemporaryFiles()
protected void internalStart()
internalStart
in class BaseFlow<JobConf>
public boolean stepsAreLocal()
protected void internalClean(boolean stop)
internalClean
in class BaseFlow<JobConf>
protected void internalShutdown()
internalShutdown
in class BaseFlow<JobConf>
protected int getMaxNumParallelSteps()
getMaxNumParallelSteps
in class BaseFlow<JobConf>
protected long getTotalSliceCPUMilliSeconds()
getTotalSliceCPUMilliSeconds
in class BaseFlow<JobConf>
Copyright © 2007-2015 Concurrent, Inc. All Rights Reserved.