13 · The Tea Factory: Batch Pipelines

On Tuesday afternoon, at 15:10, Priya posted a screenshot in the team chat. It showed the city scoreboard, the report that the city managers read every morning. Revenue for 1 to 27 September: $4,379,529.

Either September was twice as good as we thought, she wrote, or something is broken. I think it is broken.

Two minutes later, Dana forwarded it to Mia with one line: Are your numbers wrong too?

Mia’s case board, her week-over-week numbers and yesterday’s draft slide all used Steep’s data. She opened her notebook and wrote one question: Which of my numbers touched the city scoreboard’s table?

Theo was already walking to her desk. He had read the chat.

“A backfill,” he said, before she could ask. “It ran twice.”

“A what?”

“Running the job again for past days. Since we built each September day, 1,131 of those orders were refunded. Our nightly job only rebuilds yesterday. So this morning we rebuilt all of September. At 14:30, it ran again.”

“And the second run replaced the first?”

“No. It added to it.” He took a napkin and drew a row of boxes joined by arrows, from left to right. “Our nightly jobs replace a day’s data. The last step of this one was written to add rows. Run it once, and September is right. Run it twice, and September is double.”

Mia looked at the napkin for a while. “Then the real question is not who pressed the button the second time. It is why pressing it twice could change the answer at all.”

For the first time that day, Theo smiled. “That is the question I wish more people asked.”

Mia wrote at the top of a new page: Suspect: the pipeline. Which of my numbers came from that table?

ImportantThe big idea

A good pipeline gives the same result no matter how many times it runs.

A cutaway drawing of a small tea factory that reads from left to right. A hopper pours tea leaves onto a pile on a conveyor belt. A round glass drum washes leaves in water sprayed from a pipe. A sorting machine with a funnel on top sends leaves out onto three separate belts. A tall cabinet dries leaves on four heated trays. A shelf holds fifteen glass jars of finished tea. At the far right, a steaming cup of tea sits on a small table. Teal pipes join the machines, lamps hang from the ceiling, and tea plants stand at both ends.

Look at the picture: a tea factory, drawn from left to right. Leaves pour in from a hopper. A drum washes them. A sorter sends them onto three belts. Racks dry them. Jars on a shelf hold the finished tea, and at the end there is a cup. Each machine does one job and hands its work to the next.

A data pipeline is the same idea: a chain of steps that turns raw data into tables people can use. Steep’s pipeline loads raw data (the hopper), cleans it (the drum), sorts and sums it (the sorter and its belts), and stores the results (the jars), so that a dashboard (the cup) can serve them. The picture hides one thing. A real factory runs on a timetable, and someone must decide when each machine starts, and what happens when one of them jams.

ETL and ELT

People who build pipelines often describe them with three letters. ETL stands for extract, transform, load: take data out of a source, change it on a separate machine, then load the finished result into the warehouse. ELT swaps the last two steps: load the raw data first, then transform it inside the warehouse, usually with SQL.

Steep works the ELT way. The ODS layer of Chapter 12 is the raw load, and DWD, DWS and ADS are SQL that runs inside the warehouse. ELT became common when warehouses became cheap and strong enough to do the heavy work themselves. It keeps the raw copy, so you can transform it again when a rule changes, like Mia’s new definition of orders on Monday. ETL still makes sense when data must be cleaned or masked before it may be stored at all, for example to remove personal details.

The night shift: schedulers and DAGs

Every night at 02:00, a program called a scheduler starts Steep’s pipeline, named steep_daily. The pipeline has seven tasks, one for each table it fills. A task may start only when the tasks it needs have finished: dwd_order_detail cannot run before ods_orders_load has loaded the orders. A rule like “this one first” is a dependency.

Draw each task as a box and each dependency as an arrow, and you get a DAG, a directed acyclic graph. “Directed” means every arrow points one way. “Acyclic” means no path of arrows leads back to where it started. A loop would mean that a task waits for itself, forever.

