nictownsend commented on code in PR #28863:
URL: https://github.com/apache/flink/pull/28863#discussion_r3912580714
##########
docs/content/docs/concepts/glossary.md:
##########
@@ -46,137 +48,331 @@ Cluster](#flink-cluster) is bound to the lifetime of the
Flink Application.
#### ApplicationResultStore
The ApplicationResultStore is a Flink component that persists the results of
terminated
-(i.e. finished, cancelled or failed) applications to a filesystem, allowing
the results to outlive
-a terminated application. Each result contains the application's identifier,
final state, name,
-etc. These results are then used by Flink to determine whether applications
should
-be subject to recovery in highly-available clusters.
+(i.e. finished, cancelled or failed) Applications to a filesystem, allowing
the results to outlive
+a terminated Application. Each result contains the Application's identifier,
final state, name,
+etc. These results are then used by Flink to determine whether Applications
should
+be subject to recovery in highly-available Clusters.
+
+#### Channel
+
+Also called *Stream Partitions*.
+
+A Channel is the physical link between a [Sub-Task](#sub-task) and a
downstream Sub-Task, and the
+edge of a [Physical Graph](#physical-graph). Parts of the documentation refer
to Channels as *Stream
+Partitions*, in the sense of internal, physical Partitions.
+
+Channels carry data records as well as signals such as
[Watermarks](#watermark), Watermark Status
+updates and Checkpoint barriers. Transmission over a Channel is always
unidirectional (upstream to
+downstream) and asynchronous.
+
+A Sub-Task may have one or more input Channels and one or more output
Channels. Source Sub-Tasks have
+no input Channels, since they begin the graph, and Sink Sub-Tasks have no
output Channels, since they
+end it.
+
+A Sub-Task routes each record to one of its output Channels according to the
[Physical
+Partitioning](#partition) of the stream. Hash partitioning (`keyBy()` in the
DataStream API, `GROUP
+BY` in SQL) routes a record to the Channel connected to the downstream
Sub-Task that handles the
+record's key, whereas `rebalance()` or `rescale()` may round-robin records
across output Channels.
+
+A Channel is *local* when both Sub-Tasks run in the same [Flink
+TaskManager](#flink-taskmanager), in which case records are handed over
through an in-memory buffer,
+or *remote* when the Sub-Tasks run in different TaskManagers, in which case
the data crosses the
+network.
+
+#### Checkpoint
+
+A consistent snapshot of the State of a [Flink Job](#flink-job) at a logical
point in time, taken
+with a variant of the Chandy-Lamport algorithm and written to [Checkpoint
+Storage](#checkpoint-storage).
+
+A Checkpoint contains the [State](#managed-state) of all stateful
[Operators](#operator). This also
+includes source positions (for example Kafka partition offsets), assignment of
[Source Splits](#source-split)
+to [Sub-Tasks](#sub-task), and Sink transaction metadata. Async I/O in-flight
data and buffered data
+of some asynchronous Sink connectors are also part of the
[Operator](#operator) [State](#managed-state)
+and are saved in the Checkpoint.
+When [Unaligned Checkpoints]({{< ref
"docs/concepts/stateful-stream-processing" >}}#unaligned-checkpointing)
+are enabled, it may also contain data in flight between Sub-Tasks.
+
+Checkpoints are triggered automatically and periodically while the Job is
running, and are used to
+recover from failures such as a TaskManager crash or a network problem: the
Job restarts from the
+latest completed Checkpoint. They are designed for low overhead and run mostly
asynchronously,
+without blocking record processing, apart from a synchronous phase in each
Sub-Task.
+Transactional [Sources and Sinks](#operator) tie their transactions to the
Checkpoint; the
+Kafka Sink, for instance, commits its Kafka transactions when a Checkpoint
completes.
+
+Checkpoints are only used in the `STREAMING` [Execution
Mode](#runtime-execution-mode). In `BATCH`
+mode, Flink recovers instead by backtracking to previous processing stages
whose intermediate results
+are still available, so that potentially only the failed [Tasks](#task) and
their predecessors are
+restarted. As a consequence, Sinks that rely on Checkpoints to commit their
transactions do not work
+in `BATCH` mode unless they are implemented with the Unified Sink API, which
commits once the whole
+input has been processed.
+
+Compare to [Savepoint](#savepoint).
+
+Checkpoints and [Savepoints](#savepoint) are also referred to, collectively,
as *State Snapshots* or
+*Snapshots*.
#### Checkpoint Storage
-The location where the [State Backend](#state-backend) will store its snapshot
during a checkpoint (Java Heap of [JobManager](#flink-jobmanager) or
Filesystem).
+The durable location where [Checkpoints](#checkpoint) and
[Savepoints](#savepoint) are saved. It can
+be either the Java Heap of the [Flink JobManager](#flink-jobmanager) or a
filesystem. Production
+deployments use a filesystem, typically remote object storage, since
Checkpoint Storage is what makes
+State survive the loss of a [TaskManager](#flink-taskmanager) or of the whole
[Flink Cluster](#flink-cluster).
+
+The relationship between the State Backend and Checkpoint Storage changes with
[Disaggregated
+State]({{< ref "docs/ops/state/disaggregated_state" >}}), where remote storage
becomes the primary
+location of the State and the local State Backend acts as a cache, the two
being synchronized
+asynchronously.
#### Flink Cluster
A distributed system consisting of (typically) one
[JobManager](#flink-jobmanager) and one or more
-[Flink TaskManager](#flink-taskmanager) processes.
+[Flink TaskManager](#flink-taskmanager) processes. Each of these processes
runs in a separate JVM,
+usually on a separate container or machine, although this is not a requirement.
+
+See also [Flink Architecture: Anatomy of a Flink Cluster]({{< ref
"docs/concepts/flink-architecture" >}}#anatomy-of-a-flink-cluster).
#### Event
-An event is a statement about a change of the state of the domain modelled by
the
-application. Events can be input and/or output of a stream or batch processing
application.
-Events are special types of [records](#Record).
+An Event is a statement about a change of the state of the domain modeled by
the
+Application. Events can be input and/or output of a stream or batch processing
Application.
+Events are special types of records.
+
+#### Execution Graph
-#### ExecutionGraph
+Also called *ExecutionGraph*.
-see [Physical Graph](#physical-graph)
+See [Physical Graph](#physical-graph)
#### Function
-Functions are implemented by the user and encapsulate the
+Functions are implemented by the user, in Java or Python, and encapsulate the
application logic of a Flink program. Most Functions are wrapped by a
corresponding
-[Operator](#operator).
+[Operator](#operator). In the DataStream API, Functions are passed to the
+[Transformations](#transformation) they implement. In the Table API and SQL,
they are declared
+separately as [User-Defined Functions]({{< ref "docs/dev/table/functions/udfs"
>}}) (UDF) or
+[Process Table Functions]({{< ref "docs/dev/table/functions/ptfs" >}}) (PTF).
#### History Server
The History Server is a standalone service that serves the detailed history of
completed Flink
-applications and jobs, using archives generated by the JobManager. Unlike the
-[ApplicationResultStore](#applicationresultstore) and
[JobResultStore](#jobresultstore), which store
-minimal metadata for internal recovery decisions in highly-available clusters,
the History Server
-provides detailed archives for analysis via Web UI or REST API after the
cluster has been shut down.
+Applications and Jobs, using archives generated by the JobManager. Unlike the
+[ApplicationResultStore](#applicationresultstore) and
[JobResultStore](#jobresultstore), which store
+minimal metadata for internal recovery decisions in highly-available Clusters,
the History Server
+provides detailed archives for analysis via Web UI or REST API after the
Cluster has been shut down.
#### Instance
The term *instance* is used to describe a specific instance of a specific type
(usually
-[Operator](#operator) or [Function](#function)) during runtime. As Apache
Flink is mostly written in
+[Operator](#operator) or [Function](#function)) at runtime. As Apache Flink is
mostly written in
Java, this corresponds to the definition of *Instance* or *Object* in Java. In
the context of Apache
Flink, the term *parallel instance* is also frequently used to emphasize that
multiple instances of
the same [Operator](#operator) or [Function](#function) type are running in
parallel.
#### Flink Job
-A Flink Job is the runtime representation of a [logical graph](#logical-graph)
-(also often called dataflow graph) that is created and submitted by calling
-`execute()` in a [Flink Application](#flink-application).
+A Flink Job is the unit of data processing execution in Flink: a Job as a
whole is submitted,
+started, stopped and resumed, although under some conditions Flink may restart
a Job only partially
+(See [Restart Pipelined Region Failover Strategy]({{< ref
"docs/ops/state/task_failure_recovery"
>}}#restart-pipelined-region-failover-strategy)).
+
+A Job is submitted either by a [Flink Application](#flink-application), by
calling `execute()` on an
+execution environment, or as a single [Flink SQL
Statement](#flink-sql-statement) or [Statement
+Set](#statement-set).
+
+A Flink Job is the runtime representation of a [Logical Graph](#logical-graph)
(also often called
+*Dataflow Graph*). The Logical Graph is optimized into a [Job
Graph](#job-graph), from which the
+[Physical Graph](#physical-graph) that actually runs in a [Flink
Cluster](#flink-cluster) is derived.
#### Flink Job Cluster
A Flink Job Cluster is a dedicated [Flink Cluster](#flink-cluster) that only
executes a single [Flink Job](#flink-job). The lifetime of the
-[Flink Cluster](#flink-cluster) is bound to the lifetime of the Flink Job.
-This deployment mode has been deprecated since Flink 1.15.
+[Flink Cluster](#flink-cluster) is bound to the lifetime of the Flink Job.
+This deployment mode has been deprecated since Flink 1.15.
+
+#### Job Graph
-#### JobGraph
+Also called *JobGraph* or *Optimized Dataflow*.
-see [Logical Graph](#logical-graph)
+A Job Graph is the optimized representation of a [Logical
Graph](#logical-graph), and the
+representation that a [Flink Application](#flink-application) submits to the
[Flink
+Cluster](#flink-cluster).
+
+Producing the Job Graph is mainly a matter of chaining [Operators](#operator):
consecutive
+[Operators](#operator) that are not separated by a repartitioning are merged
into a single
+[Task](#task). The nodes of a Job Graph are therefore [Tasks](#task), each
implementing one Operator
+or one [Operator Chain](#operator-chain).
+
+The Job Graph is translated into a [Physical Graph](#physical-graph) for
execution.
#### Flink JobManager
-The JobManager is the orchestrator of a [Flink Cluster](#flink-cluster). It
contains three distinct
-components: Flink Resource Manager, Flink Dispatcher and one [Flink
JobMaster](#flink-jobmaster)
-per running [Flink Job](#flink-job).
+Also called *Job Manager*.
+
+The JobManager is the orchestrator of a [Flink Cluster](#flink-cluster). It
does not process any
+data itself: it translates the submitted [Job Graph](#job-graph) into a
[Physical
+Graph](#physical-graph), schedules the resulting [Sub-Tasks](#sub-task) on the
+[TaskManagers](#flink-taskmanager), and coordinates [Checkpoints](#checkpoint)
and
+[Savepoints](#savepoint). It contains three distinct components: Flink
Resource Manager, Flink
+Dispatcher and one [Flink JobMaster](#flink-jobmaster) per running [Flink
Job](#flink-job).
+
+See also [Flink Architecture: JobManager]({{< ref
"docs/concepts/flink-architecture" >}}#jobmanager).
#### Flink JobMaster
JobMasters are one of the components running in the
[JobManager](#flink-jobmanager). A JobMaster is
-responsible for supervising the execution of the [Tasks](#task) of a single
job.
+responsible for supervising the execution of the [Sub-Tasks](#sub-task) of a
single Job. It derives
+the [Physical Graph](#physical-graph) from the Job's [Job Graph](#job-graph),
requests the slots
+needed to run it, deploys the Sub-Tasks to the
[TaskManagers](#flink-taskmanager), and triggers the
+Job's [Checkpoints](#checkpoint).
#### JobResultStore
The JobResultStore is a Flink component that persists the results of globally
terminated
-(i.e. finished, cancelled or failed) jobs to a filesystem, allowing the
results to outlive
-a finished job. Each result contains the job's identifier, final state, name,
the application it
-belongs to, etc. These results are then used by Flink to determine whether
jobs should
-be subject to recovery in highly-available clusters.
+(i.e. finished, cancelled or failed) Jobs to a filesystem, allowing the
results to outlive
+a finished Job. Each result contains the Job's identifier, final state, name,
the Application it
+belongs to, etc. These results are then used by Flink to determine whether
Jobs should
+be subject to recovery in highly-available Clusters.
Review Comment:
I realise Application is a new concept in 2.3, though I think it's valuable
to highlight the differences
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]