Package | Description |
---|---|
org.apache.flink.api.common.accumulators | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.client | |
org.apache.flink.runtime.deployment | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.jobgraph | |
org.apache.flink.runtime.messages.accumulators | |
org.apache.flink.runtime.messages.checkpoint |
This package contains the messages that are sent between
JobManager
and TaskManager to coordinate the checkpoint snapshots of the
distributed dataflow. |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.streaming.connectors.kafka | |
org.apache.flink.streaming.connectors.kafka.internal | |
org.apache.flink.streaming.connectors.kafka.internals | |
org.apache.flink.util |
Modifier and Type | Method and Description |
---|---|
static Map<String,Object> |
AccumulatorHelper.deserializeAccumulators(Map<String,SerializedValue<Object>> serializedAccumulators,
ClassLoader loader)
Takes the serialized accumulator results and tries to deserialize them using the provided
class loader.
|
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
KeyGroupState.getKeyGroupState() |
SerializedValue<StateHandle<?>> |
SubtaskState.getState() |
Modifier and Type | Method and Description |
---|---|
Map<Integer,SerializedValue<StateHandle<?>>> |
TaskState.getUnwrappedKvStates(Set<Integer> keyGroupPartition)
Retrieve the set of key-value state key groups specified by the given key group partition set.
|
Modifier and Type | Method and Description |
---|---|
PendingCheckpoint.TaskAcknowledgeResult |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID executionAttemptId,
SerializedValue<StateHandle<?>> state,
long stateSize,
Map<Integer,SerializedValue<StateHandle<?>>> kvState) |
Modifier and Type | Method and Description |
---|---|
PendingCheckpoint.TaskAcknowledgeResult |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID executionAttemptId,
SerializedValue<StateHandle<?>> state,
long stateSize,
Map<Integer,SerializedValue<StateHandle<?>>> kvState) |
Constructor and Description |
---|
KeyGroupState(SerializedValue<StateHandle<?>> keyGroupState,
long stateSize,
long duration) |
SubtaskState(SerializedValue<StateHandle<?>> state,
long stateSize,
long duration) |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
SerializedJobExecutionResult.getSerializedAccumulatorResults() |
Constructor and Description |
---|
SerializedJobExecutionResult(JobID jobID,
long netRuntime,
Map<String,SerializedValue<Object>> accumulators)
Creates a new SerializedJobExecutionResult.
|
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
TaskDeploymentDescriptor.getOperatorState() |
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 |
---|
TaskDeploymentDescriptor(SerializedValue<JobInformation> serializedJobInformation,
SerializedValue<TaskInformation> serializedTaskInformation,
ExecutionAttemptID executionAttemptId,
int subtaskIndex,
int attemptNumber,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors) |
TaskDeploymentDescriptor(SerializedValue<JobInformation> serializedJobInformation,
SerializedValue<TaskInformation> serializedTaskInformation,
ExecutionAttemptID executionAttemptId,
int subtaskIndex,
int attemptNumber,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors) |
TaskDeploymentDescriptor(SerializedValue<JobInformation> serializedJobInformation,
SerializedValue<TaskInformation> serializedTaskInformation,
ExecutionAttemptID executionAttemptId,
int subtaskIndex,
int attemptNumber,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<ExecutionConfig> |
JobInformation.getSerializedExecutionConfig() |
SerializedValue<JobInformation> |
ExecutionGraph.getSerializedJobInformation() |
SerializedValue<TaskInformation> |
ExecutionJobVertex.getSerializedTaskInformation() |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
ExecutionGraph.getAccumulatorsSerialized()
Gets a serialized accumulator map.
|
Modifier and Type | Method and Description |
---|---|
void |
Execution.setInitialState(SerializedValue<StateHandle<?>> initialState,
Map<Integer,SerializedValue<StateHandle<?>>> initialKvState) |
Modifier and Type | Method and Description |
---|---|
void |
Execution.setInitialState(SerializedValue<StateHandle<?>> initialState,
Map<Integer,SerializedValue<StateHandle<?>>> initialKvState) |
Constructor and Description |
---|
ExecutionGraph(Executor futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
scala.concurrent.duration.FiniteDuration timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
Scheduler scheduler,
ClassLoader userClassLoader,
MetricGroup metricGroup) |
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<BlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<ExecutionConfig> |
JobGraph.getSerializedExecutionConfig()
Returns the
ExecutionConfig |
Modifier and Type | Method and Description |
---|---|
Map<String,SerializedValue<Object>> |
AccumulatorResultsFound.result() |
Constructor and Description |
---|
AccumulatorResultsFound(JobID jobID,
Map<String,SerializedValue<Object>> result) |
Modifier and Type | Method and Description |
---|---|
SerializedValue<StateHandle<?>> |
AcknowledgeCheckpoint.getState() |
Constructor and Description |
---|
AcknowledgeCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
SerializedValue<StateHandle<?>> state,
long stateSize) |
Constructor and Description |
---|
Task(JobInformation jobInformation,
TaskInformation taskInformation,
ExecutionAttemptID executionAttemptID,
int subtaskIndex,
int attemptNumber,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
MemoryManager memManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
BroadcastVariableManager bcVarManager,
ActorGateway taskManagerActor,
ActorGateway jobManagerActor,
scala.concurrent.duration.FiniteDuration actorAskTimeout,
LibraryCacheManager libraryCache,
FileCache fileCache,
TaskManagerRuntimeInfo taskManagerConfig,
TaskMetricGroup metricGroup)
IMPORTANT: This constructor may not start any work that would need to
be undone in the case of a failing task deployment.
|
Modifier and Type | Method and Description |
---|---|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer09.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer09.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the
data, and emits it into the data streams.
|
protected abstract AbstractFetcher<T,?> |
FlinkKafkaConsumerBase.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext)
Creates the fetcher that connect to the Kafka brokers, pulls data, deserialized the
data, and emits it into the data streams.
|
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer08.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
protected AbstractFetcher<T,?> |
FlinkKafkaConsumer08.createFetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> thisSubtaskPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext) |
Constructor and Description |
---|
Kafka09Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
boolean useMetrics) |
Kafka09Fetcher(SourceFunction.SourceContext<T> sourceContext,
List<KafkaTopicPartition> assignedPartitions,
SerializedValue<AssignerWithPeriodicWatermarks<T>> watermarksPeriodic,
SerializedValue<AssignerWithPunctuatedWatermarks<T>> watermarksPunctuated,
StreamingRuntimeContext runtimeContext,
KeyedDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
boolean useMetrics) |
Modifier and Type | Method and Description |
---|---|
static <T> SerializedValue<T> |
SerializedValue.fromBytes(byte[] serializedData) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.