TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Same GROUP BY: One Batch Row per Key vs. a 229-Line Streaming Changelog

Continue in TT Lab

Goal

Run the same aggregation in batch and in streaming, and confirm that the final state is the same while streaming emits a changelog. Read, from the execution plan and the rejection error, how the change kind (+I · -U · +U · -D) is determined by the operator and the sink.

Why it matters

A streaming result is not a value that is written once and done but a table that keeps being rewritten. If the receiving side cannot handle retractions (-U), numbers get inflated, and a store that only appends cannot receive it at all. You have to be able to read from the plan in advance which queries emit updates and which sinks can receive them, to prevent accidents at the design stage. The grader of this lab does not ask the cluster anything — it reads the sql-client output and the sink files you saved, and compares the changelog by reproducing it line by line from the source CSV (if the input is append-only and the parallelism is 1, even the log order is determined).

Steps

  1. Start the cluster with flink-up, write SQL in /root/flink/dynamic/batch.sql that, in batch mode, reads /opt/lab/fixtures/data/dynamic_clicks.csv and produces clicks (count) and dwell (sum of dwell) per user_id, and save the output to /root/flink/dynamic/batch.out.
  2. Create /root/flink/dynamic/stream.sql, which runs the same query in streaming mode, and save the output to /root/flink/dynamic/stream.out.
  3. Create /root/flink/dynamic/nested.sql, which runs in streaming an aggregation over an aggregation that produces the number of users (users) per number of clicks (clicks), and save the output to /root/flink/dynamic/nested.out.
  4. In /root/flink/dynamic/explain.sql, write two EXPLAIN CHANGELOG_MODE statements (the per-user COUNT aggregation, and a filter that selects only the user_id, url of rows with dwell > 60) and save the output to /root/flink/dynamic/explain.out.
  5. In /root/flink/dynamic/fs-sink.sql, create a filesystem connector table user_clicks_fs(user_id STRING, clicks BIGINT), write an INSERT INTO of the per-user COUNT, and save the output containing the rejection error to /root/flink/dynamic/fs-sink.out.
  6. In /root/flink/dynamic/to-changelog.sql, convert the per-user COUNT view with TO_CHANGELOG and write it to a filesystem sink (path /root/flink/dynamic/changelog, columns change, user_id, clicks), and save the output to /root/flink/dynamic/to-changelog.out.
  7. In /root/flink/dynamic/upsert.sql, create a print connector table retract_sink and a blackhole connector table upsert_sink with PRIMARY KEY (user_id) NOT ENFORCED, run an EXPLAIN CHANGELOG_MODE INSERT INTO ... for each that feeds the same aggregation to the two sinks, and save the output to /root/flink/dynamic/upsert.out.
  8. In /root/flink/dynamic/report.json, write users, retract_messages, upsert_messages, and nested_deletes.

Notes

Count all at once in batch

Start the cluster with flink-up, and in /root/flink/dynamic/batch.sql write SET 'execution.runtime-mode' = 'batch';, a CREATE TABLE that reads the source CSV, and a SELECT that produces clicks (COUNT(*)) and dwell (SUM(dwell)) per user_id, and save the output to /root/flink/dynamic/batch.out.

For the filesystem connector, the path is file:///opt/lab/fixtures/data/dynamic_clicks.csv and the format is csv. Give the result columns the aliases AS clicks and AS dwell. A batch result has no op column and comes out as one line per user.

Run the same query in streaming

Create /root/flink/dynamic/stream.sql, which runs the same query as in step 1 with SET 'execution.runtime-mode' = 'streaming';, and save the output to /root/flink/dynamic/stream.out. The final state after applying the log to the end must be the same as the batch result.

The result table changes every time a row arrives. A user seen for the first time gets one +I line, and a user that already existed gets two lines, -U which erases the old value and +U with the new value. The number of log lines must be the number of users + 2 × (number of clicks − number of users). Leave the settings at their defaults.