Show the code
bk.setup()
fig, ax = plt.subplots(figsize=(8, 4.2), layout="constrained")
ax.set_xlim(0, 100)
ax.set_ylim(0, 60)
ax.axis("off")
ax.set_title(f"steep_daily: {WORDS[n_task_names]} tasks, and arrows that say \"wait for this one first\"")
W, H = 22, 9
POS = {"ods_orders_load": (11.5, 38), "ods_events_load": (11.5, 18),
       "dwd_order_detail": (37, 38), "dwd_event_detail": (37, 18),
       "dws_city_platform_day": (63, 47), "dws_user_day": (63, 8),
       "ads_ceo_dashboard_day": (88.5, 28)}
LAYER_COLOR = {"ods": bk.MUSTARD, "dwd": bk.TEAL, "dws": bk.TOMATO, "ads": bk.INK}
for x, label in [(11.5, "ODS"), (37, "DWD"), (63, "DWS"), (88.5, "ADS")]:
    ax.text(x, 57, label, ha="center", va="center", fontsize=11, fontweight="semibold", color=bk.INK)
for task, needs in TASK_PARENTS.items():
    for need in needs:
        (x0, y0), (x1, y1) = POS[need], POS[task]
        ax.annotate("", xy=(x1 - W / 2 - 0.4, y1), xytext=(x0 + W / 2 + 0.4, y0),
                    arrowprops=dict(arrowstyle="-|>", color=bk.INK, linewidth=1.1,
                                    connectionstyle="arc3,rad=0.0"))
for task, (x, y) in POS.items():
    color = LAYER_COLOR[task.split("_")[0]]
    ax.add_patch(FancyBboxPatch((x - W / 2, y - H / 2), W, H, boxstyle="round,pad=0.25,rounding_size=1",
                                linewidth=1.6, edgecolor=color, facecolor=bk.PAPER, zorder=3))
    minutes = median_minutes[task]
    when = "under 1 min" if minutes < 1 else f"about {minutes:.0f} min"
    ax.text(x, y + 1.3, task, ha="center", va="center", fontsize=8.6, family="monospace",
            color=bk.INK, zorder=4)
    ax.text(x, y - 2.3, when, ha="center", va="center", fontsize=8.5, color=bk.MUTED, zorder=4)
plt.show()
A diagram of seven boxes in four columns, labelled ODS, DWD, DWS and ADS. Arrows run from ods_orders_load to dwd_order_detail, and from ods_events_load to dwd_event_detail. From dwd_order_detail, arrows go to dws_city_platform_day, dws_user_day and ads_ceo_dashboard_day. From dwd_event_detail, arrows go to dws_user_day and ads_ceo_dashboard_day. Each box shows how many minutes the task usually takes, from under a minute to about four minutes.
Figure 1: Steep’s nightly pipeline, steep_daily. Each box is one task; an arrow means “wait for this one first”. Times are the median from the run log.

The scheduler’s run log records every task it starts: when, for which day, and how it ended. On Tuesday afternoon, Steep’s log held every night from 1 June to 28 September: 120 nightly runs of seven tasks, and 850 task runs in all. (That is 10 more than 120 × 7. You will see why in a moment.) Steep’s scheduler runs one task at a time, in an order that respects the arrows.

Three tools do most of this work in the real world. Apache Airflow is the best-known open-source scheduler; each DAG is a Python file. Apache DolphinScheduler is popular in China, and many teams there draw their DAGs in a web page. On Alibaba Cloud, DataWorks schedules jobs for MaxCompute, Alibaba Cloud’s data warehouse service, and for other engines. The words differ (a “workflow”, a “node”) but the ideas are the same.

One idea confuses almost everyone. The run that starts at 02:00 on 25 September works on the data of 24 September. Airflow calls that day the run’s logical date. DataWorks calls it the data timestamp, and Chinese teams often say 业务日期 (business date). The run waits until the day is over, so that all of the day’s data has arrived. Some tools date the run by its start time instead. Check which day your run really loads.

Each run then writes exactly one partition (Chapter 9), named by that date: dt=2026-09-24. This one habit makes everything else in this chapter easier. A run, a date and a partition line up, one to one.

When a task fails: retries

Machines fail for boring reasons: a short network break, a busy database, a full disk. Most of these failures go away on their own. So a scheduler can retry a failed task: wait a little, then run it again. In Airflow the settings are called retries and retry_delay. DolphinScheduler has the same two settings on every task: how many times to retry, and how long to wait in between.

