The Quietest Bug in Your Data Pipeline
Most engineering bugs announce themselves. A service crashes, a test goes red, an exception lands in your inbox. You know something is wrong because something stopped working.
Data pipelines have a more dangerous failure mode: the job stays green, runs on schedule, the dashboard lights up with fresh numbers - and the numbers are wrong. Nobody finds out until someone reconciles against the source weeks later and notices the totals don’t add up. By then the wrong number has been in three reports and one board deck.
I build pipelines on the assumption that this is the failure that matters. Everything below is a set of mental models for designing batch pipelines that fail loudly when they fail at all, and stay correct when they don’t. It’s deliberately about how to think, not which tool to click.
One question decides your write strategy
When you load data into a table, you reach for one of two strategies: overwrite (wipe and rewrite) or merge/upsert (match rows by key and update them in place). Most “which one should I use” guides hand you a list of rules. You don’t need a list. You need one question:
For the unit I’m about to overwrite, can I fully reconstruct its correct contents from the data I have on hand?
If yes, overwrite. If no, merge.
That’s the whole decision. Everything else is a special case of it.
Consider immutable events partitioned by time - clicks, logs, transactions. A click that happened at 9am on the 29th has an event_date of the 29th forever; it never moves and never changes. So when you need to reprocess the 29th, you pull all of the 29th’s events from the source, rebuild the day completely, and overwrite. You can reconstruct the whole unit, so overwrite is safe - and you get idempotency for free, because rerunning the 29th just wipes and rewrites the 29th with the same result.
Now take an orders table where an order changes status over time. An order created on June 1st lives in the June 1st partition forever, but its status changes today. To update it you’d have to overwrite the June 1st partition - but today’s change set contains only that one order, not every other order from June 1st. Overwriting the partition would delete thousands of orders you can’t rewrite. You cannot reconstruct the whole unit, so you must merge: match the single row by order_id and update it.
Notice the rule was never “mutable entities can’t be overwritten.” Feed that same orders table a full nightly snapshot - every order at its latest status - and overwrite comes right back, because the snapshot is the complete contents of the table. The entity didn’t change; your ability to reconstruct the unit did.
Partitions are the unit of reprocessing
The word “unit” is doing a lot of work above, and partitions are what make it small.
A partition is a physical split of your table by some column. “Partition by event_date” means every row for the 29th lives together, usually in its own folder of files, isolated from every other day. That isolation is the point: you can overwrite the 29th without touching the 28th.
This is what makes “reconstruct the whole unit” affordable. Without partitions, reprocessing the 29th means rebuilding and overwriting the entire table - expensive and risky. With partitions, the unit you have to reconstruct shrinks from the whole history down to a single day. Partitioning is what turns reprocessing from a once-a-quarter ordeal into something you can do every morning.
So the real question, stated precisely, is: the unit I plan to overwrite - one partition, or the whole table - can I rebuild it completely from what I have? Partitions exist to make that unit as small as the data’s natural grain allows.
Two guarantees that look like one
Here’s the trap that catches people. Suppose you overwrite the 29th from a message queue, and you rerun the pipeline five times. Are you safe from duplicate data?
You’re safe from one kind of duplication and not another, and conflating them is how correctness bugs slip through.
Overwrite gives you idempotency of the write: rerunning the job any number of times produces the same partition, because each run wipes and rewrites from scratch and never sees the previous run’s output. This is real and valuable - reruns never make things worse.
But idempotency of the write says nothing about whether the contents are correct. If your source data already contains duplicates, overwrite will faithfully write those duplicates every single time. Perfectly idempotent, and consistently wrong.
These two properties are orthogonal:
- Idempotency of the write (overwrite handles this): rerunning doesn’t multiply your data.
- Correctness of the content (something else has to handle this): no real-world event is counted twice inside the partition.
The number of times you run the pipeline, and the duplicates baked into your source, are different problems. Overwrite solves the first. The second needs its own tool.
At-least-once delivery, and the natural key
Where do source duplicates come from? Usually from how your messaging layer delivers data.
Most queues guarantee at-least-once delivery. The mechanism is worth understanding precisely, because it explains why the fix lives where it does. The queue sends event X. You process it and send back an acknowledgement. But the ack gets lost - a network blip, or your consumer dies right after processing but before reporting success. The queue waits, hears nothing, concludes “they probably didn’t get it,” and sends X again. Not because it’s malicious, but because it would rather deliver twice than lose a message. That is what at-least-once means: “I guarantee you receive it at least once, not exactly once.”
So the duplicate isn’t two copies in one delivery. It’s the same event delivered twice, possibly seconds or hours apart, scattered into your data stream. You can’t catch it by inspecting a single batch. You catch it with a natural key.
A natural key identifies an event by its business meaning, not by its row. For a click it might be (user_id, page, event_timestamp); for an order, order_id; for a payment, transaction_id. The idea: two rows with the same natural key are the same real-world event, no matter how many times they arrived. Deduplication means grouping by natural key and keeping one row per key before you write.
This is why dedup lives in your data layer, not in the queue. The queue can’t know whether the X it’s resending is the same X from before - resending is the correct behavior from its point of view. Only you, at the destination, can define “these two rows are the same real event.” Accept the duplicates, then kill them by key. That’s how you get exactly-once effects without an exactly-once queue - which is the pragmatic way real systems achieve it.
Choosing the key is a business decision, not a technical one
This is the hardest and most underrated part. “Unique in what sense?” is not a question your engine can answer. The engine only compares bytes. You, with whoever understands the domain, define what counts as a single distinct event - and that definition shifts with the business.
Get it wrong in either direction and you get a silent failure, in opposite shapes.
Key too narrow - say you dedup clicks by user_id alone. A real user clicks Home in the morning, Products at noon, Checkout at night: three distinct, real events, all sharing one user_id. Group by user_id, keep one row, and you’ve collapsed three real events into one. You didn’t remove duplicates; you removed data. This is over-deduplication, and it under-counts.
Key too wide - say you fold processing_timestamp (when your pipeline handled the row) into the key. The original and the resent copy were processed at different times, so they get different keys, so dedup doesn’t recognize them as the same event. The duplicate sails through. This over-counts.
The same payments table can need different keys depending on the rules. If every transaction has a gateway-issued transaction_id, that’s your key. But if the business allows multiple payments against one order - installments, or a top-up after an underpayment - then order_id isn’t unique and (order_id, payment_attempt) is. Same data, different definition of “unique,” because the business logic differs.
So the process for choosing a key is really one question to the domain: what makes two records the same event, and what makes them different events? The answer is the natural key. The rule that falls out:
A natural key must contain enough to distinguish two genuinely different events (too little -> over-dedup -> under-count), and nothing that isn’t intrinsic to the event itself - never anything the pipeline generates, like processing time or row id (too much -> duplicates leak -> over-count).
And because a bad key never crashes anything, this is a silent bug by construction. No exception, no alert. Just numbers that are quietly off.
Incremental models and the lookback window
The last piece ties it all together in the place these tools actually get used.
The opposite of an incremental model is a full refresh: recompute the entire table from scratch every run. Fine for small data; absurd when each morning means recomputing years of history. An incremental model processes only what’s new or changed since the last run, then folds it into the existing table. Efficient - and fragile, in a way that pulls in everything above.
The naive version tracks a watermark: “only take rows with a timestamp greater than the highest I saw last time.” The problem is late data. An event happens yesterday but arrives today, after the watermark has already moved past yesterday. It never gets processed. It vanishes, silently.
The fix is a lookback window: instead of taking only what’s strictly newer than the watermark, each run reprocesses a trailing window - say “always recompute the last 3 days.” A late event landing inside that window still gets caught. You trade a little extra work for not losing late data.
But the window creates exactly the duplication problem from earlier: if each run reprocesses the last 3 days and appends, those overlapping days stack on top of each other and multiply. So a lookback window forces idempotent writes - either partition-overwrite the days in the window, or merge/upsert with dedup by natural key. The whole Week 1 toolkit, in one pattern.
Lookback is the price of not watching the data arrive
Here’s the mental model I keep coming back to. A lookback window is the cost of triggering on a schedule instead of on data.
A scheduler (cron, a daily Airflow run) fires on time. It wakes at 2am and declares “time to process yesterday” - whether or not yesterday’s data has actually arrived. It can’t know; it only knows the clock. The lookback window is the answer to that blindness: “since I can’t tell whether the last few days are complete, I’ll rescan them every run, just in case.” It’s a defensive buffer against latency the scheduler can’t observe.
An event-driven trigger is the opposite. It fires because data arrived - a file lands in S3, a message hits a queue. The late event, whenever it shows up, generates its own trigger and gets processed right then. No need to rescan three days on the off chance something was late, because lateness announces itself.
The analogy: a scheduler is someone collecting mail at a fixed hour, who - not knowing when letters arrive - has to re-check the last few days’ box to be sure nothing came late. Event-driven is a doorbell: a letter arrives, it rings, you grab that one. No re-checking.
This reframes how you choose the window length N. With a scheduler, N is a bet on the maximum lateness of your source. Too short, and any event later than the window falls outside every scan and is lost forever - a correctness failure. Too long, and you grind over partitions that settled long ago - a cost failure. Choose N from the actual lateness distribution of the source (the high percentile, plus a safety margin), not from a round number.
And here’s the failure in full: a source that’s occasionally 2 days late, behind a daily scheduler with a 1-day lookback. The day the late record finally arrives, the pipeline rescans only yesterday and today - but the record belongs to two days ago, outside the window, past the watermark. No run ever looks at it again. It’s gone. The job is green, the dashboard is up, the aggregate is quietly short. A watermark and a lookback set wrong turn a data bug into an invisible one.
The throughline
Every section here is the same lesson wearing different clothes. The decision table, the partition grain, the orthogonal guarantees, the natural key, the lookback length - in each case the dangerous mistake doesn’t throw an error. It succeeds, on schedule, with wrong numbers.
So design for it. Ask whether you can reconstruct the unit before you overwrite it. Assume your source has duplicates and define the key that kills them. Assume data arrives late and size your window for it. Build pipelines that assume failure - because the worst failures are the ones that look exactly like success.
Hầu hết các bug kỹ thuật đều tự thông báo. Một service crash, một test đỏ, một exception rơi vào hộp thư của bạn. Bạn biết có gì đó sai vì có thứ gì đó đã ngừng hoạt động.
Data pipeline có một kiểu thất bại nguy hiểm hơn: job vẫn xanh, chạy đúng lịch, dashboard sáng lên với những con số mới - và những con số đó sai. Không ai phát hiện ra cho đến khi ai đó đối chiếu với nguồn vài tuần sau và nhận ra tổng số không khớp. Lúc đó con số sai đã xuất hiện trong ba báo cáo và một deck trình cho hội đồng quản trị.
Tôi xây dựng pipeline với giả định rằng đây là loại thất bại quan trọng nhất. Tất cả nội dung dưới đây là một tập hợp các mental model để thiết kế batch pipeline biết thất bại ồn ào khi thất bại, và giữ đúng khi không thất bại. Bài viết cố ý nói về cách tư duy, không phải về việc click vào công cụ nào.
Một câu hỏi quyết định chiến lược ghi dữ liệu
Khi bạn load dữ liệu vào một bảng, bạn sẽ chọn một trong hai chiến lược: overwrite (xóa và ghi lại) hoặc merge/upsert (khớp hàng theo key và cập nhật tại chỗ). Hầu hết các hướng dẫn “nên dùng cái nào” đưa cho bạn một danh sách quy tắc. Bạn không cần danh sách. Bạn cần một câu hỏi:
Với đơn vị tôi sắp overwrite, tôi có thể tái tạo hoàn toàn nội dung đúng của nó từ dữ liệu tôi đang có không?
Nếu có, overwrite. Nếu không, merge.
Đó là toàn bộ quyết định. Mọi thứ khác đều là trường hợp đặc biệt của nó.
Hãy xét các immutable event được partition theo thời gian - clicks, logs, transaction. Một click xảy ra lúc 9 giờ sáng ngày 29 có event_date là ngày 29 mãi mãi; nó không bao giờ di chuyển và không bao giờ thay đổi. Vì vậy khi bạn cần reprocess ngày 29, bạn lấy toàn bộ event của ngày 29 từ nguồn, rebuild hoàn toàn ngày đó, và overwrite. Bạn có thể tái tạo toàn bộ đơn vị, nên overwrite là an toàn - và bạn có được idempotency miễn phí, vì rerun ngày 29 chỉ xóa và ghi lại ngày 29 với kết quả như nhau.
Bây giờ xét một bảng orders trong đó một đơn hàng thay đổi trạng thái theo thời gian. Một đơn hàng được tạo ngày 1 tháng 6 nằm trong partition ngày 1 tháng 6 mãi mãi, nhưng trạng thái của nó thay đổi hôm nay. Để cập nhật nó bạn phải overwrite partition ngày 1 tháng 6 - nhưng tập thay đổi hôm nay chỉ chứa một đơn hàng đó, không phải mọi đơn hàng khác từ ngày 1 tháng 6. Overwrite partition sẽ xóa hàng nghìn đơn hàng bạn không thể ghi lại. Bạn không thể tái tạo toàn bộ đơn vị, vì vậy bạn phải merge: khớp hàng đơn bằng order_id và cập nhật nó.
Lưu ý quy tắc chưa bao giờ là “các entity có thể thay đổi không thể bị overwrite.” Cung cấp cho bảng orders đó một full snapshot hàng đêm - mọi đơn hàng ở trạng thái mới nhất - và overwrite quay trở lại ngay, vì snapshot chính là toàn bộ nội dung của bảng. Entity không thay đổi; khả năng tái tạo đơn vị của bạn mới thay đổi.
Partition là đơn vị của reprocessing
Từ “đơn vị” đang làm rất nhiều việc ở trên, và partition là thứ làm cho nó nhỏ lại.
Một partition là một phân tách vật lý của bảng theo một cột nào đó. “Partition by event_date” nghĩa là mọi hàng của ngày 29 nằm cùng nhau, thường trong thư mục riêng của các file, cách ly khỏi mọi ngày khác. Sự cách ly đó mới là điểm quan trọng: bạn có thể overwrite ngày 29 mà không chạm vào ngày 28.
Đây là điều làm cho “tái tạo toàn bộ đơn vị” trở nên hợp lý. Không có partition, việc reprocess ngày 29 có nghĩa là rebuild và overwrite toàn bộ bảng - tốn kém và rủi ro. Với partition, đơn vị bạn phải tái tạo thu nhỏ từ toàn bộ lịch sử xuống còn một ngày. Partitioning là thứ biến reprocessing từ một công việc khổ sai mỗi quý thành điều bạn có thể làm mỗi sáng.
Vậy câu hỏi thực sự, phát biểu chính xác, là: đơn vị tôi định overwrite - một partition, hay toàn bộ bảng - tôi có thể rebuild hoàn toàn từ những gì tôi có không? Partition tồn tại để làm cho đơn vị đó nhỏ nhất theo grain tự nhiên của dữ liệu cho phép.
Hai đảm bảo trông như một
Đây là cái bẫy mà mọi người hay mắc vào. Giả sử bạn overwrite ngày 29 từ một message queue, và bạn rerun pipeline năm lần. Bạn có an toàn khỏi dữ liệu trùng lặp không?
Bạn an toàn khỏi một loại trùng lặp nhưng không phải loại khác, và việc đồng nhất chúng là cách các bug về tính đúng đắn lọt qua.
Overwrite cho bạn idempotency của việc ghi: rerun job bao nhiêu lần cũng tạo ra cùng một partition, vì mỗi lần chạy xóa và ghi lại từ đầu và không bao giờ thấy output của lần chạy trước. Điều này thực sự có giá trị - reruns không bao giờ làm mọi thứ tệ hơn.
Nhưng idempotency của việc ghi không nói gì về việc nội dung có đúng hay không. Nếu dữ liệu nguồn của bạn đã chứa duplicate, overwrite sẽ trung thành ghi những duplicate đó mỗi lần. Hoàn toàn idempotent, và nhất quán sai.
Hai thuộc tính này là trực giao:
- Idempotency của việc ghi (overwrite xử lý cái này): rerun không nhân thêm dữ liệu của bạn.
- Tính đúng đắn của nội dung (thứ khác phải xử lý cái này): không có sự kiện thực tế nào được đếm hai lần bên trong partition.
Số lần bạn chạy pipeline, và các duplicate đã nằm trong nguồn, là các vấn đề khác nhau. Overwrite giải quyết cái đầu. Cái thứ hai cần công cụ riêng của nó.
At-least-once delivery và natural key
Duplicate trong nguồn đến từ đâu? Thường là từ cách messaging layer của bạn deliver dữ liệu.
Hầu hết các queue đảm bảo delivery at-least-once. Cơ chế này đáng hiểu chính xác, vì nó giải thích tại sao cách khắc phục lại nằm ở chỗ nó nằm. Queue gửi event X. Bạn xử lý nó và gửi lại một acknowledgement. Nhưng ack bị mất - một sự cố mạng, hoặc consumer của bạn chết ngay sau khi xử lý nhưng trước khi báo cáo thành công. Queue chờ đợi, không nghe thấy gì, kết luận “họ có thể không nhận được,” và gửi X lại. Không phải vì ác ý, mà vì nó thích deliver hai lần hơn là mất một tin nhắn. Đó là ý nghĩa của at-least-once: “Tôi đảm bảo bạn nhận được nó ít nhất một lần, không phải đúng một lần.”
Vì vậy duplicate không phải là hai bản sao trong một lần delivery. Đó là cùng một event được deliver hai lần, có thể cách nhau vài giây hoặc vài giờ, phân tán vào data stream của bạn. Bạn không thể bắt nó bằng cách kiểm tra một batch đơn lẻ. Bạn bắt nó bằng một natural key.
Một natural key xác định một event theo ý nghĩa nghiệp vụ của nó, không phải theo hàng của nó. Với một click nó có thể là (user_id, page, event_timestamp); với một đơn hàng, order_id; với một thanh toán, transaction_id. Ý tưởng: hai hàng có cùng natural key là cùng một sự kiện thực tế, bất kể chúng đến bao nhiêu lần. Deduplication nghĩa là group theo natural key và giữ một hàng trên mỗi key trước khi ghi.
Đây là lý do tại sao dedup nằm trong data layer của bạn, không phải trong queue. Queue không thể biết liệu X nó đang gửi lại có phải là X cũ hay không - việc gửi lại là hành vi đúng từ quan điểm của nó. Chỉ bạn, ở điểm đích, mới có thể định nghĩa “hai hàng này là cùng một sự kiện thực tế.” Chấp nhận các duplicate, rồi tiêu diệt chúng bằng key. Đó là cách bạn đạt được hiệu quả exactly-once mà không cần một queue exactly-once - đây là cách thực dụng mà các hệ thống thực sự đạt được điều đó.
Chọn key là quyết định nghiệp vụ, không phải kỹ thuật
Đây là phần khó nhất và bị đánh giá thấp nhất. “Unique theo nghĩa nào?” không phải là câu hỏi mà engine của bạn có thể trả lời. Engine chỉ so sánh bytes. Bạn, cùng với người hiểu domain, định nghĩa thế nào là một sự kiện riêng biệt - và định nghĩa đó thay đổi theo nghiệp vụ.
Sai theo một trong hai hướng và bạn nhận được một silent failure, với hình dạng ngược nhau.
Key quá hẹp - giả sử bạn dedup clicks bằng chỉ user_id. Một user thực click Home vào buổi sáng, Products vào buổi trưa, Checkout vào buổi tối: ba sự kiện riêng biệt, thực sự, tất cả chia sẻ một user_id. Group theo user_id, giữ một hàng, và bạn đã thu gọn ba sự kiện thực thành một. Bạn không xóa duplicate; bạn xóa dữ liệu. Đây là over-deduplication, và nó under-count.
Key quá rộng - giả sử bạn đưa processing_timestamp (khi pipeline của bạn xử lý hàng) vào key. Bản gốc và bản gửi lại được xử lý vào các thời điểm khác nhau, nên chúng nhận được các key khác nhau, nên dedup không nhận ra chúng là cùng một event. Duplicate vượt qua. Điều này over-count.
Cùng một bảng payments có thể cần các key khác nhau tùy theo quy tắc. Nếu mỗi giao dịch có transaction_id do gateway cấp, đó là key của bạn. Nhưng nếu nghiệp vụ cho phép nhiều thanh toán cho một đơn hàng - trả góp, hoặc nạp thêm sau khi thanh toán thiếu - thì order_id không unique và (order_id, payment_attempt) mới là. Cùng dữ liệu, định nghĩa khác nhau về “unique,” vì logic nghiệp vụ khác nhau.
Vậy quy trình chọn key thực sự là một câu hỏi cho domain: điều gì làm cho hai bản ghi là cùng một event, và điều gì làm cho chúng là các event khác nhau? Câu trả lời chính là natural key. Quy tắc suy ra:
Một natural key phải chứa đủ để phân biệt hai sự kiện thực sự khác nhau (quá ít -> over-dedup -> under-count), và không chứa bất cứ thứ gì không gắn với chính sự kiện đó - không bao giờ là thứ gì pipeline tạo ra, như processing time hay row id (quá nhiều -> duplicate lọt qua -> over-count).
Và vì một key sai không bao giờ crash bất cứ thứ gì, đây là một silent bug theo cấu trúc. Không có exception, không có cảnh báo. Chỉ là những con số âm thầm sai.
Incremental model và lookback window
Phần cuối cùng kết nối tất cả lại ở nơi những công cụ này thực sự được sử dụng.
Ngược lại của một incremental model là full refresh: tính toán lại toàn bộ bảng từ đầu mỗi lần chạy. Ổn với dữ liệu nhỏ; vô lý khi mỗi buổi sáng có nghĩa là tính toán lại nhiều năm lịch sử. Một incremental model chỉ xử lý những gì mới hoặc thay đổi kể từ lần chạy cuối, rồi fold nó vào bảng hiện có. Hiệu quả - và mong manh, theo cách kéo theo mọi thứ ở trên.
Phiên bản naive theo dõi một watermark: “chỉ lấy các hàng có timestamp lớn hơn cao nhất tôi thấy lần trước.” Vấn đề là late data. Một event xảy ra hôm qua nhưng đến hôm nay, sau khi watermark đã di chuyển qua hôm qua. Nó không bao giờ được xử lý. Nó biến mất, trong im lặng.
Cách khắc phục là một lookback window: thay vì chỉ lấy những gì nghiêm ngặt mới hơn watermark, mỗi lần chạy reprocess một trailing window - giả sử “luôn luôn tính toán lại 3 ngày cuối.” Một late event rơi vào trong window đó vẫn được bắt. Bạn đánh đổi một chút công việc thêm để không mất late data.
Nhưng window tạo ra đúng vấn đề trùng lặp từ trước đó: nếu mỗi lần chạy reprocess 3 ngày cuối và append, những ngày chồng lấp đó xếp chồng lên nhau và nhân lên. Vì vậy một lookback window buộc phải có idempotent writes - hoặc partition-overwrite các ngày trong window, hoặc merge/upsert với dedup theo natural key. Toàn bộ bộ công cụ từ đầu, trong một pattern.
Lookback là cái giá của việc không xem dữ liệu đến
Đây là mental model tôi tiếp tục quay lại. Một lookback window là chi phí của việc trigger theo lịch thay vì theo dữ liệu.
Một scheduler (cron, một Airflow run hàng ngày) fire theo thời gian. Nó thức dậy lúc 2 giờ sáng và tuyên bố “đến lúc xử lý hôm qua” - dù dữ liệu của hôm qua có thực sự đến hay không. Nó không thể biết; nó chỉ biết đồng hồ. Lookback window là câu trả lời cho sự mù quáng đó: “vì tôi không thể biết vài ngày cuối có đầy đủ chưa, tôi sẽ quét lại chúng mỗi lần chạy, phòng khi.” Đó là một bộ đệm phòng thủ chống lại độ trễ mà scheduler không thể quan sát.
Một event-driven trigger thì ngược lại. Nó fire vì dữ liệu đến - một file rơi vào S3, một tin nhắn đến queue. Late event, dù đến lúc nào, tạo ra trigger riêng của nó và được xử lý ngay lúc đó. Không cần quét lại ba ngày vì có thể có thứ gì đó đến muộn, vì sự muộn tự thông báo.
Phép ẩn dụ: một scheduler giống người thu thư vào một giờ cố định, người - không biết khi nào thư đến - phải kiểm tra lại hộp thư của vài ngày trước để chắc không có thứ gì đến muộn. Event-driven là chuông cửa: một lá thư đến, nó reo, bạn lấy cái đó. Không cần kiểm tra lại.
Điều này định hình lại cách bạn chọn độ dài window N. Với scheduler, N là một cược về độ trễ tối đa của nguồn. Quá ngắn, và bất kỳ event nào muộn hơn window sẽ nằm ngoài mọi lần quét và bị mất mãi mãi - đó là lỗi tính đúng đắn. Quá dài, và bạn xử lý lại các partition đã ổn định từ lâu - đó là lỗi chi phí. Chọn N từ phân phối độ trễ thực tế của nguồn (percentile cao, cộng thêm safety margin), không phải từ một con số tròn.
Và đây là lỗi đầy đủ: một nguồn đôi khi muộn 2 ngày, đằng sau một daily scheduler với lookback 1 ngày. Ngày record muộn cuối cùng đến, pipeline chỉ quét lại hôm qua và hôm nay - nhưng record thuộc về hai ngày trước, ngoài window, đã qua watermark. Không có lần chạy nào nhìn vào nó nữa. Nó biến mất. Job màu xanh, dashboard hoạt động, aggregate âm thầm thiếu. Một watermark và lookback đặt sai biến một data bug thành một bug vô hình.
Sợi chỉ xuyên suốt
Mỗi phần ở đây là cùng một bài học mặc những bộ quần áo khác nhau. Bảng quyết định, partition grain, các đảm bảo trực giao, natural key, độ dài lookback - trong mỗi trường hợp, sai lầm nguy hiểm không throw error. Nó thành công, đúng lịch, với những con số sai.
Vì vậy hãy thiết kế cho điều đó. Hỏi xem bạn có thể tái tạo đơn vị trước khi overwrite nó không. Giả định nguồn của bạn có duplicate và định nghĩa key tiêu diệt chúng. Giả định dữ liệu đến muộn và điều chỉnh window cho phù hợp. Xây dựng pipeline giả định thất bại - vì những thất bại tệ nhất là những thất bại trông chính xác như thành công.