See -D in an aggregation over an aggregation

Create /root/flink/dynamic/nested.sql, which runs SELECT clicks, COUNT(*) AS users FROM (사용자별 COUNT(*) AS clicks) GROUP BY clicks in streaming mode (the parenthesized part stands for the per-user count), and save the output to /root/flink/dynamic/nested.out.

When a user moves from click bucket 1 to bucket 2, the inner aggregation emits -U(1) and +U(2). The outer aggregation receives the -U and reduces the number of users in bucket 1 by one, and when that bucket reaches 0, the result row has to disappear, so it emits -D. Name the column of the inner query clicks.

Read the change kind from the plan

In /root/flink/dynamic/explain.sql, write two statements that start with EXPLAIN CHANGELOG_MODE — the per-user COUNT aggregation, and a filter that selects only the user_id, url of rows with dwell > 60 — and save the output to /root/flink/dynamic/explain.out.

EXPLAIN does not run the job but only prints the plan. If you add CHANGELOG_MODE, each node of the Optimized Physical Plan gets a changelogMode=[...]. I is an insert, UB is before an update, UA is after an update, and D is a delete. What will the node of a query with only a filter emit?

Put updates into a file sink and get rejected

In /root/flink/dynamic/fs-sink.sql, create a filesystem connector table user_clicks_fs(user_id STRING, clicks BIGINT) (path /root/flink/dynamic/fs-out, format csv), write INSERT INTO user_clicks_fs SELECT user_id, COUNT(*) FROM clicks GROUP BY user_id;, and then save the output to /root/flink/dynamic/fs-sink.out. It is normal for it to end with an error.

A file cannot fix a line once it has been written. The optimizer first asks the sink which change kinds it can receive, and if it cannot receive the UB/UA that the aggregation emits, it rejects before submitting the job. The error sentence says which sink cannot receive what from which node.

Pull the change kind out into a column and write it to a file

In /root/flink/dynamic/to-changelog.sql, make the per-user COUNT into a view (two columns, user_id and clicks), and INSERT the result of TO_CHANGELOG(input => TABLE 뷰, op => DESCRIPTOR(change)) (where the placeholder stands for the view) into a filesystem sink (path /root/flink/dynamic/changelog, columns change STRING, user_id STRING, clicks BIGINT, format csv). Save the output to /root/flink/dynamic/to-changelog.out.

TO_CHANGELOG turns each row of the update log into an insert (INSERT) and puts the original change kind into a string column — that is why even a file sink that only appends can receive it. INSERT is asynchronous, so you have to turn on table.dml-sync for the file to be left in its finalized state after the job ends. Empty the sink directory first when you run it again.

In front of a sink with a key, -U disappears

In /root/flink/dynamic/upsert.sql, create a print connector table retract_sink(user_id STRING, clicks BIGINT) and a blackhole connector table upsert_sink with the same columns and PRIMARY KEY (user_id) NOT ENFORCED. Then run two EXPLAIN CHANGELOG_MODE INSERT INTO ... statements that feed the per-user COUNT to each sink, and save the output to /root/flink/dynamic/upsert.out.

A print sink only prints what it receives, so it needs a retraction (UB) to erase an old row. A sink with a declared key can find by key and overwrite, so it does not need UB. Compare the changelogMode of the GroupAggregate in the two plans. NOT ENFORCED means that Flink does not check the key.

Report — the number of messages for retract and upsert

In /root/flink/dynamic/report.json, write as integers users (the number of users = the number of +I in stream.out), retract_messages (the total number of lines of the log in stream.out), upsert_messages (the number of messages that would have been sent if the same result had been sent to an upsert sink), and nested_deletes (the number of -D in nested.out).

An upsert sink overwrites by key, so it does not need to receive the value before the update (-U). Think about which op values remain from the same log. You can get the numbers by counting the op column of the saved output — with grep you can count shapes such as '| +I |' at the start of a line.