In Steep’s run log, 10 task runs failed on their first try. All 10 succeeded on the second try, each one 10 minutes after the failure. That explains the extra 10 task runs. Two of the failures fell in the weeks of the case:

Table 1: The two failed tasks in the case weeks, from Steep’s run log. Times are Steep’s local time.
Data date Task First try ended (failed) Retry ended (success) Rows in the partition
5 September ods_events_load 6 September, 02:01 6 September, 02:12 68,072
7 September ods_orders_load 8 September, 02:00 8 September, 02:10 5,486

In both cases, the retry wrote exactly as many rows as the partition holds: one copy, not two. These failed tries wrote no rows. So even an appending task would have left one copy, by luck, not by design. Every task in steep_daily overwrites its partition, and overwrite makes a retry safe even when a try fails halfway. A retry is only safe if running the task again cannot hurt. That property has a name.

Idempotency: the same result every time

Press the call button for an elevator five times, and one elevator comes. The first press calls it; the other four change nothing. An action like that is idempotent (you met the word in Chapter 8): doing it twice has the same effect as doing it once.

A task that writes one day of data can do it in two ways.

  • Append: add the new rows to whatever is already there. Run it twice, and you get two copies.
  • Overwrite the partition: replace the whole day with the new rows. Run it twice, and you still get one copy.

In Hive’s SQL, the difference is one word. INSERT INTO appends to the table or partition. INSERT OVERWRITE replaces what is there:

-- Idempotent: run it once or ten times, 24 September holds one copy.
INSERT OVERWRITE TABLE dws_city_platform_day PARTITION (order_date = '2026-09-24')
SELECT city, platform, count(*) AS orders, ...
FROM dwd_order_detail
WHERE order_date = '2026-09-24' AND NOT is_refunded
GROUP BY city, platform;

-- Not idempotent: every run adds one more copy of 24 September.
INSERT INTO TABLE dws_city_platform_day PARTITION (order_date = '2026-09-24')
SELECT ...;

Airflow’s own guide to good practice says the same thing in other words. Treat a task like a transaction in a database. A task should produce the same outcome on every re-run. Do not use a plain INSERT in a task that may run again, because it can create duplicate rows. And read and write a fixed partition, never “the latest data”, because the latest data changes between runs.

Overwriting a partition is the most common way to be idempotent, but not the only one:

  • Delete, then insert, in one transaction. Remove the day’s rows and write them again, so that nobody ever sees the gap between the two steps.
  • Merge on a key (an upsert: update the row if the key exists, insert it if not). This works only if the key is truly unique, such as (order_date, city, platform).

Idempotency also needs fixed inputs. A task that reads “everything up to now” gives a different answer tomorrow, even if its code never changes.

Backfills: running the past again

A backfill runs a pipeline again for dates in the past. Chinese teams call it 补数 (“filling in the numbers”). There are good reasons to do it: a bug was fixed, a new column needs history, a new table needs its first months, or the source changed after the fact, like Steep’s late refunds (Chapter 11).

Every scheduler supports it. In Airflow 3 it is a command, airflow backfill create --dag-id ... --from-date ... --to-date .... DolphinScheduler calls it complement, and lets you run the dates one after another (serial) or several at once (parallel). DataWorks calls it data backfill, and can run the chosen task together with the tasks that depend on it.

A scheduler can also fill some gaps by itself. Suppose it was switched off for three nights. When it starts again, Airflow can create the three missed runs, one for each missed data date. Airflow calls this catchup, and it is off unless you turn it on. Either way, a late run must still work on its own data date, not on “today”. That only works if every task reads and writes the partition for the date it is given.

A safe backfill has six checks.

  1. Start at the first table whose input changed. If the orders changed, start at ODS. Rebuilding DWS alone would copy the old DWD again.
  2. Every task in it overwrites its partition. Check the write mode, not the name of the job.
  3. The tasks downstream run too. If you fix DWS, the ADS tables built from it are still old.
  4. Serial or parallel is a real choice. Steep’s daily partitions do not depend on each other, so they can run side by side. A zipper table (Chapter 12) builds each night from the night before, so its nights must run in order.
  5. Control totals before and after (see below).
  6. Tell the people who use the table. A table that is being rebuilt can be half-empty for a while.

