Delivery Guarantees and Latency in Confluent Cloud for Apache Flink
Confluent Cloud for Apache Flink® provides exactly-once semantics end-to-end by default, which means that the output of a statement reflects every input message exactly-once and that every output message is delivered exactly once.
To achieve this, Confluent Cloud for Apache Flink relies on Apache Flink®’s checkpointing mechanism and Kafka transactions. While checkpointing and fault tolerance fall into Confluent’s responsibility, you should understand the implications of how Flink reads from and writes to Kafka:
Flink statements write to Kafka by using transactions. Flink commits transactions periodically, approximately every minute.
Flink by default only reads committed messages from Kafka. For more information, see isolation.level.
This implies that depending on the delivery guarantees required by your use case, you can currently achieve different end-to-end latencies with Flink.
Exactly-Once: If you require exactly-once, each statement whose output a consumer reads typically adds up to roughly one minute of latency, depending on the interval at which Flink commits transactions. For pipelines of more than one statement, see Latency in multi-statement pipelines. Ensure that all consumers of the output topics of Flink statements use
isolation.level: read_committedor set the Flink table option'kafka.consumer.isolation-level' = 'read-committed'.At-Least-Once: If at-least-once is sufficient for your use case, you can read from the output topics with
isolation.level: read_uncommitted, which is the default in Kafka, or set the Flink table option'kafka.consumer.isolation-level' = 'read-uncommitted'. With this configuration, depending on the logic of your query, you can achieve an end-to-end latency below 100 ms, but you can see some output messages twice. This happens when Flink needs to abort a transaction that your consumer has already read. If the statement is deterministic, each message is correct, but you can read the same message more than once. Whether that causes incorrect results depends on how the consumer handles duplicates. If the statement doesn’t resume after the failure, for example because you delete it, you can read messages from a transaction that never commits. For more information, see Latency in multi-statement pipelines.
Latency in multi-statement pipelines
In a pipeline, each point where a consumer reads the output of an Flink statement
is a hop. The consumer can be another Flink statement or a client outside Flink
that reads with read_committed. A pipeline of two statements whose output an
application reads has two hops: one where the second statement reads the
output of the first, and one where the application reads the output of the
second.
The commit interval applies at every hop. When a statement writes to a table
that a downstream statement reads with the default read-committed isolation
level, the downstream statement doesn’t see a row until the upstream
transaction commits. Because transactions commit
approximately every minute, each hop in the chain typically adds between zero
and one commit interval of latency, or about 30 seconds on average. A hop can
take longer than one commit interval when checkpoints take longer to complete,
for example under backpressure or with large state.
For example, consider a pipeline that joins two input topics in one statement, writes the result to an intermediate table, and aggregates that table in a second statement. With exactly-once, the two commit waits dominate the end-to-end latency of this pipeline, not the cost of the join or the aggregation. Tuning the join or the aggregation doesn’t reduce this latency. Each statement that you add to the chain adds a commit wait, and each one that you remove saves one. For example, if you run the join and the aggregation in a single statement, or define the join as a view that the aggregation reads, the pipeline has one commit wait instead of two.
If your pipeline has a latency target of a few seconds or less, decide which
delivery guarantee each hop needs before you tune the SQL. If a hop and every
statement and client downstream of it tolerate duplicates, read its input
table with read-uncommitted. With
read-uncommitted, a downstream statement sees rows as soon as the upstream
statement writes them, so the hop adds no commit wait.
A hop tolerates duplicates if it’s stateless, such as a projection or a
filter. It also tolerates duplicates if the table it reads has a primary key
and uses the upsert
changelog mode. In that
case, the replay rewrites each row, and the row ends up with the same value.
A row can briefly show an earlier value while the replay catches up.
Both cases make two assumptions about the upstream statement. First, the statement is deterministic, so it writes the same rows again after a failure. Second, the statement resumes after the failure. If you delete or replace the upstream statement instead, a reader can keep messages from a transaction that never commits.
A hop that aggregates an append-only input, for example with SUM or
COUNT, doesn’t tolerate duplicates, even though its output upserts on the
grouping key. If the upstream statement aborts a transaction after a failure
and writes the rows again, the aggregation counts those rows twice. The
incorrect result remains in the statement’s state.
Duplicates don’t stay in the hop that reads them. A stateless hop that reads
with read-uncommitted writes any duplicates to its output table, where they
become committed messages that downstream statements read even with
read-committed, and that clients outside Flink read even with
read_committed. Lower the isolation level of a hop only if every statement
and every client downstream of it also tolerates duplicates.
Set the isolation level
The kafka.consumer.isolation-level option controls which messages a statement reads from a table. You can’t set it for a whole session with a SET statement. Instead, set it for each table or for each query.
Set it when you create a table. The option applies to every statement that reads the table, so set it this way only if all those statements, and every statement and client downstream of them, tolerate duplicates:
CREATE TABLE orders_enriched ( order_id STRING, amount DOUBLE ) WITH ( 'kafka.consumer.isolation-level' = 'read-uncommitted' );
Change it for an existing table. As with
CREATE TABLE, the option applies to every statement that reads the table:ALTER TABLE orders_enriched SET ( 'kafka.consumer.isolation-level' = 'read-uncommitted' );
Statements that are already running keep the isolation level that they started with. To apply the change to a running statement, replace the statement with a new one. If the statement is stateless, you can use Carry-over Offsets to resume from where the previous statement stopped. Carry-over offsets don’t apply to statements that aggregate, use
LAG, windows, or pattern matching, or write to an upsert sink. A replacement for such a statement starts from the position that scan.startup.mode sets, which isearliest-offsetby default, so it reprocesses the table and writes its results to the output table again.Override it for a single query with a dynamic table option hint, without changing the table definition:
SELECT * FROM orders_enriched /*+ OPTIONS('kafka.consumer.isolation-level' = 'read-uncommitted') */;
The isolation level applies to the table that a statement reads from. To lower
the latency of a hop, set read-uncommitted on the intermediate table that
the hop reads, or in a hint on the query, but only if the hop and every
statement and client downstream of it tolerate duplicates. The option doesn’t
affect consumers outside Flink. For those consumers, set the isolation.level
consumer configuration instead.
Choose a delivery guarantee
Many use cases that seem to require exactly-once work correctly with at-least-once. Consider the following questions when you design a pipeline:
Does the final sink upsert on a key? If the sink writes each result by its primary key, as a database or a document store does, and the statements are deterministic, the sink ends up with the same rows. While a replay catches up, a row can briefly show an earlier value, because the replay rewrites updates that the sink already applied. The final result is the same as with exactly-once only if every statement downstream of a hop that reads with
read-uncommittedalso tolerates duplicates, for example because it doesn’t aggregate an append-only input.Is the downstream logic idempotent? If a consumer already checks whether a message is an insert or an update based on a technical key, it tolerates duplicates under the same conditions as a sink that upserts on a key: the statements are deterministic, and no statement between the consumer and a hop that reads with
read-uncommittedaggregates an append-only input.What is the latency target? Under exactly-once, each hop typically adds up to one commit interval of approximately one minute. If the target is below the sum of the commit intervals across all hops of the pipeline, exactly-once can’t guarantee it.
Do duplicates cause incorrect results? If a consumer appends every message, for example to count events or to trigger a payment, duplicates cause incorrect results. Use exactly-once for this consumer and for every hop upstream of it, because a hop that reads with
read-uncommittedcan commit duplicates to its output.
Note
Confluent is actively working on reducing the latency under exactly-once semantics. If your use case requires a lower latency, reach out to Support or your account manager.