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 |
---|---|
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).
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskManagerGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
JobMasterOperatorEventGateway.sendOperatorEventToCoordinator(ExecutionAttemptID task,
OperatorID operatorID,
SerializedValue<OperatorEvent> event) |
CompletableFuture<Acknowledge> |
JobMaster.sendOperatorEventToCoordinator(ExecutionAttemptID task,
OperatorID operatorID,
SerializedValue<OperatorEvent> serializedEvent) |
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
Modifier and Type | Class and Description |
---|---|
class |
AcknowledgeCheckpointEvent
An
OperatorEvent sent from a subtask to its OperatorCoordinator to signal that
the checkpoint of an individual task is completed. |
Modifier and Type | Method and Description |
---|---|
void |
RecreateOnResetOperatorCoordinator.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event) |
void |
OperatorCoordinator.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event)
Hands an OperatorEvent coming from a parallel Operator instance (one attempt of the parallel
subtasks).
|
void |
OperatorCoordinatorHolder.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event) |
void |
OperatorEventHandler.handleOperatorEvent(OperatorEvent evt) |
CompletableFuture<Acknowledge> |
OperatorCoordinator.SubtaskGateway.sendEvent(OperatorEvent evt)
Sends an event to the parallel subtask with the given subtask index.
|
void |
OperatorEventGateway.sendEventToCoordinator(OperatorEvent event)
Sends the given event to the coordinator, where it will be handled by the
OperatorCoordinator.handleEventFromOperator(int, int, OperatorEvent) method. |
Modifier and Type | Method and Description |
---|---|
void |
DefaultOperatorCoordinatorHandler.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecutionId,
OperatorID operatorId,
OperatorEvent evt) |
void |
SchedulerNG.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecution,
OperatorID operator,
OperatorEvent evt)
Delivers the given OperatorEvent to the
OperatorCoordinator with the given OperatorID . |
void |
SchedulerBase.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecutionId,
OperatorID operatorId,
OperatorEvent evt) |
void |
OperatorCoordinatorHandler.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecutionId,
OperatorID operatorId,
OperatorEvent event)
Delivers an OperatorEvent to a
OperatorCoordinator . |
Modifier and Type | Method and Description |
---|---|
void |
AdaptiveScheduler.deliverOperatorEventToCoordinator(ExecutionAttemptID taskExecution,
OperatorID operator,
OperatorEvent evt) |
Modifier and Type | Method and Description |
---|---|
void |
SourceCoordinator.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event) |
Modifier and Type | Class and Description |
---|---|
class |
AddSplitEvent<SplitT>
A source event that adds splits to a source reader.
|
class |
NoMoreSplitsEvent
A source event sent from the SplitEnumerator to the SourceReader to indicate that no more splits
will be assigned to the source reader anymore.
|
class |
ReaderRegistrationEvent
An
OperatorEvent that registers a SourceReader to the SourceCoordinator. |
class |
ReportedWatermarkEvent
Reports last emitted
Watermark from a subtask to the SourceCoordinator . |
class |
RequestSplitEvent
An event to request splits, sent typically from the Source Reader to the Source Enumerator.
|
class |
SourceEventWrapper
A wrapper operator event that contains a custom defined operator event.
|
class |
WatermarkAlignmentEvent
Signals source operators the maximum watermark that emitted records can have.
|
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) |
Modifier and Type | Method and Description |
---|---|
void |
Task.deliverOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> evt)
Dispatches an operator event to the invokable task.
|
Modifier and Type | Method and Description |
---|---|
void |
SourceOperator.handleOperatorEvent(OperatorEvent event) |
Modifier and Type | Class and Description |
---|---|
class |
CollectSinkAddressEvent
An
OperatorEvent that passes the socket server address in the sink to the coordinator. |
Modifier and Type | Method and Description |
---|---|
void |
CollectSinkOperatorCoordinator.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event) |
void |
CollectSinkOperator.handleOperatorEvent(OperatorEvent evt) |
Modifier and Type | Method and Description |
---|---|
abstract void |
OperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
FinishedOperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
RegularOperatorChain.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
void |
StreamTask.dispatchOperatorEvent(OperatorID operator,
SerializedValue<OperatorEvent> event) |
Modifier and Type | Method and Description |
---|---|
void |
DynamicFilteringDataCollectorOperatorCoordinator.handleEventFromOperator(int subtask,
int attemptNumber,
OperatorEvent event) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.