What happened on 29 September

The run log tells the story without anyone’s memory. The backfill job ran twice, at 10:05 and at 14:30. Each run reloaded the orders into ODS and DWD, then appended them to DWS. With the delete before them and the repair the next morning, the log of the backfill pipeline, steep_backfill_dws, shows four runs:

  • At 09:50 on 29 September, it deleted all 27 daily partitions of dws_city_platform_day for 1 to 27 September.
  • From 10:05 to 10:23 on 29 September, it reloaded the orders into ODS and DWD for 1 to 27 September (ods_orders_load, then dwd_order_detail), overwriting each day, then it rebuilt dws_city_platform_day in append mode. September was complete, and up to date.
  • From 14:30 to 14:49 on 29 September, it reloaded the orders into ODS and DWD for 1 to 27 September (ods_orders_load, then dwd_order_detail), overwriting each day, then it rebuilt dws_city_platform_day in append mode. Every September partition of the summary table now held two copies.
  • From 09:20 to 09:36 on 30 September, it rebuilt dws_city_platform_day in overwrite mode. One copy again.

Every ODS and DWD reload overwrote its day, so running those steps twice did no harm. Only the last step appended: 27 drops, 108 reloads, 54 appends and 27 repairs in all.

Show the code
fig, ax = bk.figure(8, 4.0)
ax.step(timeline.t, timeline.orders / 1000, where="post", color=bk.INK, linewidth=1.8)
ax.axhline(control_orders / 1000, color=bk.TEAL, linestyle="--", linewidth=1.2)
ax.fill_between(timeline.t, timeline.orders / 1000, control_orders / 1000, step="post",
                where=timeline.orders > 1.5 * control_orders, color=bk.TOMATO, alpha=0.25, linewidth=0)
ax.text(timeline.t.iloc[-1], control_orders / 1000 * 0.92, "the right total", color=bk.TEAL,
        ha="right", va="top", fontsize=9.5)
labels = [(drop_start, "drop, reload,\nappend 1"), (append2.start, "append 2"),
          (repair.start, "overwrite")]
for t, text in labels:
    ax.annotate(text, xy=(t, doubled_orders / 1000 * 1.04), ha="left", fontsize=9, color=bk.INK)
ax.set_ylim(0, doubled_orders / 1000 * 1.15)
ax.set_ylabel("Orders, thousands")
ax.xaxis.set_major_locator(mdates.HourLocator(byhour=[0, 6, 12, 18]))
ax.xaxis.set_major_formatter(mdates.DateFormatter("%d %b\n%H:%M"))
ax.set_title(f"September doubled at {clock(append2.end)} and stayed doubled for "
             f"{hours_doubled:.0f} hours")
plt.show()
A step chart from the morning of 29 September to the next morning. The line starts near the correct total, falls to zero for a few minutes in the morning, climbs back to the correct total, then in the afternoon climbs to twice the correct total. It stays doubled overnight and falls back to the correct total the next morning. A dashed teal line marks the correct total.
Figure 2: Orders for 1–27 September in dws_city_platform_day, replayed from the run log and the binlog. Each step is one partition dropped, appended or overwritten. The dashed line is dwd_order_detail after the morning reload.

For about 19 hours, everyone who read dws_city_platform_day saw September twice. When Priya posted, it said 347,279 orders and $4,379,529. dwd_order_detail, reloaded that morning, said 173,647 orders and $2,189,853. Not exactly twice: the second run reloaded DWD about 4 hours later, and a few more refunds had landed in between. Each day should hold 12 rows, one for each city and platform; each September day held 24.

“Who ran it the second time?” Mia asked.

“An engineer who did not see that the morning run had finished,” said Theo. “People do that. We do not hunt for the person. We fix the job: from today, a backfill job must overwrite, or it does not run.” His team writes a short report after every incident, without blame, so that people tell the whole truth next time.

“Why not repair it now?”

“Because a repair written in a hurry is how a second incident starts. Tonight’s run does not touch September. Tomorrow at 09:20, after a second engineer has read the repair job, we overwrite September once.” Until then, a note at the top of the city scoreboard said that September was being repaired.

Did it touch Mia’s numbers?

