public class ContinuousHiveSplitEnumerator<T extends Comparable<T>> extends Object implements org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>
SplitEnumerator for hive source.| Constructor and Description |
|---|
ContinuousHiveSplitEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumeratorContext,
T currentReadOffset,
Collection<List<String>> seenPartitionsSinceOffset,
org.apache.flink.connector.file.src.assigners.FileSplitAssigner splitAssigner,
long discoveryInterval,
org.apache.hadoop.mapred.JobConf jobConf,
org.apache.flink.table.catalog.ObjectPath tablePath,
org.apache.flink.table.filesystem.ContinuousPartitionFetcher<org.apache.hadoop.hive.metastore.api.Partition,T> fetcher,
HiveTableSource.HiveContinuousPartitionFetcherContext<T> fetcherContext) |
| Modifier and Type | Method and Description |
|---|---|
void |
addReader(int subtaskId) |
void |
addSplitsBack(List<HiveSourceSplit> splits,
int subtaskId) |
void |
close() |
void |
handleSourceEvent(int subtaskId,
org.apache.flink.api.connector.source.SourceEvent sourceEvent) |
void |
handleSplitRequest(int subtaskId,
String hostName) |
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> |
snapshotState() |
void |
start() |
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitpublic ContinuousHiveSplitEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumeratorContext, T currentReadOffset, Collection<List<String>> seenPartitionsSinceOffset, org.apache.flink.connector.file.src.assigners.FileSplitAssigner splitAssigner, long discoveryInterval, org.apache.hadoop.mapred.JobConf jobConf, org.apache.flink.table.catalog.ObjectPath tablePath, org.apache.flink.table.filesystem.ContinuousPartitionFetcher<org.apache.hadoop.hive.metastore.api.Partition,T> fetcher, HiveTableSource.HiveContinuousPartitionFetcherContext<T> fetcherContext)
public void start()
start in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>public void handleSplitRequest(int subtaskId,
@Nullable
String hostName)
handleSplitRequest in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>public void handleSourceEvent(int subtaskId,
org.apache.flink.api.connector.source.SourceEvent sourceEvent)
handleSourceEvent in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>public void addSplitsBack(List<HiveSourceSplit> splits, int subtaskId)
addSplitsBack in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>public void addReader(int subtaskId)
addReader in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>public org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> snapshotState() throws Exception
snapshotState in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>Exceptionpublic void close()
throws IOException
close in interface AutoCloseableclose in interface org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>>IOExceptionCopyright © 2014–2021 The Apache Software Foundation. All rights reserved.