public abstract class CommonExecLegacyTableSourceScan extends ExecNodeBase<RowData> implements MultipleTransformationTranslator<RowData>
ExecNode
to read data from an external source defined by a StreamTableSource
.Modifier and Type | Field and Description |
---|---|
protected List<String> |
qualifiedName |
protected TableSource<?> |
tableSource |
FIELD_NAME_CONFIGURATION, FIELD_NAME_DESCRIPTION, FIELD_NAME_ID, FIELD_NAME_INPUT_PROPERTIES, FIELD_NAME_OUTPUT_TYPE, FIELD_NAME_STATE, FIELD_NAME_TYPE
Constructor and Description |
---|
CommonExecLegacyTableSourceScan(int id,
ExecNodeContext context,
ReadableConfig persistedConfig,
TableSource<?> tableSource,
List<String> qualifiedName,
RowType outputType,
String description) |
Modifier and Type | Method and Description |
---|---|
protected int[] |
computeIndexMapping(boolean isStreaming) |
protected abstract Transformation<RowData> |
createConversionTransformationIfNeeded(StreamExecutionEnvironment streamExecEnv,
ExecNodeConfig config,
ClassLoader classLoader,
Transformation<?> sourceTransform,
org.apache.calcite.rex.RexNode rowtimeExpression) |
protected abstract <IN> Transformation<IN> |
createInput(StreamExecutionEnvironment env,
InputFormat<IN,? extends InputSplit> inputFormat,
TypeInformation<IN> typeInfo) |
protected boolean |
needInternalConversion(int[] fieldIndexes) |
protected Transformation<RowData> |
translateToPlanInternal(org.apache.flink.table.planner.delegation.PlannerBase planner,
ExecNodeConfig config)
Internal method, translates this node into a Flink operator.
|
accept, createFormattedTransformationDescription, createFormattedTransformationName, createTransformationDescription, createTransformationMeta, createTransformationMeta, createTransformationName, createTransformationUid, getContextFromAnnotation, getDescription, getId, getInputEdges, getInputProperties, getOutputType, getPersistedConfig, getSimplifiedName, getTransformation, inputsContainSingleton, replaceInputEdge, setCompiled, setInputEdges, supportFusionCodegen, translateToFusionCodegenSpec, translateToFusionCodegenSpecInternal, translateToPlan
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
translateToPlan
protected final TableSource<?> tableSource
public CommonExecLegacyTableSourceScan(int id, ExecNodeContext context, ReadableConfig persistedConfig, TableSource<?> tableSource, List<String> qualifiedName, RowType outputType, String description)
protected Transformation<RowData> translateToPlanInternal(org.apache.flink.table.planner.delegation.PlannerBase planner, ExecNodeConfig config)
ExecNodeBase
translateToPlanInternal
in class ExecNodeBase<RowData>
planner
- The planner.config
- per-ExecNode
configuration that contains the merged configuration from
various layers which all the nodes implementing this method should use, instead of
retrieving configuration from the planner
. For more details check ExecNodeConfig
.protected abstract <IN> Transformation<IN> createInput(StreamExecutionEnvironment env, InputFormat<IN,? extends InputSplit> inputFormat, TypeInformation<IN> typeInfo)
protected abstract Transformation<RowData> createConversionTransformationIfNeeded(StreamExecutionEnvironment streamExecEnv, ExecNodeConfig config, ClassLoader classLoader, Transformation<?> sourceTransform, @Nullable org.apache.calcite.rex.RexNode rowtimeExpression)
protected boolean needInternalConversion(int[] fieldIndexes)
protected int[] computeIndexMapping(boolean isStreaming)
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.