Mia answered her own question with the lineage map she had drawn on Monday. Her case numbers came from the orders database and from dwd_order_detail. The dashboard’s line came from dwd_event_detail. None of them reads dws_city_platform_day. Her case board was safe. So was Dana’s wall: its new database line, agreed on Monday, was not live yet, and the old event line reads dwd_event_detail.

She checked what would have happened if it had not been. Computed from the doubled table, the week of the drop would have shown a change of +1.1% instead of −4.9%, the database’s number that afternoon (the case board uses final statuses, see Chapter 1). A small rise instead of a fall. The backfill covered September only. So 31 August had one copy, and the other 13 days had two. A broken table does not always look broken. Sometimes it looks like good news.

There was one more thing, and it worried her. Her new metric views from Monday, orders and net_revenue, read dws_city_platform_day. Her slide numbers had been computed on Monday, before the backfill, so the backfill did not touch them. But anyone who refreshed them that afternoon would have seen double. One place to define a metric is also one place to break it. That is not an argument against a metric layer. It is an argument for checking the layer below it every night.

Control totals: counting between steps

How should the pipeline itself have noticed? With the oldest control in auditing: count things on both sides of every step.

  • A record count compares the number of rows. Each day of dws_city_platform_day should hold 12 rows. On the evening of 29 September, September’s days held 24 each.
  • A control total compares a sum. The orders in dws_city_platform_day for a day must equal the completed, not refunded orders in dwd_order_detail for that day. For 1 to 27 September, dwd_order_detail said 173,647 that afternoon. The doubled table said 347,279.
  • A grain check tests the table’s promise, one row per key (Chapter 12). For 1 to 27 September, the doubled table had 648 rows but only 324 different (order_date, city, platform) keys.

The backfill wrote one day per task. Any one of these checks, run after each task, would have failed at about 14:32, after the first doubled day, and stopped the job. Dashboards that read the table could still have shown that one bad day. Safer still: write to a staging table, a private copy that no report reads, check it there, then swap it in. Instead, the whole table was doubled by 14:49, and Priya noticed it 21 minutes later, by eye. A good pipeline does not depend on someone’s eye.

Promises about time: SLAs and freshness

Correct data that arrives too late is still a problem. Theo’s team has a written promise, pinned at the top of the data team’s page: yesterday’s data is in the warehouse by 07:00, Steep time, every morning. A promise like this is a service level agreement, or SLA. Freshness is how old the newest data in a table is. It is how you check the promise: if at 07:00 the newest day in dws_city_platform_day is not yesterday, the SLA is broken.

Steep keeps its promise with a wide margin. Over 120 nights, the pipeline finished between 02:09 and 02:31. Even the nights with a retry finished long before 07:00.

Schedulers can watch the promise for you: you give a task a time by which it must finish, and the scheduler warns you when it will be late. DataWorks calls this baseline monitoring; the box below has the details for the other tools.

The data date. For a daily pipeline, the run for data date \(D\) covers the period of data from \(D\) 00:00 to \(D+1\) 00:00, its data interval. It starts after the interval ends, at \(D+1\), 02:00 for Steep. So the data date is the day before the day the run starts. Airflow’s documentation says the same: a run is usually scheduled after its data interval has ended. One trap: many teams set the timetable with a cron schedule, a short code such as "0 2 * * *" that means “every day at 02:00”. In Airflow 3, a cron schedule has no data interval by default, and the logical date is the time the run starts: 25 September, 02:00, not 24 September. Turn data intervals on (the setting is create_cron_data_intervals), or compute the day before yourself. This is a classic off-by-one bug.

Deadlines in each tool. Airflow 2 had an SLA feature; Airflow 3.0 removed it, and Airflow 3.1 replaced it with deadline alerts. DataWorks baseline monitoring predicts a miss from how long the tasks usually take, and warns you early. DolphinScheduler can raise a timeout alarm when one task runs too long.

Why the doubled table showed a rise. Let \(W_1\) and \(W_2\) be the true weekly totals, and let \(a\) be the orders on 31 August, the only day of the two weeks that was not doubled. If both runs read the same DWD, the doubled table gives

\[\frac{2W_2}{a + 2(W_1 - a)} = \frac{2W_2}{2W_1 - a}.\]

