| Package | Description |
|---|---|
| org.apache.flink.connectors.hive | |
| org.apache.flink.connectors.hive.read |
| Modifier and Type | Method and Description |
|---|---|
HiveSourceSplit |
HiveSourceSplitSerializer.deserialize(int version,
byte[] serialized) |
| Modifier and Type | Method and Description |
|---|---|
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext) |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext) |
static List<HiveSourceSplit> |
HiveSourceFileEnumerator.createInputSplits(int minNumSplits,
List<HiveTablePartition> partitions,
org.apache.hadoop.mapred.JobConf jobConf) |
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> |
ContinuousHivePendingSplitsCheckpointSerializer.deserialize(int version,
byte[] serialized) |
org.apache.flink.core.io.SimpleVersionedSerializer<org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.getEnumeratorCheckpointSerializer() |
org.apache.flink.core.io.SimpleVersionedSerializer<HiveSourceSplit> |
HiveSource.getSplitSerializer() |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext,
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> checkpoint) |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext,
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> checkpoint) |
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> |
ContinuousHiveSplitEnumerator.snapshotState() |
| Modifier and Type | Method and Description |
|---|---|
byte[] |
HiveSourceSplitSerializer.serialize(HiveSourceSplit split) |
| Modifier and Type | Method and Description |
|---|---|
void |
ContinuousHiveSplitEnumerator.addSplitsBack(List<HiveSourceSplit> splits,
int subtaskId) |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.createEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext) |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext,
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> checkpoint) |
org.apache.flink.api.connector.source.SplitEnumerator<HiveSourceSplit,org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.restoreEnumerator(org.apache.flink.api.connector.source.SplitEnumeratorContext<HiveSourceSplit> enumContext,
org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> checkpoint) |
byte[] |
ContinuousHivePendingSplitsCheckpointSerializer.serialize(org.apache.flink.connector.file.src.PendingSplitsCheckpoint<HiveSourceSplit> checkpoint) |
| Constructor and Description |
|---|
ContinuousHivePendingSplitsCheckpoint(Collection<HiveSourceSplit> splits,
Comparable<?> currentReadOffset,
Collection<List<String>> seenPartitionsSinceOffset) |
ContinuousHivePendingSplitsCheckpointSerializer(org.apache.flink.core.io.SimpleVersionedSerializer<HiveSourceSplit> splitSerDe) |
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 |
|---|---|
org.apache.flink.connector.file.src.reader.BulkFormat.Reader<org.apache.flink.table.data.RowData> |
HiveBulkFormatAdapter.createReader(org.apache.flink.configuration.Configuration config,
HiveSourceSplit split) |
org.apache.flink.connector.file.src.reader.BulkFormat.Reader<org.apache.flink.table.data.RowData> |
HiveBulkFormatAdapter.restoreReader(org.apache.flink.configuration.Configuration config,
HiveSourceSplit split) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.