public class HiveTableSource extends Object implements org.apache.flink.table.sources.StreamTableSource<org.apache.flink.table.dataformat.BaseRow>, org.apache.flink.table.sources.PartitionableTableSource, org.apache.flink.table.sources.ProjectableTableSource<org.apache.flink.table.dataformat.BaseRow>, org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.dataformat.BaseRow>
| 构造器和说明 |
|---|
HiveTableSource(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.table.catalog.ObjectPath tablePath,
org.apache.flink.table.catalog.CatalogTable catalogTable) |
| 限定符和类型 | 方法和说明 |
|---|---|
org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> |
applyLimit(long limit) |
org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> |
applyPartitionPruning(List<Map<String,String>> remainingPartitions) |
String |
explainSource() |
org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.dataformat.BaseRow> |
getDataStream(org.apache.flink.streaming.api.environment.StreamExecutionEnvironment execEnv) |
List<Map<String,String>> |
getPartitions() |
org.apache.flink.table.types.DataType |
getProducedDataType() |
org.apache.flink.table.api.TableSchema |
getTableSchema() |
boolean |
isBounded() |
boolean |
isLimitPushedDown() |
org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> |
projectFields(int[] fields) |
public HiveTableSource(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.table.catalog.ObjectPath tablePath,
org.apache.flink.table.catalog.CatalogTable catalogTable)
public boolean isBounded()
isBounded 在接口中 org.apache.flink.table.sources.StreamTableSource<org.apache.flink.table.dataformat.BaseRow>public org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.dataformat.BaseRow> getDataStream(org.apache.flink.streaming.api.environment.StreamExecutionEnvironment execEnv)
getDataStream 在接口中 org.apache.flink.table.sources.StreamTableSource<org.apache.flink.table.dataformat.BaseRow>public org.apache.flink.table.api.TableSchema getTableSchema()
getTableSchema 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow>public org.apache.flink.table.types.DataType getProducedDataType()
getProducedDataType 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow>public boolean isLimitPushedDown()
isLimitPushedDown 在接口中 org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.dataformat.BaseRow>public org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> applyLimit(long limit)
applyLimit 在接口中 org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.dataformat.BaseRow>public List<Map<String,String>> getPartitions()
getPartitions 在接口中 org.apache.flink.table.sources.PartitionableTableSourcepublic org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> applyPartitionPruning(List<Map<String,String>> remainingPartitions)
applyPartitionPruning 在接口中 org.apache.flink.table.sources.PartitionableTableSourcepublic String explainSource()
explainSource 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow>public org.apache.flink.table.sources.TableSource<org.apache.flink.table.dataformat.BaseRow> projectFields(int[] fields)
projectFields 在接口中 org.apache.flink.table.sources.ProjectableTableSource<org.apache.flink.table.dataformat.BaseRow>Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.