The denominator is a little less than \(2W_1\), so the ratio is a little larger than \(W_2/W_1\). Here it is large enough to turn a fall of about 5% into a small rise.

An idempotent merge. On engines with MERGE, an upsert on the grain key looks like this:

MERGE INTO dws_city_platform_day AS t
USING todays_rows AS s
  ON t.order_date = s.order_date AND t.city = s.city AND t.platform = s.platform
WHEN MATCHED THEN UPDATE SET orders = s.orders, net_revenue = s.net_revenue
WHEN NOT MATCHED THEN INSERT (order_date, city, platform, orders, net_revenue)
                      VALUES (s.order_date, s.city, s.platform, s.orders, s.net_revenue);

Run it twice and the second run updates every row to the same values. It does not, however, remove a row that should no longer exist (say, a city with no orders today). Overwriting the whole partition does.

Try it

The first playground is a small copy of Steep’s nightly pipeline. Its partitions hold the real number of orders for each September day, and its control total comes from dwd_order_detail.

The pipeline step-through. The bars are the partitions of dws_city_platform_day for 1–27 September: one bar per day, teal for the first copy, tomato for any extra copy. Days 1 to 20 are loaded. Run the nights that follow, break a task, retry it, and backfill. Every task overwrites its day, except that you can switch the DWS task to append.

Things to try:

  • Run two nights. Each task turns teal, and each new day holds one copy.
  • Make dwd_order_detail fail with 0 retries. The tasks that need it are marked “upstream failed”, but dwd_event_detail still runs: the scheduler follows the arrows. Tonight’s day is missing.
  • Set the backfill to start at dws_city_platform_day and press “Backfill all loaded days”. The gap stays: DWS cannot run without its input. Now start at dwd_order_detail. The gap fills, and nothing else moves. Press it again. Still nothing moves: that is idempotency.
  • Replay 29 September: switch to append, drop the loaded days, then backfill twice. The first backfill looks perfect. The second doubles everything, and the control total turns red. Repair it with one overwrite backfill.
  • In append mode, make dws_city_platform_day fail once with 1 retry. In this toy, as a loader without transactions can, the failed try leaves half its rows behind, and the retry adds a full copy on top. A retry is only as safe as the write it repeats.

This toy is simpler than a real scheduler in one way: it runs one task at a time, like Steep’s. Airflow, DolphinScheduler and DataWorks can run independent tasks side by side.

The second playground is real: Steep’s run log, the doubled table and the correct table, in a small database inside your browser. Two things differ from Tuesday afternoon. The log runs to the end of the data, 25 October, so it also shows the repair and the nights after it. And the copy of the doubled table was rebuilt later from the final data, which already holds the refunds that arrived after 29 September. So it has 345,970 September orders, not the 347,279 that Priya saw. Every September day is still there twice.

The run log, in your browser. Pick a starting query, change it if you like, and press “Run query”. Times in the log are UTC; the queries convert them to Steep’s local time.

The first query shows every failure and its retry: 11 in the whole log, 10 of them before Priya’s message. Try the others too, then change a date, a table name or a threshold, and run your own.

Common traps

  • Append in a job that can run twice. Retries, backfills and a second click all run a job again. If it appends, each run adds a copy.
  • Reading “the latest” data. A task that reads whatever is newest gives a different answer every time it reruns. Read the partition for the run’s data date.
  • Confusing the run date with the data date. The run on 25 September loads 24 September. Off-by-one bugs here are common and quiet.
  • Backfilling one table, not its children. Fixing DWS leaves ADS built from the old DWS.
  • A parallel backfill of a table whose days depend on each other. A zipper table needs its nights in order.
  • No checks between steps. If the first alarm is a person reading a dashboard, the pipeline has no controls. Count rows, add up totals, and test the grain after every write.
TipAudit Instinct · Batch controls and duplicate postings

Batch controls. Before computers did it for them, clerks sent invoices to data entry in batches, with a slip on top: the number of documents (a record count), the total amount (a control total), and the sum of a field that means nothing on its own, such as all the invoice numbers added together (a hash total). After entry, the system’s totals had to match the slip. Auditors still test these input controls. A pipeline needs exactly the same slip at every step: rows in, rows out, totals in, totals out.

