public class HiveTableInputFormat extends org.apache.flink.api.java.hadoop.common.HadoopInputFormatCommonBase<org.apache.flink.table.dataformat.BaseRow,HiveTableInputSplit>
| 限定符和类型 | 字段和说明 |
|---|---|
protected SplitReader |
reader |
| 构造器和说明 |
|---|
HiveTableInputFormat(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.table.catalog.CatalogTable catalogTable,
List<HiveTablePartition> partitions,
int[] projectedFields,
long limit,
String hiveVersion,
boolean useMapRedReader) |
| 限定符和类型 | 方法和说明 |
|---|---|
void |
close() |
void |
configure(org.apache.flink.configuration.Configuration parameters) |
HiveTableInputSplit[] |
createInputSplits(int minNumSplits) |
org.apache.flink.core.io.InputSplitAssigner |
getInputSplitAssigner(HiveTableInputSplit[] inputSplits) |
org.apache.flink.api.common.io.statistics.BaseStatistics |
getStatistics(org.apache.flink.api.common.io.statistics.BaseStatistics cachedStats) |
org.apache.flink.table.dataformat.BaseRow |
nextRecord(org.apache.flink.table.dataformat.BaseRow reuse) |
void |
open(HiveTableInputSplit split) |
boolean |
reachedEnd() |
getCredentialsFromUGI, read, write@VisibleForTesting protected transient SplitReader reader
public HiveTableInputFormat(org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.table.catalog.CatalogTable catalogTable,
List<HiveTablePartition> partitions,
int[] projectedFields,
long limit,
String hiveVersion,
boolean useMapRedReader)
public void configure(org.apache.flink.configuration.Configuration parameters)
public void open(HiveTableInputSplit split) throws IOException
IOExceptionpublic boolean reachedEnd()
throws IOException
IOExceptionpublic org.apache.flink.table.dataformat.BaseRow nextRecord(org.apache.flink.table.dataformat.BaseRow reuse)
throws IOException
IOExceptionpublic void close()
throws IOException
IOExceptionpublic HiveTableInputSplit[] createInputSplits(int minNumSplits) throws IOException
IOExceptionpublic org.apache.flink.api.common.io.statistics.BaseStatistics getStatistics(org.apache.flink.api.common.io.statistics.BaseStatistics cachedStats)
public org.apache.flink.core.io.InputSplitAssigner getInputSplitAssigner(HiveTableInputSplit[] inputSplits)
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.