Modifier and Type | Method and Description |
---|---|
String |
AbstractRuntimeUDFContext.getAllocationIDAsString() |
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 | Field and Description |
---|---|
static String |
PythonDriverEnvUtils.PYFLINK_PY_ARCHIVES |
static String |
PythonDriverEnvUtils.PYFLINK_PY_EXECUTABLE |
static String |
PythonDriverEnvUtils.PYFLINK_PY_FILES |
static String |
PythonDriverEnvUtils.PYFLINK_PY_REQUIREMENTS |
Modifier and Type | Method and Description |
---|---|
static Map<String,String> |
ConfigurationUtils.parseTmResourceDynamicConfigs(String dynamicConfigsStr) |
static Map<String,String> |
ConfigurationUtils.parseTmResourceJvmParams(String jvmParamsStr) |
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 |
---|---|
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(RefCountedFile file,
int bufferSize) |
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 |
---|---|
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 | Method and Description |
---|---|
boolean |
AbstractServerBase.isEventGroupShutdown() |
boolean |
Client.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 |
TransientBlobCache.getStorageLocation(JobID jobId,
BlobKey key)
Returns a file handle to the file associated with the given blob key on the blob
server.
|
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.
|
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.
|
CompletableFuture<CompletedCheckpoint> |
CheckpointCoordinator.triggerCheckpoint(long timestamp,
CheckpointProperties props,
String externalSavepointLocation,
boolean isPeriodic,
boolean advanceToEndOfTime) |
Constructor and Description |
---|
CheckpointCoordinator(JobID job,
CheckpointCoordinatorConfiguration chkConfig,
ExecutionVertex[] tasksToTrigger,
ExecutionVertex[] tasksToWaitFor,
ExecutionVertex[] tasksToCommitTo,
CheckpointIDCounter checkpointIDCounter,
CompletedCheckpointStore completedCheckpointStore,
StateBackend checkpointStateBackend,
Executor executor,
ScheduledExecutor timer,
SharedStateRegistryFactory sharedStateRegistryFactory,
CheckpointFailureManager failureManager,
Clock clock) |
Modifier and Type | Class and Description |
---|---|
class |
SavepointV2Serializer
(De)serializer for checkpoint metadata format version 2.
|
Modifier and Type | Method and Description |
---|---|
static KeyedStateHandle |
SavepointV2Serializer.deserializeKeyedStateHandle(DataInputStream dis) |
static KeyedStateHandle |
SavepointV1Serializer.deserializeKeyedStateHandle(DataInputStream dis) |
static OperatorStateHandle |
SavepointV2Serializer.deserializeOperatorStateHandle(DataInputStream dis) |
static OperatorStateHandle |
SavepointV1Serializer.deserializeOperatorStateHandle(DataInputStream dis) |
static StreamStateHandle |
SavepointV1Serializer.deserializeStreamStateHandle(DataInputStream dis) |
static void |
SavepointV2Serializer.serializeKeyedStateHandle(KeyedStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointV1Serializer.serializeKeyedStateHandle(KeyedStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointV2Serializer.serializeOperatorStateHandle(OperatorStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointV1Serializer.serializeOperatorStateHandle(OperatorStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointV2Serializer.serializeStreamStateHandle(StreamStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointV1Serializer.serializeStreamStateHandle(StreamStateHandle stateHandle,
DataOutputStream dos) |
static void |
SavepointSerializers.setFailWhenLegacyStateDetected(boolean fail)
This is only visible as a temporary solution to keep the stateful job migration it cases working from binary
savepoints that still contain legacy state (<= Flink 1.1).
|
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.
|
Constructor and Description |
---|
BlobLibraryCacheManager(PermanentBlobService blobService,
FlinkUserCodeClassLoaders.ResolveOrder classLoaderResolveOrder,
String[] alwaysParentFirstPatterns) |
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 |
---|---|
protected CompletableFuture<?> |
AdaptedRestartPipelinedRegionStrategyNG.cancelTasks(Set<ExecutionVertexID> vertices) |
protected void |
AdaptedRestartPipelinedRegionStrategyNG.restartTasks(Set<ExecutionVertexID> verticesToRestart) |
Modifier and Type | Method and Description |
---|---|
FailoverRegion |
RestartPipelinedRegionStrategy.getFailoverRegion(ExecutionVertexID vertexID)
Returns the failover region that contains the given execution vertex.
|
Constructor and Description |
---|
RestartPipelinedRegionStrategy(FailoverTopology<?,?> 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 |
---|---|
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. |
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,
ResultPartitionID id,
ResultPartitionType type,
int numberOfSubpartitions,
int maxParallelism,
FunctionWithException<BufferPoolOwner,BufferPool,IOException> bufferPoolFactory) |
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.
|
Buffer |
RemoteInputChannel.getNextReceivedBuffer() |
void |
RemoteInputChannel.requestSubpartition(int subpartitionIndex)
Requests a remote subpartition.
|
Constructor and Description |
---|
BufferOrEvent(AbstractEvent event,
int channelIndex) |
BufferOrEvent(Buffer buffer,
int channelIndex) |
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 |
---|---|
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 |
---|---|
int |
KvStateEntry.getCacheSize() |
Modifier and Type | Method and Description |
---|---|
abstract Collection<ResourceProfile> |
ResourceManager.startNewWorker(ResourceProfile resourceProfile)
Allocates a resource using the resource profile.
|
Modifier and Type | Method and Description |
---|---|
int |
SlotManagerImpl.getNumberAssignedPendingTaskManagerSlots() |
void |
SlotManagerImpl.unregisterTaskManagersAndReleaseResources() |
void |
SlotManager.unregisterTaskManagersAndReleaseResources() |
Modifier and Type | Method and Description |
---|---|
void |
MetricStore.add(MetricDump metric) |
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 | 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.
|
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 |
JobLeaderService.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 | Method and Description |
---|---|
static void |
Task.setupPartitionsAndGates(ResultPartitionWriter[] producedPartitions,
InputGate[] inputGates) |
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() |
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() |
<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 |
---|---|
int |
InternalTimeServiceManager.numEventTimeTimers() |
int |
InternalTimerServiceImpl.numEventTimeTimers() |
int |
AbstractStreamOperator.numEventTimeTimers() |
int |
InternalTimerServiceImpl.numEventTimeTimers(N namespace) |
int |
InternalTimeServiceManager.numProcessingTimeTimers() |
int |
InternalTimerServiceImpl.numProcessingTimeTimers() |
int |
AbstractStreamOperator.numProcessingTimeTimers() |
int |
InternalTimerServiceImpl.numProcessingTimeTimers(N namespace) |
Modifier and Type | Method and Description |
---|---|
TwoInputStreamOperator<IN1,IN2,OUT> |
TwoInputTransformation.getOperator() |
StreamSource<T,?> |
SourceTransformation.getOperator() |
OneInputStreamOperator<IN,OUT> |
OneInputTransformation.getOperator() |
StreamSink<T> |
SinkTransformation.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.NextTransactionalIdHintSerializer
|
static class |
FlinkKafkaProducer.TransactionStateSerializer
TypeSerializer for
FlinkKafkaProducer.KafkaTransactionState . |
static class |
FlinkKafkaProducer011.ContextStateSerializer
|
static class |
FlinkKafkaProducer011.KafkaTransactionContext
Context associated to this instance of the
FlinkKafkaProducer011 . |
static class |
FlinkKafkaProducer011.NextTransactionalIdHintSerializer
|
static class |
FlinkKafkaProducer011.TransactionStateSerializer
TypeSerializer for
KafkaTransactionState . |
Modifier and Type | Method and Description |
---|---|
protected <K,V> org.apache.kafka.clients.producer.KafkaProducer<K,V> |
FlinkKafkaProducerBase.getKafkaProducer(Properties props)
Used for testing only.
|
protected long |
FlinkKafkaProducerBase.numPendingRecords() |
Modifier and Type | Method and Description |
---|---|
protected org.apache.kafka.clients.consumer.ConsumerRecords<byte[],byte[]> |
KafkaConsumerThread.getRecordsFromKafka()
Get records from Kafka.
|
int |
FlinkKafkaInternalProducer.getTransactionCoordinatorId() |
int |
FlinkKafkaProducer.getTransactionCoordinatorId() |
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 |
---|---|
Evictor<? super IN,? super W> |
EvictingWindowOperator.getEvictor() |
KeySelector<IN,K> |
WindowOperator.getKeySelector() |
StateDescriptor<? extends AppendingState<IN,ACC>,?> |
WindowOperator.getStateDescriptor() |
StateDescriptor<? extends AppendingState<IN,Iterable<IN>>,?> |
EvictingWindowOperator.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 | Class and Description |
---|---|
protected static class |
StreamTask.AsyncCheckpointRunnable
This runnable executes the asynchronous parts of all involved backend snapshots for the subtask.
|
Modifier and Type | Method and Description |
---|---|
static <OUT> RecordWriterDelegate<SerializationDelegate<StreamRecord<OUT>>> |
StreamTask.createRecordWriterDelegate(StreamConfig configuration,
Environment environment) |
Constructor and Description |
---|
OneInputStreamTask(Environment env,
TimerService timeProvider)
Constructor for initialization, possibly with initial state (recovery / savepoint / etc).
|
Modifier and Type | Method and Description |
---|---|
boolean |
MailboxProcessor.isDefaultActionUnavailable() |
Constructor and Description |
---|
TaskMailboxImpl() |
Modifier and Type | Method and Description |
---|---|
Planner |
TableEnvironmentImpl.getPlanner() |
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.Table |
HiveCatalog.getHiveTable(ObjectPath tablePath) |
String |
HiveCatalog.getHiveVersion() |
protected static org.apache.hadoop.hive.metastore.api.Table |
HiveCatalog.instantiateHiveTable(ObjectPath tablePath,
CatalogBaseTable table) |
Constructor and Description |
---|
HiveCatalog(String catalogName,
String defaultDatabase,
org.apache.hadoop.hive.conf.HiveConf hiveConf,
String hiveVersion) |
Constructor and Description |
---|
CliClient(org.jline.terminal.Terminal terminal,
String sessionId,
Executor executor)
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 |
---|---|
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 |
---|---|
List<MemorySegment> |
BaseHybridHashTable.getFreedMemory() |
Modifier and Type | Method and Description |
---|---|
void |
BinaryExternalSorter.write(MutableObjectIterator<BinaryRow> iterator) |
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 |
---|---|
TypeSerializer |
BaseArraySerializer.getEleSer() |
TypeSerializer |
BaseMapSerializer.getKeySerializer() |
TypeSerializer |
BaseMapSerializer.getValueSerializer() |
Modifier and Type | Method and Description |
---|---|
int |
AbstractCloseableRegistry.getNumberOfRegisteredCloseables() |
boolean |
AbstractCloseableRegistry.isCloseableRegistered(Closeable c) |
Modifier and Type | Method and Description |
---|---|
static Configuration |
FlinkYarnSessionCli.setLogConfigFileInConfig(Configuration configuration,
String configurationDirectory) |
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.