Duplicate postings. Posting the same batch twice is a classic accounting error, and a classic way to hide fraud. Good systems make it harmless: each batch and each document has a unique number, and the ledger refuses a number it has already seen. That refusal is idempotency. Doubled September is a duplicate posting, and overwriting a partition is the control that makes posting twice harmless.

NoteInterview Corner

Q1. How do you make a pipeline idempotent? (如何保证任务幂等?)

Make every task a function of its data date: it reads fixed input partitions for that date and fully replaces its output partition. In Hive, use INSERT OVERWRITE ... PARTITION (dt = '...'), never INSERT INTO. Elsewhere, delete and insert the partition in one transaction, or MERGE on a unique key. Never read “the latest” data, and never let a task depend on what an earlier run left behind. Then retries and backfills are safe. Prove it with checks after each write: row counts, control totals against the layer below, and a uniqueness test on the grain key.

Q2. How do you run a safe backfill? (补数怎么做?)

Start at the first table whose input changed, and rerun every task downstream of it, so that DWS and ADS are rebuilt from the new DWD. Check that every task in the backfill overwrites its partition. Choose the date range. Decide serial or parallel: partitions that do not depend on each other can run in parallel; anything that builds on the day before, such as a zipper table or a running total, must run in order. Record control totals before and after, and compare them with the layer below. Tell the people who read the tables, and lower the backfill’s priority so it does not delay the nightly run. Airflow (airflow backfill create in Airflow 3, airflow dags backfill in Airflow 2), DolphinScheduler (“complement”, serial or parallel) and DataWorks (“data backfill”) all support this.

Q3. What is the difference between ETL and ELT?

ETL transforms data on its way into the warehouse, on a separate engine, and loads only the result. ELT loads raw data first and transforms it inside the warehouse with SQL. ELT keeps a raw copy, so you can re-transform it when rules change, and it uses the warehouse’s own power; it is the common choice today, and layered warehouses (ODS first, then DWD and above) follow it. ETL is still useful when data must be filtered or masked before it may be stored, or when the target cannot do heavy work.

Reported change: −12.0% orders, the week of 7 September compared with the week before (CEO dashboard).

Explained so far: 0 of the 12 points. About 7 of them sit between the orders database and the dashboard (Clue 1).

Suspects: the app, and the service that writes to Kafka (Chapter 8). The iOS app, version 3.2.0, is the main suspect (Chapter 6). Not proved.

Ruled out: the matcha menu (Chapter 4), Kafka (Chapter 8), storage (Chapter 9), computing (Chapter 10), refunds and restatements (Chapter 11), the warehouse layers (Chapter 12), and now the pipeline. The 29 September incident touched one table, dws_city_platform_day, from 14:30 that day, two weeks after the drop was first seen. The case numbers come from the orders database, dwd_order_detail and dwd_event_detail.

Open questions: are the missing order_completed events only late, still on their way? (Chapter 14.) Did iOS 3.2.1, released on 24 September, fix it? How many of the 12 points does it explain? (Chapter 15.)

New evidence: in the case weeks, two nightly tasks failed and were retried, and each retry left exactly one copy of its partition.

Recap

  • A pipeline is a DAG of tasks that a scheduler runs, usually one partition per data date. Retries and backfills run tasks again, so running again must be safe.
  • Idempotent tasks overwrite their partition (or delete-and-insert, or merge on a unique key). Appending turned one extra click into a doubled September.
  • Count between steps: record counts, control totals and grain checks catch a broken load before a person does. An SLA and a freshness check make sure correct data is also on time.
English 中文
data pipeline 数据管道
ETL / ELT 抽取-转换-加载 / 抽取-加载-转换
scheduler 调度器
task 任务
DAG (directed acyclic graph) 有向无环图
dependency 依赖
run log 运行日志
logical date / data date 逻辑日期 / 业务日期
partition 分区
retry 重试
idempotent 幂等
append / overwrite 追加写 / 覆盖写
upsert 更新插入
backfill 补数 / 回填
serial / parallel complement 串行补数 / 并行补数
SLA (service level agreement) 服务等级协议
freshness 新鲜度 / 时效性
record count 记录数
control total 控制总数
hash total 哈希总数
batch control 批次控制

Further reading