Modifier and Type | Method and Description |
---|---|
static Map<String,OptionalFailure<Object>> |
AccumulatorHelper.deserializeAccumulators(Map<String,SerializedValue<OptionalFailure<Object>>> serializedAccumulators,
ClassLoader loader)
Takes the serialized accumulator results and tries to deserialize them using the provided
class loader.
|
static Map<String,Object> |
AccumulatorHelper.deserializeAndUnwrapAccumulators(Map<String,SerializedValue<OptionalFailure<Object>>> serializedAccumulators,
ClassLoader loader)
Takes the serialized accumulator results and tries to deserialize them using the provided
class loader, and then try to unwrap the value unchecked.
|
Modifier and Type | Class and Description |
---|---|
class |
CachedSerializedValue<T>
An extension of SerializedValue which caches the deserialized data.
|
Modifier and Type | Method and Description |
---|---|
static <T> Either<SerializedValue<T>,PermanentBlobKey> |
BlobWriter.serializeAndTryOffload(T value,
JobID jobId,
BlobWriter blobWriter)
Serializes the given value and offloads it to the BlobServer if its size exceeds the minimum
offloading size of the BlobServer.
|
static <T> Either<SerializedValue<T>,PermanentBlobKey> |
BlobWriter.tryOffload(SerializedValue<T> serializedValue,
JobID jobId,
BlobWriter blobWriter) |
Modifier and Type | Method and Description |
---|---|
static <T> Either<SerializedValue<T>,PermanentBlobKey> |
BlobWriter.tryOffload(SerializedValue<T> serializedValue,
JobID jobId,
BlobWriter blobWriter) |
Modifier and Type | Method and Description |
---|---|
static SerializedValue<TaskStateSnapshot> |
TaskStateSnapshot.serializeTaskStateSnapshot(TaskStateSnapshot subtaskState) |
Modifier and Type | Method and Description |
---|---|
void |
CheckpointCoordinatorGateway.acknowledgeCheckpoint(JobID jobID,
ExecutionAttemptID executionAttemptID,
long checkpointId,
CheckpointMetrics checkpointMetrics,
SerializedValue<TaskStateSnapshot> subtaskState) |
static TaskStateSnapshot |
TaskStateSnapshot.deserializeTaskStateSnapshot(SerializedValue<TaskStateSnapshot> subtaskState,
ClassLoader classLoader) |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<OptionalFailure<Object>>> |
SerializedJobExecutionResult.getSerializedAccumulatorResults() |
Constructor and Description |
---|
SerializedJobExecutionResult(JobID jobID,
long netRuntime,
Map<String,SerializedValue<OptionalFailure<Object>>> accumulators)
Creates a new SerializedJobExecutionResult.
|
Modifier and Type | Field and Description |
---|---|
SerializedValue<T> |
TaskDeploymentDescriptor.NonOffloaded.serializedValue
The serialized value.
|
Modifier and Type | Method and Description |
---|---|
SerializedValue<JobInformation> |
TaskDeploymentDescriptor.getSerializedJobInformation()
Return the sub task's serialized job information.
|
SerializedValue<TaskInformation> |
TaskDeploymentDescriptor.getSerializedTaskInformation()
Return the sub task's serialized task information.
|
Constructor and Description |
---|
NonOffloaded(SerializedValue<T> serializedValue) |
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 |
---|---|
SerializedValue<ExecutionConfig> |
JobInformation.getSerializedExecutionConfig() |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<OptionalFailure<Object>>> |
AccessExecutionGraph.getAccumulatorsSerialized()
Returns a map containing the serialized values of user-defined accumulators.
|
Map<String,SerializedValue<OptionalFailure<Object>>> |
DefaultExecutionGraph.getAccumulatorsSerialized()
Gets a serialized accumulator map.
|
Map<String,SerializedValue<OptionalFailure<Object>>> |
ArchivedExecutionGraph.getAccumulatorsSerialized() |
Either<SerializedValue<JobInformation>,PermanentBlobKey> |
DefaultExecutionGraph.getJobInformationOrBlobKey() |
Either<SerializedValue<JobInformation>,PermanentBlobKey> |
InternalExecutionGraphAccessor.getJobInformationOrBlobKey() |
Either<SerializedValue<TaskInformation>,PermanentBlobKey> |
ExecutionJobVertex.getTaskInformationOrBlobKey() |
Modifier and Type | Method and Description |
---|---|
protected OperatorCoordinatorHolder |
ExecutionJobVertex.createOperatorCoordinatorHolder(SerializedValue<OperatorCoordinator.Provider> provider,
ClassLoader classLoader,
CoordinatorStore coordinatorStore) |
protected OperatorCoordinatorHolder |
SpeculativeExecutionJobVertex.createOperatorCoordinatorHolder(SerializedValue<OperatorCoordinator.Provider> provider,
ClassLoader classLoader,
CoordinatorStore coordinatorStore) |
CompletableFuture<Acknowledge> |
Execution.sendOperatorEvent(OperatorID operatorId,
SerializedValue<OperatorEvent> event)
Sends the operator event to the Task on the Task Executor.
|
Constructor and Description |
---|
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<PermanentBlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
Constructor and Description |
---|
ArchivedExecutionGraph(JobID jobID,
String jobName,
Map<JobVertexID,ArchivedExecutionJobVertex> tasks,
List<ArchivedExecutionJobVertex> verticesInCreationOrder,
long[] stateTimestamps,
JobStatus state,
ErrorInfo failureCause,
String jsonPlan,
StringifiedAccumulatorResult[] archivedUserAccumulators,
Map<String,SerializedValue<OptionalFailure<Object>>> serializedUserAccumulators,
ArchivedExecutionConfig executionConfig,
boolean isStoppable,
CheckpointCoordinatorConfiguration jobCheckpointingConfiguration,
CheckpointStatsSnapshot checkpointStatsSnapshot,
String stateBackendName,
String checkpointStorageName,
TernaryBoolean stateChangelogEnabled,
String changelogStorageName) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<ExecutionConfig> |
JobGraph.getSerializedExecutionConfig()
Returns the
ExecutionConfig . |
Modifier and Type | Method and Description |
---|---|
List<SerializedValue<OperatorCoordinator.Provider>> |
JobVertex.getOperatorCoordinators() |
Modifier and Type | Method and Description |
---|---|
void |
JobVertex.addOperatorCoordinator(SerializedValue<OperatorCoordinator.Provider> serializedCoordinatorProvider) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<CheckpointStorage> |
JobCheckpointingSettings.getDefaultCheckpointStorage() |
SerializedValue<StateBackend> |
JobCheckpointingSettings.getDefaultStateBackend() |
SerializedValue<MasterTriggerRestoreHook.Factory[]> |
JobCheckpointingSettings.getMasterHooks() |
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 |
---|---|
Map<String,SerializedValue<OptionalFailure<Object>>> |
JobResult.getAccumulatorResults() |
Modifier and Type | Method and Description |
---|---|
JobResult.Builder |
JobResult.Builder.accumulatorResults(Map<String,SerializedValue<OptionalFailure<Object>>> accumulatorResults) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<CoordinationResponse> |
MiniCluster.deliverCoordinationRequestToCoordinator(JobID jobId,
OperatorID operatorId,
SerializedValue<CoordinationRequest> serializedRequest) |
Modifier and Type | Method and Description |
---|---|
static OperatorCoordinatorHolder |
OperatorCoordinatorHolder.create(SerializedValue<OperatorCoordinator.Provider> serializedProvider,
ExecutionJobVertex jobVertex,
ClassLoader classLoader,
CoordinatorStore coordinatorStore,
boolean supportsConcurrentExecutionAttempts) |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<OptionalFailure<Object>>> |
JobAccumulatorsInfo.getSerializedUserAccumulators() |
Constructor and Description |
---|
JobAccumulatorsInfo(List<JobAccumulatorsInfo.JobAccumulator> jobAccumulators,
List<JobAccumulatorsInfo.UserTaskAccumulator> userAccumulators,
Map<String,SerializedValue<OptionalFailure<Object>>> serializedUserAccumulators) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<CoordinationRequest> |
ClientCoordinationRequestBody.getSerializedCoordinationRequest() |
SerializedValue<CoordinationResponse> |
ClientCoordinationResponseBody.getSerializedCoordinationResponse() |
Constructor and Description |
---|
ClientCoordinationRequestBody(SerializedValue<CoordinationRequest> serializedCoordinationRequest) |
ClientCoordinationResponseBody(SerializedValue<CoordinationResponse> serializedCoordinationResponse) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<?> |
SerializedValueDeserializer.deserialize(org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonParser p,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.DeserializationContext ctxt) |
Modifier and Type | Method and Description |
---|---|
void |
SerializedValueSerializer.serialize(SerializedValue<?> value,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonGenerator gen,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.SerializerProvider provider) |
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.
|
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.
|
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics)
Deprecated.
|
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> subscribedPartitionsToStartOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup kafkaMetricGroup,
boolean useMetrics)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the data, and
emits it into the data streams.
|
Constructor and Description |
---|
AbstractFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> seedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaShuffleFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics,
TypeSerializer<T> typeSerializer,
int producerParallelism) |
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaShuffleConsumer.createFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
StreamingRuntimeContext runtimeContext,
OffsetCommitMode offsetCommitMode,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
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) |
Modifier and Type | Class and Description |
---|---|
class |
CompressedSerializedValue<T>
An extension of
SerializedValue that compresses the value after the serialization. |
Modifier and Type | Method and Description |
---|---|
static <T> SerializedValue<T> |
SerializedValue.fromBytes(byte[] serializedData)
Constructs serialized value from serialized data.
|
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.