Modifier and Type | Method and Description |
---|---|
String |
AbstractRuntimeUDFContext.getAllocationIDAsString() |
Modifier and Type | Method and Description |
---|---|
protected Connection |
JDBCInputFormat.getDbConn()
Deprecated.
|
protected PreparedStatement |
JDBCInputFormat.getStatement()
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
TypeSerializer<Object>[] |
TupleSerializerBase.getFieldSerializers() |
Modifier and Type | Method and Description |
---|---|
com.esotericsoftware.kryo.Kryo |
KryoSerializer.getKryo() |
Modifier and Type | Method and Description |
---|---|
Collection<State<T>> |
NFA.getStates() |
Modifier and Type | Method and Description |
---|---|
int |
SharedBuffer.getEventsBufferCacheSize() |
int |
SharedBuffer.getEventsBufferSize() |
int |
SharedBuffer.getSharedBufferNodeCacheSize() |
int |
SharedBuffer.getSharedBufferNodeSize() |
Modifier and Type | Method and Description |
---|---|
<M extends MessageHeaders<R,P,U>,U extends MessageParameters,R extends RequestBody,P extends ResponseBody> |
RestClusterClient.sendRequest(M messageHeaders,
U messageParameters,
R request) |
Modifier and Type | Method and Description |
---|---|
static Map<String,String> |
ConfigurationUtils.parseJvmArgString(String jvmParamsStr) |
static Map<String,String> |
ConfigurationUtils.parseTmResourceDynamicConfigs(String dynamicConfigsStr) |
Modifier and Type | Method and Description |
---|---|
T |
FutureCompletingBlockingQueue.take()
Warning: This is a dangerous method and should only be used for testing convenience.
|
Modifier and Type | Method and Description |
---|---|
HBaseOptions |
HBaseUpsertTableSink.getHBaseOptions() |
HBaseOptions |
HBaseDynamicTableSink.getHBaseOptions() |
HBaseTableSchema |
HBaseUpsertTableSink.getHBaseTableSchema() |
HBaseTableSchema |
HBaseDynamicTableSink.getHBaseTableSchema() |
HBaseWriteOptions |
HBaseUpsertTableSink.getWriteOptions() |
HBaseWriteOptions |
HBaseDynamicTableSink.getWriteOptions() |
Modifier and Type | Method and Description |
---|---|
HBaseTableSchema |
HBaseDynamicTableSource.getHBaseTableSchema() |
HBaseTableSchema |
HBaseTableSource.getHBaseTableSchema() |
String |
HBaseRowDataLookupFunction.getHTableName() |
String |
HBaseLookupFunction.getHTableName() |
Modifier and Type | Method and Description |
---|---|
protected Connection |
JdbcInputFormat.getDbConn() |
protected PreparedStatement |
JdbcInputFormat.getStatement() |
Modifier and Type | Method and Description |
---|---|
AbstractJdbcCatalog |
JdbcCatalog.getInternal() |
Modifier and Type | Method and Description |
---|---|
Connection |
AbstractJdbcOutputFormat.getConnection() |
Modifier and Type | Method and Description |
---|---|
Connection |
JdbcRowDataLookupFunction.getDbConnection() |
Connection |
JdbcLookupFunction.getDbConnection() |
Modifier and Type | Field and Description |
---|---|
protected SplitReader |
HiveTableInputFormat.reader |
Modifier and Type | Method and Description |
---|---|
void |
RocksDBKeyedStateBackend.compactState(StateDescriptor<?,?> stateDesc) |
PredefinedOptions |
RocksDBStateBackend.getPredefinedOptions()
Gets the currently set predefined options for RocksDB.
|
int |
RocksDBKeyedStateBackend.numKeyValueStateEntries() |
Modifier and Type | Method and Description |
---|---|
int |
RefCountedFile.getReferenceCounter() |
Modifier and Type | Method and Description |
---|---|
static MemorySegment |
MemorySegmentFactory.allocateOffHeapUnsafeMemory(int size) |
Modifier and Type | Method and Description |
---|---|
static ClassLoader |
PluginLoader.createPluginClassLoader(PluginDescriptor pluginDescriptor,
ClassLoader parentClassLoader,
String[] alwaysParentFirstPatterns) |
Constructor and Description |
---|
PluginLoader(ClassLoader pluginClassLoader) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.parquet.filter2.predicate.FilterPredicate |
ParquetInputFormat.getPredicate() |
Modifier and Type | Method and Description |
---|---|
int |
RefCountedBufferingFileStream.getReferenceCounter() |
Constructor and Description |
---|
RefCountedBufferingFileStream(RefCountedFileWithStream file,
int bufferSize) |
Modifier and Type | Method and Description |
---|---|
static String |
FlinkConfMountDecorator.getFlinkConfConfigMapName(String clusterId) |
Modifier and Type | Method and Description |
---|---|
protected OrcRowInputFormat |
OrcTableSource.buildOrcInputFormat() |
org.apache.orc.RecordReader |
OrcSplitReader.getRecordReader() |
Modifier and Type | Method and Description |
---|---|
static Tuple2<Long,Long> |
OrcShimV200.getOffsetAndLengthForSplit(long splitStart,
long splitLength,
List<org.apache.orc.StripeInformation> stripes) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.orc.OrcFile.WriterOptions |
OrcBulkWriterFactory.getWriterOptions() |
Modifier and Type | Method and Description |
---|---|
org.apache.beam.runners.fnexecution.control.JobBundleFactory |
AbstractPythonFunctionRunner.createJobBundleFactory(org.apache.beam.vendor.grpc.v1p21p0.com.google.protobuf.Struct pipelineOptions) |
Modifier and Type | Field and Description |
---|---|
static String |
ProcessPythonEnvironmentManager.PYTHON_REQUIREMENTS_CACHE |
static String |
ProcessPythonEnvironmentManager.PYTHON_REQUIREMENTS_FILE |
static String |
ProcessPythonEnvironmentManager.PYTHON_REQUIREMENTS_INSTALL_DIR |
static String |
ProcessPythonEnvironmentManager.PYTHON_WORKING_DIR |
Modifier and Type | Field and Description |
---|---|
static String |
PythonEnvironmentManagerUtils.PYFLINK_UDF_RUNNER_DIR |
Modifier and Type | Method and Description |
---|---|
boolean |
Client.isEventGroupShutdown() |
boolean |
AbstractServerBase.isEventGroupShutdown() |
Modifier and Type | Class and Description |
---|---|
class |
VoidBlobWriter
BlobWriter which does not support writing BLOBs to a store.
|
Modifier and Type | Method and Description |
---|---|
byte[] |
BlobKey.getHash()
Returns the hash component of this key.
|
File |
BlobServer.getStorageLocation(JobID jobId,
BlobKey key)
Returns a file handle to the file associated with the given blob key on the blob server.
|
File |
PermanentBlobCache.getStorageLocation(JobID jobId,
BlobKey key)
Returns a file handle to the file associated with the given blob key on the blob server.
|
File |
TransientBlobCache.getStorageLocation(JobID jobId,
BlobKey key)
Returns a file handle to the file associated with the given blob key on the blob server.
|
Constructor and Description |
---|
PermanentBlobKey()
Constructs a new BLOB key.
|
TransientBlobKey()
Constructs a new BLOB key.
|
Modifier and Type | Method and Description |
---|---|
static void |
StateAssignmentOperation.extractIntersectingState(Collection<? extends KeyedStateHandle> originalSubtaskStateHandles,
KeyGroupRange rangeToExtract,
List<KeyedStateHandle> extractedStateCollector)
Extracts certain key group ranges from the given state handles and adds them to the
collector.
|
static CheckpointOptions |
CheckpointOptions.forCheckpointWithDefaultLocation() |
CompletableFuture<CompletedCheckpoint> |
CheckpointCoordinator.triggerCheckpoint(CheckpointProperties props,
String externalSavepointLocation,
boolean isPeriodic) |
Constructor and Description |
---|
CheckpointCoordinator(JobID job,
CheckpointCoordinatorConfiguration chkConfig,
ExecutionVertex[] tasksToTrigger,
ExecutionVertex[] tasksToWaitFor,
ExecutionVertex[] tasksToCommitTo,
Collection<OperatorCoordinatorCheckpointContext> coordinatorsToCheckpoint,
CheckpointIDCounter checkpointIDCounter,
CompletedCheckpointStore completedCheckpointStore,
StateBackend checkpointStateBackend,
Executor executor,
ScheduledExecutor timer,
SharedStateRegistryFactory sharedStateRegistryFactory,
CheckpointFailureManager failureManager,
Clock clock) |
CheckpointOptions(CheckpointType checkpointType,
CheckpointStorageLocationReference targetLocation) |
Modifier and Type | Method and Description |
---|---|
static akka.actor.ActorSystem |
BootstrapTools.startRemoteActorSystem(Configuration configuration,
String externalAddress,
String externalPortRange,
org.slf4j.Logger logger)
Starts a remote ActorSystem at given address and specific port range.
|
Constructor and Description |
---|
TaskExecutorProcessSpec(CPUResource cpuCores,
MemorySize frameworkHeapSize,
MemorySize frameworkOffHeapSize,
MemorySize taskHeapSize,
MemorySize taskOffHeapSize,
MemorySize networkMemSize,
MemorySize managedMemorySize,
MemorySize jvmMetaspaceSize,
MemorySize jvmOverheadSize) |
Modifier and Type | Field and Description |
---|---|
static ResourceProfile |
ResourceProfile.ANY
A ResourceProfile that indicates infinite resource that matches any resource requirement, for
testability purpose only.
|
Modifier and Type | Method and Description |
---|---|
static ResourceProfile |
ResourceProfile.fromResources(double cpuCores,
int taskHeapMemoryMB) |
static SlotProfile |
SlotProfile.noLocality(ResourceProfile resourceProfile)
Returns a slot profile for the given resource profile, without any locality requirements.
|
static SlotProfile |
SlotProfile.noRequirements()
Returns a slot profile that has no requirements.
|
static SlotProfile |
SlotProfile.preferredLocality(ResourceProfile resourceProfile,
Collection<TaskManagerLocation> preferredLocations)
Returns a slot profile for the given resource profile and the preferred locations.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Collection<TaskManagerLocation>> |
Execution.calculatePreferredLocations(LocationPreferenceConstraint locationPreferenceConstraint)
Calculates the preferred locations based on the location preference constraint.
|
void |
ExecutionVertex.deployToSlot(LogicalSlot slot) |
JobStatus |
ExecutionGraph.waitUntilTerminal() |
Modifier and Type | Method and Description |
---|---|
SchedulingPipelinedRegion |
RestartPipelinedRegionFailoverStrategy.getFailoverRegion(ExecutionVertexID vertexID)
Returns the failover region that contains the given execution vertex.
|
Constructor and Description |
---|
RestartPipelinedRegionFailoverStrategy(SchedulingTopology topology)
Creates a new failover strategy to restart pipelined regions that works on the given
topology.
|
Modifier and Type | Method and Description |
---|---|
boolean |
EmbeddedLeaderService.isShutdown() |
Modifier and Type | Method and Description |
---|---|
NettyShuffleEnvironmentConfiguration |
NettyShuffleEnvironment.getConfiguration() |
ConnectionManager |
NettyShuffleEnvironment.getConnectionManager() |
Optional<InputGate> |
NettyShuffleEnvironment.getInputGate(InputGateID id) |
NetworkBufferPool |
NettyShuffleEnvironment.getNetworkBufferPool() |
ResultPartitionManager |
NettyShuffleEnvironment.getResultPartitionManager() |
Modifier and Type | Field and Description |
---|---|
static String |
RecordWriter.DEFAULT_OUTPUT_FLUSH_THREAD_NAME
Default name for the output flush thread, if no name with a task reference is given.
|
Modifier and Type | Method and Description |
---|---|
Meter |
RecordWriter.getIdleTimeMsPerSecond() |
Modifier and Type | Method and Description |
---|---|
Buffer |
BufferDecompressor.decompressToOriginalBuffer(Buffer buffer)
The difference between this method and
BufferDecompressor.decompressToIntermediateBuffer(Buffer) is
that this method copies the decompressed data to the input Buffer starting from
offset 0. |
MemorySegment |
BufferBuilder.getMemorySegment() |
BufferRecycler |
BufferBuilder.getRecycler() |
Constructor and Description |
---|
NetworkBufferPool(int numberOfSegmentsToAllocate,
int segmentSize,
int numberOfSegmentsToRequest) |
Modifier and Type | Method and Description |
---|---|
abstract int |
AbstractBuffersUsageGauge.calculateTotalBuffers(SingleInputGate inputGate) |
abstract int |
AbstractBuffersUsageGauge.calculateUsedBuffers(SingleInputGate inputGate) |
Modifier and Type | Method and Description |
---|---|
ResultPartition |
ResultPartitionFactory.create(String taskNameWithSubtaskAndId,
int partitionIndex,
ResultPartitionID id,
ResultPartitionType type,
int numberOfSubpartitions,
int maxParallelism,
FunctionWithException<BufferPoolOwner,BufferPool,IOException> bufferPoolFactory) |
int |
PipelinedSubpartition.getBuffersInBacklog()
Gets the number of non-event buffers in this subpartition.
|
static DataSetMetaInfo |
DataSetMetaInfo.withNumRegisteredPartitions(int numRegisteredPartitions,
int numTotalPartitions) |
Constructor and Description |
---|
ResultPartitionID() |
Modifier and Type | Method and Description |
---|---|
void |
SingleInputGate.assignExclusiveSegments()
Assign the exclusive buffers to all remote input channels directly for credit-based mode.
|
protected InputChannel |
SingleInputGateFactory.createKnownInputChannel(SingleInputGate inputGate,
int index,
NettyShuffleDescriptor inputChannelDescriptor,
SingleInputGateFactory.ChannelStatistics channelStatistics,
InputChannelMetrics metrics) |
Buffer |
RemoteInputChannel.getNextReceivedBuffer() |
int |
RemoteInputChannel.getNumberOfAvailableBuffers() |
protected int |
RecoveredInputChannel.getNumberOfQueuedBuffers() |
int |
RemoteInputChannel.getNumberOfRequiredBuffers() |
int |
RemoteInputChannel.getSenderBacklog() |
void |
RemoteInputChannel.requestSubpartition(int subpartitionIndex)
Requests a remote subpartition.
|
Constructor and Description |
---|
BufferOrEvent(AbstractEvent event,
InputChannelInfo channelInfo) |
BufferOrEvent(Buffer buffer,
InputChannelInfo channelInfo) |
Constructor and Description |
---|
IntermediateResultPartitionID()
Creates an new random intermediate result partition ID for testing.
|
Constructor and Description |
---|
CheckpointCoordinatorConfiguration(long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
CheckpointRetentionPolicy checkpointRetentionPolicy,
boolean isExactlyOnce,
boolean isUnalignedCheckpoint,
boolean isPreferCheckpointForRecovery,
int tolerableCpFailureNumber)
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
static JobManagerProcessSpec |
JobManagerProcessUtils.createDefaultJobManagerProcessSpec(int totalProcessMemoryMb) |
Constructor and Description |
---|
ScheduledUnit(Execution task) |
ScheduledUnit(JobVertexID jobVertexId,
SlotSharingGroupId slotSharingGroupId,
CoLocationConstraint coLocationConstraint) |
Modifier and Type | Method and Description |
---|---|
Collection<SlotSharingManager.MultiTaskSlot> |
SlotSharingManager.getResolvedRootSlots()
Returns a collection of all resolved root slots.
|
int |
SlotPoolImpl.AvailableSlots.size() |
protected void |
SlotPoolImpl.timeoutPendingSlotRequest(SlotRequestId slotRequestId) |
Constructor and Description |
---|
SchedulerImpl(SlotSelectionStrategy slotSelectionStrategy,
SlotPool slotPool,
Map<SlotSharingGroupId,SlotSharingManager> slotSharingManagers) |
Modifier and Type | Method and Description |
---|---|
org.apache.flink.runtime.memory.UnsafeMemoryBudget |
MemoryManager.getMemoryBudget()
Gets memory budget.
|
boolean |
MemoryManager.isShutdown()
Checks whether the MemoryManager has been shut down.
|
Constructor and Description |
---|
AcknowledgeCheckpoint(JobID jobId,
ExecutionAttemptID taskExecutionId,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
static ReporterSetup |
ReporterSetup.forReporter(String reporterName,
MetricConfig metricConfig,
MetricReporter reporter) |
static ReporterSetup |
ReporterSetup.forReporter(String reporterName,
MetricReporter reporter) |
List<MetricReporter> |
MetricRegistryImpl.getReporters() |
Modifier and Type | Method and Description |
---|---|
protected Collection<? extends DispatcherResourceManagerComponent> |
MiniCluster.createDispatcherResourceManagerComponents(Configuration configuration,
MiniCluster.RpcServiceFactory rpcServiceFactory,
HighAvailabilityServices haServices,
BlobServer blobServer,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
MetricQueryServiceRetriever metricQueryServiceRetriever,
FatalErrorHandler fatalErrorHandler) |
protected HighAvailabilityServices |
MiniCluster.createHighAvailabilityServices(Configuration configuration,
Executor executor) |
protected CompletableFuture<DispatcherGateway> |
MiniCluster.getDispatcherGatewayFuture() |
protected CompletableFuture<Void> |
MiniCluster.terminateTaskExecutor(int index) |
protected boolean |
MiniCluster.useLocalCommunication() |
Modifier and Type | Method and Description |
---|---|
static SSLContext |
SSLUtils.createRestSSLContext(Configuration config,
boolean clientMode)
Creates an SSL context for clients against the external REST endpoint.
|
Modifier and Type | Method and Description |
---|---|
protected OperatorCoordinator |
RecreateOnResetOperatorCoordinator.Provider.create(OperatorCoordinator.Context context,
long closingTimeoutMs) |
OperatorCoordinator |
RecreateOnResetOperatorCoordinator.getInternalCoordinator() |
Modifier and Type | Method and Description |
---|---|
int |
KvStateEntry.getCacheSize() |
Modifier and Type | Method and Description |
---|---|
abstract boolean |
ResourceManager.startNewWorker(WorkerResourceSpec workerResourceSpec)
Allocates a resource using the worker resource specification.
|
Modifier and Type | Method and Description |
---|---|
static ResourceProfile |
SlotManagerImpl.generateDefaultSlotResourceProfile(WorkerResourceSpec workerResourceSpec,
int numSlotsPerWorker) |
int |
SlotManagerImpl.getNumberAssignedPendingTaskManagerSlots() |
int |
SlotManagerImpl.getNumberPendingTaskManagerSlots() |
void |
SlotManager.unregisterTaskManagersAndReleaseResources() |
void |
SlotManagerImpl.unregisterTaskManagersAndReleaseResources() |
Modifier and Type | Method and Description |
---|---|
void |
MetricStore.add(MetricDump metric) |
Modifier and Type | Method and Description |
---|---|
static AkkaRpcServiceUtils.AkkaRpcServiceBuilder |
AkkaRpcServiceUtils.remoteServiceBuilder(Configuration configuration,
String externalAddress,
int externalPort) |
Constructor and Description |
---|
AkkaRpcService(akka.actor.ActorSystem actorSystem,
AkkaRpcServiceConfiguration configuration) |
Modifier and Type | Method and Description |
---|---|
ExecutionGraph |
SchedulerBase.getExecutionGraph()
ExecutionGraph is exposed to make it easier to rework tests to be based on the new scheduler.
|
Modifier and Type | Method and Description |
---|---|
SecurityConfiguration |
HadoopModule.getSecurityConfig() |
Constructor and Description |
---|
NetworkPartitionConnectionInfo(ConnectionID connectionID) |
PartitionDescriptor(IntermediateDataSetID resultId,
int totalNumberOfPartitions,
IntermediateResultPartitionID partitionId,
ResultPartitionType partitionType,
int numberOfSubpartitions,
int connectionIndex) |
ProducerDescriptor(ResourceID producerLocation,
ExecutionAttemptID producerExecutionId,
InetAddress address,
int dataPort) |
Modifier and Type | Class and Description |
---|---|
class |
JavaSerializer<T extends Serializable>
A
TypeSerializer that uses Java serialization. |
Modifier and Type | Field and Description |
---|---|
protected boolean |
DefaultOperatorStateBackendBuilder.asynchronousSnapshots
Flag to de/activate asynchronous snapshots.
|
protected CloseableRegistry |
DefaultOperatorStateBackendBuilder.cancelStreamRegistry |
protected ExecutionConfig |
DefaultOperatorStateBackendBuilder.executionConfig
The execution configuration.
|
protected Collection<OperatorStateHandle> |
DefaultOperatorStateBackendBuilder.restoreStateHandles
State handles for restore.
|
protected ClassLoader |
DefaultOperatorStateBackendBuilder.userClassloader
The user code classloader.
|
Modifier and Type | Method and Description |
---|---|
protected void |
AsyncSnapshotCallable.cancel() |
SharedStateRegistryKey |
IncrementalRemoteKeyedStateHandle.createSharedStateRegistryKeyFromFileName(StateHandleID shId)
Create a unique key to register one of our shared state handles.
|
boolean |
IncrementalRemoteKeyedStateHandle.equals(Object o)
This method is should only be called in tests! This should never serve as key in a hash map.
|
StreamCompressionDecorator |
AbstractKeyedStateBackend.getKeyGroupCompressionDecorator() |
int |
IncrementalRemoteKeyedStateHandle.hashCode()
This method should only be called in tests! This should never serve as key in a hash map.
|
abstract int |
AbstractKeyedStateBackend.numKeyValueStateEntries()
Returns the total number of state entries across all keys/namespaces.
|
int |
AbstractKeyedStateBackend.numKeyValueStatesByName() |
boolean |
AbstractKeyedStateBackend.supportsAsynchronousSnapshots() |
Constructor and Description |
---|
SharedStateRegistryKey(String keyString) |
StateSnapshotContextSynchronousImpl(long checkpointId,
long checkpointTimestamp) |
Modifier and Type | Class and Description |
---|---|
protected static class |
CopyOnWriteStateMap.StateMapEntry<K,N,S>
One entry in the
CopyOnWriteStateMap . |
Modifier and Type | Method and Description |
---|---|
LocalRecoveryConfig |
HeapKeyedStateBackend.getLocalRecoveryConfig() |
StateMap<K,N,S>[] |
StateTable.getState()
Returns the internal data structure.
|
StateTable<K,N,SV> |
AbstractHeapState.getStateTable()
This should only be used for testing.
|
int |
HeapKeyedStateBackend.numKeyValueStateEntries()
Returns the total number of state entries across all keys/namespaces.
|
int |
HeapKeyedStateBackend.numKeyValueStateEntries(Object namespace)
Returns the total number of state entries across all keys for the given namespace.
|
abstract int |
StateMap.sizeOfNamespace(Object namespace) |
int |
StateTable.sizeOfNamespace(Object namespace) |
Modifier and Type | Field and Description |
---|---|
static String |
TaskManagerServices.LOCAL_STATE_SUB_DIRECTORY_ROOT |
Modifier and Type | Method and Description |
---|---|
boolean |
DefaultJobLeaderService.containsJob(JobID jobId)
Check whether the service monitors the given job.
|
static TaskExecutorResourceSpec |
TaskExecutorResourceUtils.resourceSpecFromConfigForLocalExecution(Configuration config) |
Modifier and Type | Method and Description |
---|---|
boolean |
TaskSlotTableImpl.allocateSlot(int index,
JobID jobId,
AllocationID allocationId,
Time slotTimeout) |
boolean |
TaskSlotTable.allocateSlot(int index,
JobID jobId,
AllocationID allocationId,
Time slotTimeout)
Allocate the slot with the given index for the given job and allocation id.
|
boolean |
TaskSlotTableImpl.isClosed() |
Modifier and Type | Class and Description |
---|---|
static class |
TaskManagerLocation.DefaultHostNameSupplier
This Supplier class could retrieve the FQDN host name of the given InetAddress on demand,
extract the pure host name and cache the results for later use.
|
static class |
TaskManagerLocation.IpOnlyHostNameSupplier
This Supplier class returns the IP address of the given InetAddress directly, therefore no
reverse DNS lookup is required.
|
Modifier and Type | Method and Description |
---|---|
static void |
Task.setupPartitionsAndGates(ResultPartitionWriter[] producedPartitions,
InputGate[] inputGates) |
Constructor and Description |
---|
TaskManagerLocation(ResourceID resourceID,
InetAddress inetAddress,
int dataPort)
Constructs a new instance connection info object.
|
TaskManagerLocation(ResourceID resourceID,
InetAddress inetAddress,
int dataPort,
TaskManagerLocation.HostNameSupplier hostNameSupplier)
Constructs a new instance connection info object.
|
Modifier and Type | Field and Description |
---|---|
static String |
BashJavaUtils.EXECUTION_PREFIX |
Modifier and Type | Class and Description |
---|---|
static class |
CoGroupedStreams.UnionSerializer<T1,T2>
|
Modifier and Type | Class and Description |
---|---|
static class |
TwoPhaseCommitSinkFunction.State<TXN,CONTEXT>
State POJO class coupling pendingTransaction, context and pendingCommitTransactions.
|
static class |
TwoPhaseCommitSinkFunction.StateSerializer<TXN,CONTEXT>
Custom
TypeSerializer for the sink state. |
static class |
TwoPhaseCommitSinkFunction.TransactionHolder<TXN>
Adds metadata (currently only the start time of the transaction) to the transaction object.
|
Constructor and Description |
---|
TransactionHolder(TXN handle,
long transactionStartTime) |
Modifier and Type | Method and Description |
---|---|
long |
Buckets.getMaxPartCounter() |
Bucket<IN,BucketID> |
Buckets.onElement(IN value,
SinkFunction.Context context) |
Modifier and Type | Method and Description |
---|---|
long |
ContinuousFileMonitoringFunction.getGlobalModificationTime() |
Modifier and Type | Field and Description |
---|---|
static String |
StreamConfig.SERIALIZEDUDF |
Modifier and Type | Method and Description |
---|---|
StreamOperator<?> |
StreamNode.getOperator() |
List<StreamEdge> |
StreamGraph.getStreamEdges(int sourceId) |
List<StreamEdge> |
StreamGraph.getStreamEdges(int sourceId,
int targetId) |
List<StreamEdge> |
StreamGraph.getStreamEdgesOrThrow(int sourceId,
int targetId)
Deprecated.
|
<T extends StreamOperator<?>> |
StreamConfig.getStreamOperator(ClassLoader cl) |
Constructor and Description |
---|
StreamNode(Integer id,
String slotSharingGroup,
String coLocationGroup,
StreamOperator<?> operator,
String operatorName,
List<OutputSelector<?>> outputSelector,
Class<? extends AbstractInvokable> jobVertexClass) |
Modifier and Type | Method and Description |
---|---|
<K,N> InternalTimerService<N> |
AbstractStreamOperatorV2.getInternalTimerService(String name,
TypeSerializer<N> namespaceSerializer,
Triggerable<K,N> triggerable)
Returns a
InternalTimerService that can be used to query current processing time and
event time and to set timers. |
<K> KeyedStateBackend<K> |
AbstractStreamOperator.getKeyedStateBackend() |
<K> KeyedStateBackend<K> |
AbstractStreamOperatorV2.getKeyedStateBackend() |
OperatorStateBackend |
AbstractStreamOperator.getOperatorStateBackend() |
OperatorStateBackend |
AbstractStreamOperatorV2.getOperatorStateBackend() |
ProcessingTimeService |
AbstractStreamOperator.getProcessingTimeService()
Returns the
ProcessingTimeService responsible for getting the current processing time
and registering timers. |
ProcessingTimeService |
AbstractStreamOperatorV2.getProcessingTimeService()
Returns the
ProcessingTimeService responsible for getting the current processing time
and registering timers. |
StreamingRuntimeContext |
AbstractStreamOperator.getRuntimeContext()
Returns a context that allows the operator to query information about the execution and also
to interact with systems such as broadcast variables and managed state.
|
SourceReader<OUT,SplitT> |
SourceOperator.getSourceReader() |
int |
InternalTimerServiceImpl.numEventTimeTimers() |
int |
InternalTimeServiceManager.numEventTimeTimers() |
int |
AbstractStreamOperator.numEventTimeTimers() |
int |
AbstractStreamOperatorV2.numEventTimeTimers() |
int |
InternalTimerServiceImpl.numEventTimeTimers(N namespace) |
int |
InternalTimerServiceImpl.numProcessingTimeTimers() |
int |
InternalTimeServiceManager.numProcessingTimeTimers() |
int |
AbstractStreamOperator.numProcessingTimeTimers() |
int |
AbstractStreamOperatorV2.numProcessingTimeTimers() |
int |
InternalTimerServiceImpl.numProcessingTimeTimers(N namespace) |
Constructor and Description |
---|
StreamingRuntimeContext(AbstractStreamOperator<?> operator,
Environment env,
Map<String,Accumulator<?,?>> accumulators) |
StreamTaskStateInitializerImpl(Environment environment,
StateBackend stateBackend,
TtlTimeProvider ttlTimeProvider) |
Modifier and Type | Class and Description |
---|---|
static class |
IntervalJoinOperator.BufferEntry<T>
A container for elements put in the left/write buffer.
|
static class |
IntervalJoinOperator.BufferEntrySerializer<T>
A
serializer for the IntervalJoinOperator.BufferEntry . |
Modifier and Type | Method and Description |
---|---|
T |
IntervalJoinOperator.BufferEntry.getElement() |
boolean |
IntervalJoinOperator.BufferEntry.hasBeenJoined() |
Modifier and Type | Method and Description |
---|---|
static <T> byte[] |
CollectSinkFunction.serializeAccumulatorResult(long offset,
String version,
long lastCheckpointedOffset,
List<T> bufferedResults,
TypeSerializer<T> serializer) |
Constructor and Description |
---|
CollectResultIterator(CompletableFuture<OperatorID> operatorIdFuture,
TypeSerializer<T> serializer,
String accumulatorName,
int retryMillis) |
Modifier and Type | Method and Description |
---|---|
StreamSource<T,?> |
LegacySourceTransformation.getOperator() |
TwoInputStreamOperator<IN1,IN2,OUT> |
TwoInputTransformation.getOperator() |
StreamSink<T> |
SinkTransformation.getOperator() |
OneInputStreamOperator<IN,OUT> |
OneInputTransformation.getOperator() |
Modifier and Type | Method and Description |
---|---|
long |
TimeEvictor.getWindowSize() |
Modifier and Type | Method and Description |
---|---|
long |
ContinuousEventTimeTrigger.getInterval() |
long |
ContinuousProcessingTimeTrigger.getInterval() |
Trigger<T,W> |
PurgingTrigger.getNestedTrigger() |
Modifier and Type | Method and Description |
---|---|
protected org.elasticsearch.action.bulk.BulkProcessor |
ElasticsearchSinkBase.buildBulkProcessor(org.elasticsearch.action.bulk.BulkProcessor.Listener listener)
Build the
BulkProcessor . |
Modifier and Type | Method and Description |
---|---|
org.apache.flink.streaming.connectors.fs.bucketing.BucketingSink.State<T> |
BucketingSink.getState()
Deprecated.
|
Modifier and Type | Class and Description |
---|---|
static class |
FlinkKafkaProducer.ContextStateSerializer
|
static class |
FlinkKafkaProducer.KafkaTransactionContext
Context associated to this instance of the
FlinkKafkaProducer . |
static class |
FlinkKafkaProducer.KafkaTransactionState
State for handling transactions.
|
static class |
FlinkKafkaProducer.NextTransactionalIdHintSerializer
|
static class |
FlinkKafkaProducer.TransactionStateSerializer
|
static class |
FlinkKafkaProducer011.ContextStateSerializer
|
static class |
FlinkKafkaProducer011.KafkaTransactionContext
Context associated to this instance of the
FlinkKafkaProducer011 . |
static class |
FlinkKafkaProducer011.KafkaTransactionState
State for handling transactions.
|
static class |
FlinkKafkaProducer011.NextTransactionalIdHintSerializer
|
static class |
FlinkKafkaProducer011.TransactionStateSerializer
|
Modifier and Type | Method and Description |
---|---|
boolean |
FlinkKafkaConsumerBase.getEnableCommitOnCheckpoints() |
protected <K,V> org.apache.kafka.clients.producer.KafkaProducer<K,V> |
FlinkKafkaProducerBase.getKafkaProducer(Properties props)
Used for testing only.
|
protected long |
FlinkKafkaProducerBase.numPendingRecords() |
Constructor and Description |
---|
KafkaTransactionContext(Set<String> transactionalIds) |
KafkaTransactionContext(Set<String> transactionalIds) |
KafkaTransactionState(FlinkKafkaInternalProducer<byte[],byte[]> producer) |
KafkaTransactionState(FlinkKafkaProducer<byte[],byte[]> producer) |
KafkaTransactionState(String transactionalId,
FlinkKafkaInternalProducer<byte[],byte[]> producer) |
KafkaTransactionState(String transactionalId,
FlinkKafkaProducer<byte[],byte[]> producer) |
KafkaTransactionState(String transactionalId,
long producerId,
short epoch,
FlinkKafkaInternalProducer<byte[],byte[]> producer) |
KafkaTransactionState(String transactionalId,
long producerId,
short epoch,
FlinkKafkaProducer<byte[],byte[]> producer) |
Modifier and Type | Class and Description |
---|---|
static class |
KafkaShuffleFetcher.KafkaShuffleElement
An element in a KafkaShuffle.
|
static class |
KafkaShuffleFetcher.KafkaShuffleElementDeserializer<T>
Deserializer for KafkaShuffleElement.
|
static class |
KafkaShuffleFetcher.KafkaShuffleRecord<T>
One value with Type T in a KafkaShuffle.
|
static class |
KafkaShuffleFetcher.KafkaShuffleWatermark
A watermark element in a KafkaShuffle.
|
Modifier and Type | Method and Description |
---|---|
KafkaShuffleFetcher.KafkaShuffleElement |
KafkaShuffleFetcher.KafkaShuffleElementDeserializer.deserialize(org.apache.kafka.clients.consumer.ConsumerRecord<byte[],byte[]> record) |
protected org.apache.kafka.clients.consumer.ConsumerRecords<byte[],byte[]> |
KafkaConsumerThread.getRecordsFromKafka()
Get records from Kafka.
|
int |
FlinkKafkaInternalProducer.getTransactionCoordinatorId() |
int |
FlinkKafkaProducer.getTransactionCoordinatorId() |
Constructor and Description |
---|
KafkaShuffleElementDeserializer(TypeSerializer<T> typeSerializer) |
Modifier and Type | Method and Description |
---|---|
protected com.amazonaws.services.kinesis.producer.KinesisProducer |
FlinkKinesisProducer.getKinesisProducer(com.amazonaws.services.kinesis.producer.KinesisProducerConfiguration producerConfig)
Creates a
KinesisProducer . |
Modifier and Type | Method and Description |
---|---|
protected ExecutorService |
KinesisDataFetcher.createShardConsumersThreadPool(String subtaskName) |
protected void |
KinesisDataFetcher.emitWatermark()
Called periodically to emit a watermark.
|
protected long |
KinesisDataFetcher.getCurrentTimeMillis()
Return the current system time.
|
List<KinesisStreamShardState> |
KinesisDataFetcher.getSubscribedShardsState() |
Constructor and Description |
---|
KinesisDataFetcher(List<String> streams,
SourceFunction.SourceContext<T> sourceContext,
Object checkpointLock,
RuntimeContext runtimeContext,
Properties configProps,
KinesisDeserializationSchema<T> deserializationSchema,
KinesisShardAssigner shardAssigner,
AssignerWithPeriodicWatermarks<T> periodicWatermarkAssigner,
WatermarkTracker watermarkTracker,
AtomicReference<Throwable> error,
List<KinesisStreamShardState> subscribedShardsState,
HashMap<String,String> subscribedStreamsToLastDiscoveredShardIds,
KinesisDataFetcher.FlinkKinesisProxyFactory kinesisProxyFactory) |
Modifier and Type | Method and Description |
---|---|
int |
CheckpointBarrierAligner.getNumClosedChannels() |
Modifier and Type | Method and Description |
---|---|
Evictor<? super IN,? super W> |
EvictingWindowOperator.getEvictor() |
KeySelector<IN,K> |
WindowOperator.getKeySelector() |
StateDescriptor<? extends AppendingState<IN,Iterable<IN>>,?> |
EvictingWindowOperator.getStateDescriptor() |
StateDescriptor<? extends AppendingState<IN,ACC>,?> |
WindowOperator.getStateDescriptor() |
Trigger<? super IN,? super W> |
WindowOperator.getTrigger() |
WindowAssigner<? super IN,W> |
WindowOperator.getWindowAssigner() |
Modifier and Type | Class and Description |
---|---|
protected static class |
StatusWatermarkValve.InputChannelStatus
An
InputChannelStatus keeps track of an input channel's last watermark, stream
status, and whether or not the channel's current watermark is aligned with the overall
watermark output from the valve. |
Modifier and Type | Method and Description |
---|---|
protected StatusWatermarkValve.InputChannelStatus |
StatusWatermarkValve.getInputChannelStatus(int channelIndex) |
Modifier and Type | Method and Description |
---|---|
static <OUT> RecordWriterDelegate<SerializationDelegate<StreamRecord<OUT>>> |
StreamTask.createRecordWriterDelegate(StreamConfig configuration,
Environment environment) |
boolean |
StreamTask.isMailboxLoopRunning() |
boolean |
StreamTask.runMailboxStep() |
Constructor and Description |
---|
OneInputStreamTask(Environment env,
TimerService timeProvider)
Constructor for initialization, possibly with initial state (recovery / savepoint / etc).
|
Modifier and Type | Method and Description |
---|---|
Meter |
MailboxProcessor.getIdleTime() |
boolean |
MailboxProcessor.hasMail() |
boolean |
MailboxProcessor.isDefaultActionUnavailable() |
boolean |
MailboxProcessor.isMailboxLoopRunning() |
boolean |
MailboxProcessor.runMailboxStep()
Execute a single (as small as possible) step of the mailbox.
|
Constructor and Description |
---|
TaskMailboxImpl() |
Modifier and Type | Method and Description |
---|---|
protected ExplainDetail[] |
TableEnvironmentImpl.getExplainDetails(boolean extended) |
Planner |
TableEnvironmentImpl.getPlanner() |
Modifier and Type | Method and Description |
---|---|
static CatalogManager.TableLookupResult |
CatalogManager.TableLookupResult.permanent(CatalogBaseTable table,
TableSchema resolvedSchema) |
static CatalogManager.TableLookupResult |
CatalogManager.TableLookupResult.temporary(CatalogBaseTable table,
TableSchema resolvedSchema) |
Constructor and Description |
---|
ConnectorCatalogTable(TableSource<T1> tableSource,
TableSink<T2> tableSink,
TableSchema tableSchema,
boolean isBatch) |
Modifier and Type | Method and Description |
---|---|
org.apache.hadoop.hive.conf.HiveConf |
HiveCatalog.getHiveConf() |
org.apache.hadoop.hive.metastore.api.Database |
HiveCatalog.getHiveDatabase(String databaseName) |
org.apache.hadoop.hive.metastore.api.Partition |
HiveCatalog.getHivePartition(org.apache.hadoop.hive.metastore.api.Table hiveTable,
CatalogPartitionSpec partitionSpec) |
org.apache.hadoop.hive.metastore.api.Table |
HiveCatalog.getHiveTable(ObjectPath tablePath) |
String |
HiveCatalog.getHiveVersion() |
static boolean |
HiveCatalog.isEmbeddedMetastore(org.apache.hadoop.hive.conf.HiveConf hiveConf) |
Constructor and Description |
---|
HiveCatalog(String catalogName,
String defaultDatabase,
org.apache.hadoop.hive.conf.HiveConf hiveConf,
String hiveVersion,
boolean allowEmbedded) |
Constructor and Description |
---|
CliClient(org.jline.terminal.Terminal terminal,
String sessionId,
Executor executor,
Path historyFilePath)
Creates a CLI instance with a custom terminal.
|
Modifier and Type | Method and Description |
---|---|
protected ExecutionContext<?> |
LocalExecutor.getExecutionContext(String sessionId)
Get the existed
ExecutionContext from contextMap, or thrown exception if does not
exist. |
Modifier and Type | Method and Description |
---|---|
protected List<Row> |
MaterializedCollectStreamResult.getMaterializedTable() |
Constructor and Description |
---|
MaterializedCollectStreamResult(TableSchema tableSchema,
ExecutionConfig config,
InetAddress gatewayAddress,
int gatewayPort,
int maxRowCount,
int overcommitThreshold,
ClassLoader classLoader) |
Constructor and Description |
---|
StreamExecutor(StreamExecutionEnvironment executionEnvironment) |
Modifier and Type | Method and Description |
---|---|
java.time.Duration |
FileSystemLookupFunction.getCacheTTL() |
Modifier and Type | Method and Description |
---|---|
protected void |
HiveGenericUDTF.setCollector(org.apache.hadoop.hive.ql.udf.generic.Collector collector) |
Constructor and Description |
---|
BatchExecutor(StreamExecutionEnvironment executionEnvironment) |
StreamExecutor(StreamExecutionEnvironment executionEnvironment) |
Modifier and Type | Method and Description |
---|---|
boolean |
ListAggWithRetractAggFunction.ListAggWithRetractAccumulator.equals(Object o) |
boolean |
ListAggWsWithRetractAggFunction.ListAggWsWithRetractAccumulator.equals(Object o) |
Modifier and Type | Method and Description |
---|---|
static Map<String,ColumnStats> |
CatalogTableStatisticsConverter.convertToColumnStatsMap(Map<String,CatalogColumnStatisticsDataBase> columnStatisticsData) |
Modifier and Type | Method and Description |
---|---|
void |
BinaryExternalSorter.write(MutableObjectIterator<BinaryRowData> iterator) |
Modifier and Type | Method and Description |
---|---|
abstract FlinkFnApi.UserDefinedFunctions |
AbstractPythonStatelessFunctionRunner.getUserDefinedFunctionsProto()
Gets the proto representation of the Python user-defined functions to be executed.
|
Modifier and Type | Method and Description |
---|---|
FlinkFnApi.UserDefinedFunctions |
AbstractPythonScalarFunctionRunner.getUserDefinedFunctionsProto()
Gets the proto representation of the Python user-defined functions to be executed.
|
Modifier and Type | Method and Description |
---|---|
FlinkFnApi.UserDefinedFunctions |
AbstractPythonTableFunctionRunner.getUserDefinedFunctionsProto()
Gets the proto representation of the Python user-defined functions to be executed.
|
Modifier and Type | Method and Description |
---|---|
TypeSerializer |
ArrayDataSerializer.getEleSer() |
TypeSerializer |
MapDataSerializer.getKeySerializer() |
TypeSerializer |
MapDataSerializer.getValueSerializer() |
Modifier and Type | Method and Description |
---|---|
int |
AbstractCloseableRegistry.getNumberOfRegisteredCloseables() |
boolean |
AbstractCloseableRegistry.isCloseableRegistered(Closeable c) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
YarnLogConfigUtil.setLogConfigFileInConfig(Configuration configuration,
String configurationDirectory) |
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.