Apache Flink — Running Streams on a Real Engine
Removing duplicates and ranking
Goal
Build keep-the-first-row, keep-the-last-row, and Top-N with the same ROW_NUMBER pattern, and confirm the shape and volume of the changelog each emits from the output counts and the execution plans.
Why it matters
Clearing out a resent event and keeping the current state differ by one word in SQL, ASC and DESC, but the former produces an append-only result and the latter produces a retract/update log. This difference decides which sinks you can use and what the downstream aggregation receives. For Top-N, the log volume for the same result differs greatly depending on whether the row number is emitted. The grader of this lab does not ask the cluster anything; it compares the sql-client output you saved with values it computes directly from the source CSV.
Steps
- Start the cluster with
flink-upand define events and sales in /root/flink/dedup/ddl.sql (a watermark on thetsof events). Save the output of /root/flink/dedup/count.sql, which producesn_rows(all rows) andn_orders(the number of distinct orders) in batch, to /root/flink/dedup/count.out. - Save the output of the keep-the-first-row query /root/flink/dedup/first.sql (
order_id,status,amount), which keeps the earliest event for each order, to /root/flink/dedup/first.out. - Save the output of the keep-the-last-row query /root/flink/dedup/last.sql, which keeps the latest event for each order, to /root/flink/dedup/last.out.
- Save the output of /root/flink/dedup/explain.sql, which runs
EXPLAIN CHANGELOG_MODEon the queries of steps 2 and 3, to /root/flink/dedup/explain.out. - Save the output of /root/flink/dedup/board.sql, which counts
orders(the number of orders) andamount(the sum) perstatuson top of keep-the-last-row, to /root/flink/dedup/board.out. - Save the output of /root/flink/dedup/topn.sql (
category,product,total,rn), which gives the top 3 by cumulativeqtytotal per category from sales together with the row number, to /root/flink/dedup/topn.out. - Save the output of /root/flink/dedup/topn-norank.sql, which is the same Top-3 with only the
rnof the outer SELECT removed, to /root/flink/dedup/topn-norank.out. - In /root/flink/dedup/report.json, write
n_rows,n_orders,update_pairs,shipped_now,topn_log_rows,norank_log_rows, andtop_books.
Notes
- Sources (CSVs with no header, times in seconds, row order is arrival order and ts ascending):
dedup_events.csv=order_id, status, amount, ts— 1–4 events per order, and for some the same row comes once more right after (a resend).dedup_sales.csv=sale_id, category, product, qty, ts. All of them are in/opt/lab/fixtures/data/. - Run:
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1(the placeholder stands for the query file name). Streaming results have anopcolumn (+I · -U · +U · -D) at the front. - Deduplication has to follow the shape in the docs exactly:
ROW_NUMBER() OVER (PARTITION BY 키 ORDER BY 시간속성 ASC|DESC) AS rnin the inner SELECT, andWHERE rn = 1outside (where the placeholders stand for the key and the time attribute). - A common mistake: if you sort by processing time (
PROCTIME()), the result can vary with arrival time. This lab sorts by the event timets. - Official docs: Deduplication · Top-N · EXPLAIN · Dynamic Tables
Define the two sources and count the rows
After flink-up, define events and sales in /root/flink/dedup/ddl.sql (WATERMARK FOR ts AS ts on the ts of events). Run /root/flink/dedup/count.sql, which produces n_rows and n_orders in batch mode, with sql-client.sh -i ddl.sql -f count.sql and save it to /root/flink/dedup/count.out.
The number of orders is COUNT(DISTINCT order_id). The difference between the two values is exactly the number of events that arrived second or later for the same order — you will meet this number again in step 3.
Keep the first row — done once emitted
In /root/flink/dedup/first.sql, write, in streaming mode, a deduplication that keeps the one event with the earliest ts for each order (order_id), and produce order_id, status, and amount. Save the output to /root/flink/dedup/first.out.
In the inner SELECT, assign a row number with ROW_NUMBER() OVER (PARTITION BY order_id ORDER BY ts ASC), and on the outside pick only those whose row number is 1. Check whether the op in the output is all +I — a first row that has been emitted once never has a reason to change.
Keep the last row — retract and fix
In /root/flink/dedup/last.sql, write a deduplication that keeps the event with the latest ts for each order (order_id, status, amount), and save the output to /root/flink/dedup/last.out. See how the log counts (+I · -U · +U) relate to the two numbers from step 1.
Just change the sort direction. Every time a new event for the same order arrives, it takes back the row emitted earlier with -U and emits the new row with +U. A pair comes out even when a completely identical row arrives as a resend. If you sort by processing time, this count can wobble, so use ts.
Tell the two deduplications apart in the plan
In /root/flink/dedup/explain.sql, write two statements that put EXPLAIN CHANGELOG_MODE in front of the queries of steps 2 and 3, and save the output to /root/flink/dedup/explain.out.
The execution plan (Optimized Execution Plan) has a Deduplicate(keep=[...]) node, and each line of the physical plan (Optimized Physical Plan) has changelogMode=[...] attached. Compare how keep, outputInsertOnly, and changelogMode differ in the two queries.
Count by current state — an aggregation on top of deduplication
In /root/flink/dedup/board.sql, write a query that puts the keep-the-last-row of step 3 on the inside and, on the outside, produces orders (the number of orders) and amount (the sum of amount) with GROUP BY status, and save the output to /root/flink/dedup/board.out.
When an order changes from created to paid, the deduplication emits -U created and +U paid, and the aggregation subtracts from the created bucket and adds to the paid bucket. The sum of the final orders must equal the number of orders. If you count events directly without deduplication, this sum grows by the number of rows.
Top-3 — fix it every time the ranking changes
In /root/flink/dedup/topn.sql, gather sales per (category, product) as SUM(qty) AS total, and then write a query that gives the top 3 by total descending per category together with the row number rn (category, product, total, rn). Save the output to /root/flink/dedup/topn.out.
Put the GROUP BY total as a subquery, assign ROW_NUMBER() OVER (PARTITION BY category ORDER BY total DESC) on top of it, and pick rn <= 3 on the outside. The ranking moves every time a total changes, so -U and -D get mixed into the output. Looking only at the final state, there are three lines per category.
Drop the row number and the log shrinks
Run /root/flink/dedup/topn-norank.sql (category, product, total), which is the query of step 6 with only the rn of the outer SELECT removed, and save the output to /root/flink/dedup/topn-norank.out. The final top 3 must be the same and the number of log rows must be smaller than in step 6.
The row number column is part of the unique key of the result, so if you emit it, when a product's rank rises, all the rows of the ranks below it are sent again. If you drop the row number, you only need to send the row of the changed product. Count the log rows of the two files with grep -c and compare them.
Report — the volume of the log in numbers
In /root/flink/dedup/report.json, write n_rows and n_orders (step 1), update_pairs (the number of -U in last.out), shipped_now (the number of shipped orders in the final state of board.out), topn_log_rows and norank_log_rows (the number of log rows of the two Top-3 outputs), and top_books (the number 1 product of the books category).
Copy all of them from the output you saved earlier. A log row is a line that starts with the op, like '| +I |'. The final value of a table with a changelog is the last +I or +U of that key. Also check whether update_pairs agrees with the two numbers from step 1.