Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
ClusterClientJobClientAdapter.sendCoordinationRequest(OperatorID operatorId,
CoordinationRequest request) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
EmbeddedJobClient.sendCoordinationRequest(OperatorID operatorId,
CoordinationRequest request) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
ClusterClient.sendCoordinationRequest(JobID jobId,
OperatorID operatorId,
CoordinationRequest request)
Sends out a request to a specified coordinator and return the response.
|
CompletableFuture<CoordinationResponse> |
MiniClusterClient.sendCoordinationRequest(JobID jobId,
OperatorID operatorId,
CoordinationRequest request) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
RestClusterClient.sendCoordinationRequest(JobID jobId,
OperatorID operatorId,
CoordinationRequest request) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
OperatorIDPair.getGeneratedOperatorID() |
Modifier and Type | Method and Description |
---|---|
Optional<OperatorID> |
OperatorIDPair.getUserDefinedOperatorID() |
Modifier and Type | Method and Description |
---|---|
static OperatorIDPair |
OperatorIDPair.generatedIDOnly(OperatorID generatedOperatorID) |
static OperatorIDPair |
OperatorIDPair.of(OperatorID generatedOperatorID,
OperatorID userDefinedOperatorID) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
OperatorState.getOperatorID() |
Modifier and Type | Method and Description |
---|---|
Map<OperatorID,OperatorState> |
CompletedCheckpoint.getOperatorStates() |
Map<OperatorID,OperatorState> |
PendingCheckpoint.getOperatorStates() |
Set<Map.Entry<OperatorID,OperatorSubtaskState>> |
TaskStateSnapshot.getSubtaskStateMappings()
Returns the set of all mappings from operator id to the corresponding subtask state.
|
Modifier and Type | Method and Description |
---|---|
static <T> Map<OperatorInstanceID,List<T>> |
StateAssignmentOperation.applyRepartitioner(OperatorID operatorID,
OperatorStateRepartitioner<T> opStateRepartitioner,
List<List<T>> chainOpParallelStates,
int oldParallelism,
int newParallelism) |
OperatorSubtaskState |
TaskStateSnapshot.getSubtaskStateByOperatorID(OperatorID operatorID)
Returns the subtask state for the given operator id (or null if not contained).
|
OperatorSubtaskState |
TaskStateSnapshot.putSubtaskStateByOperatorID(OperatorID operatorID,
OperatorSubtaskState state)
Maps the given operator id to the given subtask state.
|
Modifier and Type | Method and Description |
---|---|
void |
DefaultCheckpointPlan.fulfillFinishedTaskStatus(Map<OperatorID,OperatorState> operatorStates) |
void |
FinishedTaskStateProvider.fulfillFinishedTaskStatus(Map<OperatorID,OperatorState> operatorStates)
Fulfills the state for the finished subtasks and operators to indicate they are finished.
|
static <T extends StateObject> |
StateAssignmentOperation.reDistributePartitionableStates(Map<OperatorID,OperatorState> oldOperatorStates,
int newParallelism,
java.util.function.Function<OperatorSubtaskState,StateObjectCollection<T>> extractHandle,
OperatorStateRepartitioner<T> stateRepartitioner,
Map<OperatorInstanceID,List<T>> result) |
Constructor and Description |
---|
FullyFinishedOperatorState(OperatorID operatorID,
int parallelism,
int maxParallelism) |
OperatorState(OperatorID operatorID,
int parallelism,
int maxParallelism) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
Dispatcher.deliverCoordinationRequestToCoordinator(JobID jobId,
OperatorID operatorId,
SerializedValue<CoordinationRequest> serializedRequest,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
Execution.sendOperatorEvent(OperatorID operatorId,
SerializedValue<OperatorEvent> event)
Sends the operator event to the Task on the Task Executor.
|
Modifier and Type | Method and Description |
---|---|
static OperatorID |
OperatorID.fromJobVertexID(JobVertexID id) |
OperatorID |
OperatorInstanceID.getOperatorId() |
Modifier and Type | Method and Description |
---|---|
Map<OperatorID,UserCodeWrapper<? extends InputFormat<?,?>>> |
InputOutputFormatContainer.getInputFormats() |
Map<OperatorID,UserCodeWrapper<? extends InputFormat<?,?>>> |
InputOutputFormatContainer.FormatUserCodeTable.getInputFormats() |
Map<OperatorID,UserCodeWrapper<? extends OutputFormat<?>>> |
InputOutputFormatContainer.getOutputFormats() |
Map<OperatorID,UserCodeWrapper<? extends OutputFormat<?>>> |
InputOutputFormatContainer.FormatUserCodeTable.getOutputFormats() |
<OT,T extends InputSplit> |
InputOutputFormatContainer.getUniqueInputFormat() |
<IT> org.apache.commons.lang3.tuple.Pair<OperatorID,OutputFormat<IT>> |
InputOutputFormatContainer.getUniqueOutputFormat() |
Modifier and Type | Method and Description |
---|---|
InputOutputFormatContainer |
InputOutputFormatContainer.addInputFormat(OperatorID operatorId,
InputFormat<?,?> inputFormat) |
InputOutputFormatContainer |
InputOutputFormatContainer.addInputFormat(OperatorID operatorId,
UserCodeWrapper<? extends InputFormat<?,?>> wrapper) |
void |
InputOutputFormatContainer.FormatUserCodeTable.addInputFormat(OperatorID operatorId,
UserCodeWrapper<? extends InputFormat<?,?>> wrapper) |
InputOutputFormatContainer |
InputOutputFormatContainer.addOutputFormat(OperatorID operatorId,
OutputFormat<?> outputFormat) |
InputOutputFormatContainer |
InputOutputFormatContainer.addOutputFormat(OperatorID operatorId,
UserCodeWrapper<? extends OutputFormat<?>> wrapper) |
void |
InputOutputFormatContainer.FormatUserCodeTable.addOutputFormat(OperatorID operatorId,
UserCodeWrapper<? extends OutputFormat<?>> wrapper) |
InputOutputFormatContainer |
InputOutputFormatContainer.addParameters(OperatorID operatorId,
Configuration parameters) |
InputOutputFormatContainer |
InputOutputFormatContainer.addParameters(OperatorID operatorId,
String key,
String value) |
String |
InputOutputFormatVertex.getFormatDescription(OperatorID operatorID) |
Configuration |
InputOutputFormatContainer.getParameters(OperatorID operatorId) |
static OperatorInstanceID |
OperatorInstanceID.of(int subtaskId,
OperatorID operatorID) |
void |
InputOutputFormatVertex.setFormatDescription(OperatorID operatorID,
String formatDescription) |
Constructor and Description |
---|
OperatorInstanceID(int subtaskId,
OperatorID operatorId) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractInvokable.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
CoordinatedTask.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
TaskOperatorEventGateway.sendOperatorEventToCoordinator(OperatorID operator,
SerializedValue<OperatorEvent> event)
Sends an event from the operator (identified by the given operator ID) to the operator
coordinator (identified by the same ID).
|
CompletableFuture<CoordinationResponse> |
TaskOperatorEventGateway.sendRequestToCoordinator(OperatorID operator,
SerializedValue<CoordinationRequest> request)
Sends a request from current operator to a specified operator coordinator which is identified
by the given operator ID and return the response.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskManagerGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
Modifier and Type | Method and Description |
---|---|
InternalOperatorMetricGroup |
UnregisteredMetricGroups.UnregisteredTaskMetricGroup.getOrAddOperator(OperatorID operatorID,
String name) |
InternalOperatorMetricGroup |
TaskMetricGroup.getOrAddOperator(OperatorID operatorID,
String operatorName) |
Modifier and Type | Method and Description |
---|---|
String[] |
OperatorScopeFormat.formatScope(TaskMetricGroup parent,
OperatorID operatorID,
String operatorName) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
MiniCluster.deliverCoordinationRequestToCoordinator(JobID jobId,
OperatorID operatorId,
SerializedValue<CoordinationRequest> serializedRequest) |
CompletableFuture<CoordinationResponse> |
MiniClusterJobClient.sendCoordinationRequest(OperatorID operatorId,
CoordinationRequest request) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
OperatorCoordinator.Context.getOperatorId()
Gets the ID of the operator to which the coordinator belongs.
|
OperatorID |
OperatorCoordinator.Provider.getOperatorId()
Gets the ID of the operator to which the coordinator belongs.
|
OperatorID |
RecreateOnResetOperatorCoordinator.Provider.getOperatorId() |
OperatorID |
OperatorInfo.operatorId() |
OperatorID |
OperatorCoordinatorHolder.operatorId() |
Modifier and Type | Method and Description |
---|---|
static Collection<OperatorID> |
OperatorInfo.getIds(Collection<? extends OperatorInfo> infos) |
Modifier and Type | Method and Description |
---|---|
OperatorEventGateway |
OperatorEventDispatcher.getOperatorEventGateway(OperatorID operatorId)
Gets the gateway through which events can be passed to the OperatorCoordinator for the
operator identified by the given OperatorID.
|
void |
OperatorEventDispatcher.registerEventHandler(OperatorID operator,
OperatorEventHandler handler)
Register a listener that is notified every time an OperatorEvent is sent from the
OperatorCoordinator (of the operator with the given OperatorID) to this subtask.
|
CompletableFuture<CoordinationResponse> |
CoordinationRequestGateway.sendCoordinationRequest(OperatorID operatorId,
CoordinationRequest request)
Send out a request to a specified coordinator and return the response.
|
Constructor and Description |
---|
Provider(OperatorID operatorID) |
Modifier and Type | Method and Description |
---|---|
protected OperatorID |
OperatorIDPathParameter.convertFromString(String value) |
Modifier and Type | Method and Description |
---|---|
protected String |
OperatorIDPathParameter.convertToString(OperatorID value) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
AdaptiveScheduler.deliverCoordinationRequestToCoordinator(OperatorID operator,
CoordinationRequest request) |
void |
AdaptiveScheduler.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecution,
OperatorID operator,
OperatorEvent evt) |
Constructor and Description |
---|
SourceCoordinatorProvider(String operatorName,
OperatorID operatorID,
Source<?,SplitT,?> source,
int numWorkerThreads,
WatermarkAlignmentParams alignmentParams,
String coordinatorListeningID)
Construct the
SourceCoordinatorProvider . |
Modifier and Type | Method and Description |
---|---|
PrioritizedOperatorSubtaskState |
TaskStateManagerImpl.prioritizedOperatorState(OperatorID operatorID) |
PrioritizedOperatorSubtaskState |
TaskStateManager.prioritizedOperatorState(OperatorID operatorID)
Returns means to restore previously reported state of an operator running in the owning task.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutor.sendOperatorEventToTask(ExecutionAttemptID executionAttemptID,
OperatorID operatorId,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutorOperatorEventGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt)
Sends an operator event to an operator in a task executed by the Task Manager (Task
Executor).
|
Modifier and Type | Method and Description |
---|---|
void |
RpcTaskOperatorEventGateway.sendOperatorEventToCoordinator(OperatorID operator,
SerializedValue<OperatorEvent> event) |
CompletableFuture<CoordinationResponse> |
RpcTaskOperatorEventGateway.sendRequestToCoordinator(OperatorID operator,
SerializedValue<CoordinationRequest> request) |
Modifier and Type | Method and Description |
---|---|
void |
Task.deliverOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> evt)
Dispatches an operator event to the invokable task.
|
Constructor and Description |
---|
OperatorSubtaskDescriptionText(OperatorID operatorId,
String operatorClass,
int subtaskIndex,
int numberOfTasks) |
Modifier and Type | Method and Description |
---|---|
default CompletableFuture<CoordinationResponse> |
RestfulGateway.deliverCoordinationRequestToCoordinator(JobID jobId,
OperatorID operatorId,
SerializedValue<CoordinationRequest> serializedRequest,
Time timeout)
Deliver a coordination request to a specified coordinator and return the response.
|
Constructor and Description |
---|
OperatorSubtaskStateReducer(OperatorID operatorID,
int maxParallelism) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
StateBootstrapWrapperOperator.getOperatorID() |
Modifier and Type | Method and Description |
---|---|
static OperatorID |
OperatorIDGenerator.fromUid(String uid)
Generate
OperatorID 's from uid 's. |
OperatorID |
StateBootstrapTransformationWithID.getOperatorID() |
OperatorID |
BootstrapTransformationWithID.getOperatorID()
Deprecated.
|
Constructor and Description |
---|
BootstrapTransformationWithID(OperatorID operatorID,
BootstrapTransformation<T> bootstrapTransformation)
Deprecated.
|
StateBootstrapTransformationWithID(OperatorID operatorID,
StateBootstrapTransformation<T> bootstrapTransformation) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
StreamConfig.getOperatorID() |
Modifier and Type | Method and Description |
---|---|
Optional<OperatorCoordinator.Provider> |
StreamNode.getCoordinatorProvider(String operatorName,
OperatorID operatorID) |
void |
StreamConfig.setOperatorID(OperatorID operatorID) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
AbstractStreamOperatorV2.getOperatorID() |
OperatorID |
AbstractStreamOperator.getOperatorID() |
OperatorID |
StreamOperator.getOperatorID() |
Modifier and Type | Method and Description |
---|---|
OperatorCoordinator.Provider |
SourceOperatorFactory.getCoordinatorProvider(String operatorName,
OperatorID operatorID) |
OperatorCoordinator.Provider |
CoordinatedOperatorFactory.getCoordinatorProvider(String operatorName,
OperatorID operatorID)
Get the operator coordinator provider for this operator.
|
StreamOperatorStateContext |
StreamTaskStateInitializerImpl.streamOperatorStateContext(OperatorID operatorID,
String operatorClassName,
ProcessingTimeService processingTimeService,
KeyContext keyContext,
TypeSerializer<?> keySerializer,
CloseableRegistry streamTaskCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction,
boolean isUsingCustomRawKeyedState) |
StreamOperatorStateContext |
StreamTaskStateInitializer.streamOperatorStateContext(OperatorID operatorID,
String operatorClassName,
ProcessingTimeService processingTimeService,
KeyContext keyContext,
TypeSerializer<?> keySerializer,
CloseableRegistry streamTaskCloseableRegistry,
MetricGroup metricGroup,
double managedMemoryFraction,
boolean isUsingCustomRawKeyedState)
Returns the
StreamOperatorStateContext for an AbstractStreamOperator that
runs in the stream task that owns this manager. |
Constructor and Description |
---|
StreamingRuntimeContext(Environment env,
Map<String,Accumulator<?,?>> accumulators,
OperatorMetricGroup operatorMetricGroup,
OperatorID operatorID,
ProcessingTimeService processingTimeService,
KeyedStateStore keyedStateStore,
ExternalResourceInfoProvider externalResourceInfoProvider) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
CollectSinkOperatorCoordinator.Provider.getOperatorId() |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<OperatorID> |
CollectSinkOperator.getOperatorIdFuture() |
Modifier and Type | Method and Description |
---|---|
OperatorCoordinator.Provider |
CollectSinkOperatorFactory.getCoordinatorProvider(String operatorName,
OperatorID operatorID) |
Constructor and Description |
---|
Provider(OperatorID operatorId,
int socketTimeout) |
Constructor and Description |
---|
CollectResultFetcher(AbstractCollectResultBuffer<T> buffer,
CompletableFuture<OperatorID> operatorIdFuture,
String accumulatorName) |
CollectResultIterator(AbstractCollectResultBuffer<T> buffer,
CompletableFuture<OperatorID> operatorIdFuture,
String accumulatorName,
int retryMillis) |
CollectResultIterator(CompletableFuture<OperatorID> operatorIdFuture,
TypeSerializer<T> serializer,
String accumulatorName,
CheckpointConfig checkpointConfig) |
Modifier and Type | Method and Description |
---|---|
OperatorID |
StreamTaskSourceInput.getOperatorID() |
Modifier and Type | Method and Description |
---|---|
OperatorID |
LatencyMarker.getOperatorId() |
Constructor and Description |
---|
LatencyMarker(long markedTime,
OperatorID operatorId,
int subtaskIndex)
Creates a latency mark with the given timestamp.
|
Modifier and Type | Method and Description |
---|---|
void |
FinishedOperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
RegularOperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
abstract void |
OperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
StreamTask.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
OperatorEventGateway |
OperatorEventDispatcherImpl.getOperatorEventGateway(OperatorID operatorId) |
void |
OperatorEventDispatcherImpl.registerEventHandler(OperatorID operator,
OperatorEventHandler handler) |
Constructor and Description |
---|
LatencyStats(MetricGroup metricGroup,
int historySize,
int subtaskIndex,
OperatorID operatorID,
LatencyStats.Granularity granularity) |
Modifier and Type | Method and Description |
---|---|
OperatorCoordinator.Provider |
DynamicFilteringDataCollectorOperatorFactory.getCoordinatorProvider(String operatorName,
OperatorID operatorID) |
Constructor and Description |
---|
Provider(OperatorID operatorID,
List<String> dynamicFilteringDataListenerIDs) |
Constructor and Description |
---|
ScriptProcessBuilder(String script,
org.apache.hadoop.mapred.JobConf hiveConf,
OperatorID operatorID) |
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.