public class HiveTableSource extends Object implements org.apache.flink.table.sources.StreamTableSource<org.apache.flink.table.data.RowData>, org.apache.flink.table.sources.PartitionableTableSource, org.apache.flink.table.sources.ProjectableTableSource<org.apache.flink.table.data.RowData>, org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.data.RowData>, org.apache.flink.table.sources.LookupableTableSource<org.apache.flink.table.data.RowData>
| 构造器和说明 |
|---|
HiveTableSource(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.configuration.ReadableConfig flinkConf,
org.apache.flink.table.catalog.ObjectPath tablePath,
org.apache.flink.table.catalog.CatalogTable catalogTable) |
| 限定符和类型 | 方法和说明 |
|---|---|
org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> |
applyLimit(long limit) |
org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> |
applyPartitionPruning(List<Map<String,String>> remainingPartitions) |
String |
explainSource() |
org.apache.flink.table.functions.AsyncTableFunction<org.apache.flink.table.data.RowData> |
getAsyncLookupFunction(String[] lookupKeys) |
org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.data.RowData> |
getDataStream(org.apache.flink.streaming.api.environment.StreamExecutionEnvironment execEnv) |
org.apache.flink.table.functions.TableFunction<org.apache.flink.table.data.RowData> |
getLookupFunction(String[] lookupKeys) |
List<Map<String,String>> |
getPartitions() |
org.apache.flink.table.types.DataType |
getProducedDataType() |
org.apache.flink.table.api.TableSchema |
getTableSchema() |
boolean |
isAsyncEnabled() |
boolean |
isBounded() |
boolean |
isLimitPushedDown() |
org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> |
projectFields(int[] fields) |
static HiveTablePartition |
toHiveTablePartition(List<String> partitionKeys,
String[] fieldNames,
org.apache.flink.table.types.DataType[] fieldTypes,
HiveShim shim,
Properties tableProps,
String defaultPartitionName,
org.apache.hadoop.hive.metastore.api.Partition partition) |
public HiveTableSource(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.configuration.ReadableConfig flinkConf,
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.data.RowData>public org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.data.RowData> getDataStream(org.apache.flink.streaming.api.environment.StreamExecutionEnvironment execEnv)
getDataStream 在接口中 org.apache.flink.table.sources.StreamTableSource<org.apache.flink.table.data.RowData>public org.apache.flink.table.api.TableSchema getTableSchema()
getTableSchema 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData>public org.apache.flink.table.types.DataType getProducedDataType()
getProducedDataType 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData>public boolean isLimitPushedDown()
isLimitPushedDown 在接口中 org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.data.RowData>public org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> applyLimit(long limit)
applyLimit 在接口中 org.apache.flink.table.sources.LimitableTableSource<org.apache.flink.table.data.RowData>public List<Map<String,String>> getPartitions()
getPartitions 在接口中 org.apache.flink.table.sources.PartitionableTableSourcepublic org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> applyPartitionPruning(List<Map<String,String>> remainingPartitions)
applyPartitionPruning 在接口中 org.apache.flink.table.sources.PartitionableTableSourcepublic org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData> projectFields(int[] fields)
projectFields 在接口中 org.apache.flink.table.sources.ProjectableTableSource<org.apache.flink.table.data.RowData>public static HiveTablePartition toHiveTablePartition(List<String> partitionKeys, String[] fieldNames, org.apache.flink.table.types.DataType[] fieldTypes, HiveShim shim, Properties tableProps, String defaultPartitionName, org.apache.hadoop.hive.metastore.api.Partition partition)
public String explainSource()
explainSource 在接口中 org.apache.flink.table.sources.TableSource<org.apache.flink.table.data.RowData>public org.apache.flink.table.functions.TableFunction<org.apache.flink.table.data.RowData> getLookupFunction(String[] lookupKeys)
getLookupFunction 在接口中 org.apache.flink.table.sources.LookupableTableSource<org.apache.flink.table.data.RowData>public org.apache.flink.table.functions.AsyncTableFunction<org.apache.flink.table.data.RowData> getAsyncLookupFunction(String[] lookupKeys)
getAsyncLookupFunction 在接口中 org.apache.flink.table.sources.LookupableTableSource<org.apache.flink.table.data.RowData>public boolean isAsyncEnabled()
isAsyncEnabled 在接口中 org.apache.flink.table.sources.LookupableTableSource<org.apache.flink.table.data.RowData>Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.