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_committed or 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 is earliest-offset by 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-uncommitted also 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-uncommitted aggregates 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-uncommitted can 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.