Skip to content
KEDBYTE
Site navigation
How Data Works
Chapter
40

Data Pipelines: Capture, Transform and Deliver

Part G · From Records to Insight|4,030 words|about 18 min read|Volume G

40.0 What this chapter gives you#

  1. A data pipeline carries information from one representation or system into another. It may copy, validate, reshape or summarise the records, but it must not silently change what they mean.
  2. You will follow a frozen export of the shop’s four agreed order lines into a small reporting store. The questions are practical: Which source was read? What was accepted? What was changed? Where does progress stop after a failure? Can the result be reproduced?
  3. The companion exercise is a local, bounded batch import. The discussion of change capture uses published PostgreSQL and Debezium documentation; it does not claim that a connector, Kafka cluster or production pipeline was deployed.

40.1 Sources and contracts#

40.1.1 PLAIN — in simple words#

  1. A pipeline needs an agreement about its input before it needs a schedule. That agreement describes what one record means, which fields identify it, which values are allowed and which changes can happen later.
  2. “Sales data” is not a sufficient description. It could mean agreed orders, dispatched goods, payments or refunds. Those are different events, with different amounts and dates.
  3. A receiving system should preserve enough context to explain each result. A total without its population, units, source boundary and transformation rule is difficult to verify, even when the arithmetic is correct.

40.1.2 PLAIN — a picture in your head#

  1. Mira sends Dev a sealed envelope containing copies of order lines. On the envelope she writes which register supplied them, the export date and the rule used to choose them.
  2. Dev writes a second note describing how he grouped the copies. Together, the two notes explain both what arrived and what he did with it.
  3. Where the comparison breaks: a digital checksum can bind a report to exact bytes, but it cannot prove who created those bytes or that the source claims are true. Authentication and business verification remain separate responsibilities.

40.1.3 PLAIN — a worked example#

  1. Our new teaching export, PIPE-01, contains the same four agreed lines introduced earlier. It is a frozen copy, not a claim about a live capture time. Its grain is one agreed order line, identified by (order_id, line_no).
  2. The quantity unit is individual items. Prices are integer paise in INR. Multiplying each quantity by its agreed unit price gives 15,100, 2,000, 7,550 and 4,000 paise.
  3. The input contract therefore predicts four distinct line keys, six units, two orders and 28,650 paise. These are useful reconciliation checks, not evidence that all possible real orders were included.
  4. The contract deliberately excludes payments, discounts, taxes, shipping and refunds. A report derived from PIPE-01 must not acquire a label such as “net cash collected” merely because that label sounds useful.

40.1.4 PLAIN — what is really happening inside#

  1. The producer serialises records according to a schema. The consumer decodes them, checks structural rules and then applies domain rules. These are different stages: valid JSON can still contain an invalid price.
  2. An identifier for the schema or transformation helps the consumer select the right interpretation. A later field rename or unit change must not be inferred solely from a familiar-looking value.
  3. Provenance records connect the input version, processing code, parameters and output version. This connection makes reproduction and investigation possible without pretending that provenance proves correctness.

40.1.5 TECHNICAL — the engineer’s version#

  1. Define grain, identity, domain types, null semantics, ordering scope, update/delete representation and compatibility policy. A contract should also specify size bounds and what happens to malformed or conflicting input.
  2. The pipeline manifest in this book records a byte digest and transformation version. The digest establishes content identity under the selected hash, not origin authenticity or semantic completeness. [S99]
  3. Source-table changes and business events must not be conflated. Logical decoding exposes database changes; deciding that a change means “a customer paid” requires additional application semantics. [S166]

40.1.6 WORDS — remember these#

  1. Data contract: the agreement about incoming information — explicit structure, semantics, identity and compatibility rules for a data interface. Provenance: the recorded origin and processing history — information linking source entities, transformations and resulting data. Transformation version: which processing rule was used — an identifier binding output to a particular implementation and its parameters.

40.2 Batch extraction#

40.2.1 PLAIN — in simple words#

  1. A batch is a bounded collection processed together. Its boundary may be a frozen file, a database snapshot or an explicitly closed interval of source events.
  2. Reading a changing database in several pieces can produce a mixture that never existed at one moment. Each piece may be individually valid while the combination is misleading.
  3. A successful extraction needs both complete coverage of its intended boundary and a consistent interpretation of that boundary. “The command finished” establishes neither by itself.

40.2.2 PLAIN — a picture in your head#

  1. Dev photographs the first half of a noticeboard. Mira then replaces several notices before he photographs the second half. Joining the photographs does not create a reliable picture of the board at one time.
  2. Freezing the board during both photographs would solve that particular problem. Keeping a versioned copy of the board would solve it differently.
  3. Where the comparison breaks: database snapshot visibility can avoid stopping all writers. The exact isolation and export mechanisms are engine-specific, and a file copied outside those mechanisms is not automatically consistent.

40.2.3 PLAIN — a worked example#

  1. Use a separate two-account teaching model, not the shop’s financial records. Initially A contains 60 units and B contains 40, for a total of 100.
  2. An extractor reads A as 60. Another operation transfers 10 from A to B. The extractor then reads B as 50 and reports 110.
  3. The source still contains 100. The error is the inconsistent observation boundary, not an addition mistake. A snapshot that sees either the complete before-state or complete after-state reports 100.
  4. PIPE-01 avoids this race by using a fixed byte sequence as its input. That makes the local exercise repeatable, but it does not establish that an arbitrary real export was originally captured consistently.

40.2.4 PLAIN — what is really happening inside#

  1. A database extractor issues queries under a particular transaction and isolation configuration. Different statements may see different committed versions unless the chosen mechanism supplies a suitable shared snapshot.
  2. Pagination is another boundary. An offset into a changing result set can skip or repeat rows as earlier rows appear or disappear. A stable ordering key helps, but a key alone does not freeze the population.
  3. A manifest should identify the selection predicate, ordering, snapshot or source cut, row count and any exclusions. When the source cannot provide a consistent cut, the limitation belongs in the result, not only in an operator’s memory.

40.2.5 TECHNICAL — the engineer’s version#

  1. PostgreSQL 17 Read Committed gives statement-level snapshots; stronger transaction-level snapshot behaviour requires an appropriate isolation choice and operational design. A multi-statement export must account for that distinction. [S66]
  2. Incremental extraction based only on an increasing row ID can miss updates and deletions to older rows. A timestamp cursor also needs tie-breaking, precision, clock and late-write assumptions.
  3. This book’s immutable batch uses (batch_id, input_digest, transform_version) as declared identity. Reusing the batch ID with different input bytes is rejected, even when someone believes the new bytes are equivalent. That is an original teaching policy, not a universal connector requirement.

40.2.6 WORDS — remember these#

  1. Extraction boundary: exactly which source state or records were selected — the snapshot, interval or immutable object defining a batch’s coverage. Consistent cut: observations that fit one permitted source history — a boundary that avoids incompatible mixtures of concurrent states. Incremental extraction: reading changes since recorded progress — a method whose cursor must account for the source’s actual change semantics.

40.3 Change capture#

40.3.1 PLAIN — in simple words#

  1. Instead of repeatedly copying everything, a pipeline can follow changes: a row was inserted, updated or deleted. This is change data capture, usually shortened to CDC.
  2. A new subscriber still needs a starting picture. Otherwise, it sees tomorrow’s changes without knowing what existed today.
  3. Joining that starting picture to the continuing changes is the difficult part. A gap loses information; an overlap can repeat information. Both need an explicit protocol rather than an approximate timestamp guess.

40.3.2 PLAIN — a picture in your head#

  1. Dev copies the entire catalogue while Mira keeps a numbered list of changes. The copy has a clear marker saying which changes are already reflected in it.
  2. Dev then follows the list from the agreed boundary. If a page of the list is missing, he cannot safely claim that his copy is current.
  3. Where the comparison breaks: database logs can describe physical or logical operations rather than human catalogue edits. Transaction boundaries, replica identity and retention rules determine what a real connector can reconstruct.

40.3.3 PLAIN — a worked example#

  1. In a separate model, a snapshot at sequence 100 says product X has quantity 5. Sequence 101 changes it to 3; sequence 102 deletes it.
  2. Applying both changes after the snapshot leaves X absent. Starting at 102 and missing 101 happens to give the same final presence here, but would lose an intermediate change required by an audit-oriented consumer.
  3. Starting at 103 misses the deletion and incorrectly retains X with quantity 5. A final-state check for this key now catches the gap, but a few spot checks cannot prove full coverage.
  4. Replaying 101 after 102 must not recreate the old state. The consumer needs a suitable version/order policy, or an ordered replay mechanism that prevents this stale update from becoming current.

40.3.4 PLAIN — what is really happening inside#

  1. A connector obtains a starting snapshot and records an associated log position. It then consumes changes from a boundary coordinated with that snapshot.
  2. Source logs are finite operational resources. A stalled consumer can require old log segments to remain available, increasing storage use or eventually invalidating the ability to resume.
  3. A delete requires enough identity to identify the destination record. A partial update may require previous state or configuration-specific fields. The consumer must understand the actual event envelope rather than treating every message as a complete new row.

40.3.5 TECHNICAL — the engineer’s version#

  1. Debezium’s PostgreSQL connector documents a coordinated initial snapshot followed by log streaming. Snapshot modes, restart behaviour and captured fields are connector/version concerns, not properties of every pipeline called CDC. [S165]
  2. PostgreSQL logical slots can resend recent changes after a crash because persisted progress may precede the last delivered position. Clients must avoid harmful repeated handling. A received message is not proof of a unique business effect. [S166]
  3. PostgreSQL logical replication does not replicate schema DDL and sequence state in the same manner as table rows. Schema evolution and failover therefore require separate planning. Do not infer a complete database clone from a table-change feed. [S164]

40.3.6 WORDS — remember these#

  1. Change data capture: following source mutations — obtaining inserts, updates and deletions for downstream processing. Snapshot-to-stream handover: joining the initial copy to later changes — a protocol preventing gaps and managing overlap at a capture boundary. Replication slot: retained source progress for a consumer — a PostgreSQL mechanism whose log retention and restart behaviour require monitoring.

40.4 Transformations#

40.4.1 PLAIN — in simple words#

  1. A transformation changes representation or derives a new value. It might convert units, choose columns, join descriptions or group lines into totals.
  2. Some transformations lose detail. Once four lines become two order totals, the original line quantities cannot be recovered from the totals alone.
  3. The processing rule must say what happens to missing, invalid and duplicate data. Silently turning all three into zero makes the output simpler by destroying distinctions the reader may need.

40.4.2 PLAIN — a picture in your head#

  1. Dev sorts four order slips into two envelopes and writes a total on each. The envelope totals are useful, but they are not substitutes for every detail on the slips.
  2. If one slip is unreadable, writing zero on it is an invented correction. Keeping it for review preserves the uncertainty instead.
  3. Where the comparison breaks: software can retain exact source references and repeat the calculation automatically. That helps only when the code, inputs and policy are recorded and the process does not depend on unrecorded outside state.

40.4.3 PLAIN — a worked example#

  1. Group PIPE-01 by order ID. O-1042 contributes 15,100 plus 2,000, giving 17,100 paise. O-1043 contributes 7,550 plus 4,000, giving 11,550.
  2. The two totals sum to 28,650. Dividing by two orders gives 14,325 paise per order. Dividing by six items instead gives 4,775 paise per item. These answer different questions.
  3. Joining the current notebook catalogue price of 8,250 into the calculation would produce the previously discussed hypothetical 30,750 total. That is not the historical agreed total and must not replace it.
  4. Joining product tags before summing can repeat line values and produce 51,300. The transformation may be syntactically valid while changing the grain and therefore the answer.

40.4.4 PLAIN — what is really happening inside#

  1. A deterministic transformation gives the same result for the same complete inputs and rules. Looking up today’s exchange rate or catalogue value creates another input, even when it is not written in the main file.
  2. Filtering changes the population. Joins can remove rows or multiply them. Rounding can make summing rounded parts differ from rounding a final exact sum.
  3. A useful transformation record states the output grain, keys, units, retained exclusions and expected invariants. Its test cases include counterexamples, not only a clean happy-path batch.

40.4.5 TECHNICAL — the engineer’s version#

  1. Treat transformation code, configuration and reference data as versioned inputs. Reproducibility requires all dependencies affecting the result, not merely a script filename.
  2. Relational joins and aggregation operate on the selected rows and multiplicities. DISTINCT is not a general repair for a wrong join: it can remove legitimate equal-valued facts. [S58] [S80]
  3. The local reporting transformation uses integer arithmetic and original line identity. It does not infer a sale’s financial, legal or physical completion from an order record. Those remain outside the teaching contract.

40.4.6 WORDS — remember these#

  1. Deterministic transformation: the same complete inputs yield the same output — a rule without uncontrolled time, randomness or external-state dependencies. Lossy transformation: some source detail is discarded — a mapping that cannot generally be inverted to recover the original records. Output grain: one resulting row means one what — the level at which a transformed record states a fact.

40.5 Delivery and checkpoints#

40.5.1 PLAIN — in simple words#

  1. A checkpoint records where processing can safely resume. It must correspond to work actually preserved, not work merely started or displayed on a progress bar.
  2. If the checkpoint moves ahead before the output is saved, a restart can skip missing output. If output is saved before progress is recorded, a restart can repeat it.
  3. A correct design either commits output and progress together within a suitable boundary, or makes repeated delivery safe. The correct choice depends on which systems participate.

40.5.2 PLAIN — a picture in your head#

  1. Dev ticks a slip after filing it. Ticking first risks losing the slip after the tick. Filing first risks filing it again after an interruption.
  2. A numbered filing system can recognise an already filed slip and check whether the repeated contents agree. That turns some repeats into harmless replays.
  3. Where the comparison breaks: a database transaction can bind several local records atomically. Two unrelated remote systems do not become one atomic filing cabinet merely because one program calls them in sequence.

40.5.3 PLAIN — a worked example#

  1. Our local batch publisher begins a transaction, inserts the four reporting lines and records PIPE-01’s identity in the same database. It commits only after all checks pass.
  2. Inject an exception after the second insert. Rollback removes the new destination lines and the uncommitted batch record. The retry starts from the same frozen input, not from an invented “half complete” marker.
  3. After a successful commit, submit PIPE-01 again with the same digest and transformation version. The stored result is returned without inserting four more lines or doubling 28,650 to 57,300.
  4. Submit different bytes under PIPE-01. The publisher rejects the identity conflict. This does not overwrite earlier output, and it does not treat a matching batch label as permission to replace history.

40.5.4 PLAIN — what is really happening inside#

  1. Local atomic publication stores destination data and the receipt for that data within one transaction. Unique keys constrain repeated insertion, while explicit payload checks distinguish replay from conflict.
  2. Remote destinations require their own acknowledgement and replay rules. A checkpoint in a source connector does not automatically prove that every downstream report has incorporated the same prefix.
  3. Partial progress also needs a defined scope. A checkpoint for partition A says nothing about partition B unless a higher-level checkpoint records both and explains their relationship.

40.5.5 TECHNICAL — the engineer’s version#

  1. Bind progress to the destination effect it claims. In the local exercise, the transaction owns both rows and batch receipt; no cross-system exactly-once guarantee is asserted. [S25] [S86]
  2. A production streaming checkpoint may capture offsets, operator state and sink-commit information. End-to-end correctness depends on compatible source and sink protocols, not just a framework’s internal recovery mechanism.
  3. Retrying an external side effect requires stable request identity and a policy for unknown outcomes. The outbox/inbox construction from Chapter 26 applies when its local atomicity and consumer deduplication assumptions hold. [S125] [S126]

40.5.6 WORDS — remember these#

  1. Checkpoint: recorded safe progress — a recovery marker bound to preserved state and explicitly scoped effects. Atomic publication: output becomes accepted as one local change — committing records and their completion evidence within one transaction boundary. Replay-safe delivery: repetition does not duplicate the intended effect — handling a repeated identity consistently while detecting conflicting payloads.

40.6 Reconciliation and replay#

40.6.1 PLAIN — in simple words#

  1. Reconciliation compares the source and result using checks chosen for the transformation. It asks whether the expected records and quantities survived, not merely whether a job reported success.
  2. A matching total is useful but incomplete. Two wrong amounts can cancel each other. Missing one record and duplicating another can preserve the row count.
  3. Replay rebuilds output from recorded inputs. It is a controlled experiment when the inputs and rules are fixed; it is a different computation when either changes.

40.6.2 PLAIN — a picture in your head#

  1. Mira checks that Dev filed four slips, then checks their identifiers and values. Counting envelopes alone would miss a slip placed twice and another omitted.
  2. She can ask him to repeat the grouping from the sealed copies. If the totals change, they can inspect the procedure rather than arguing from memory.
  3. Where the comparison breaks: digital equality can compare every included key and value, but it still cannot reveal records that the source never captured or an undisclosed exclusion outside the manifest.

40.6.3 PLAIN — a worked example#

  1. PIPE-01’s report should have exactly the four source keys, quantities totalling six and agreed amounts totalling 28,650. Check each line’s amount as well as the aggregate.
  2. In a deliberately damaged copy, add 100 paise to one line total and subtract 100 from another. The grand total stays 28,650, but keyed comparison reveals both disagreements.
  3. Rebuild the clean output with the same input digest and transformation version. Compare ordered key/value records or a canonical representation, not a hash of unspecified row order.
  4. A proposed new report definition should receive a new transformation version and its own results. Keep a comparison explaining the difference rather than relabelling the old report as though it always meant the new thing.

40.6.4 PLAIN — what is really happening inside#

  1. Reconciliation combines coverage checks, uniqueness checks, value comparisons and domain invariants. The right combination follows the contract: a filtered report should reconcile its included and excluded populations separately.
  2. A reproducible run stores parameters, versions and output evidence. This allows failures to be investigated without touching the original source records or silently changing the test expectations.
  3. Operational monitoring tracks lag, rejected input, unavailable log history, replay conflicts and destination failures. A green schedule is not enough when useful data is steadily falling behind.

40.6.5 TECHNICAL — the engineer’s version#

  1. Record source and destination watermarks with their meanings. A maximum observed source position is not necessarily a contiguous applied prefix, and offsets from different streams are not universally comparable integers.
  2. Use keyed reconciliation to supplement aggregate checks. Canonical serialisation must specify field order, type representation, text encoding and row order before a digest can meaningfully compare runs. [S45] [S99]
  3. The lab’s failure injection and replay checks establish local behaviour for frozen teaching inputs. They do not test connector failover, source-log loss, unbounded streams, external delivery or every possible malformed-file attack.

40.6.6 WORDS — remember these#

  1. Reconciliation: comparing expected and observed information — a set of coverage, identity, value and invariant checks across representations. Contiguous progress: every required earlier position is included — a stronger condition than merely having seen a large position number. Reprocessing: applying rules again to retained inputs — a controlled rebuild whose interpretation depends on input and transformation versions.

40.97 Practice and worked answers#

  1. Question: Why is a pipeline labelled “sales” ambiguous? Answer: Orders, dispatches, payments and refunds describe different facts. Specify which event or state is included before assigning the metric a name.
  2. Question: An extractor reads 60 from A and later 50 from B after a 10-unit transfer. What failed? Answer: The observation boundary mixed before and after states. Addition is correct; the claimed population is not one consistent state.
  3. Question: Can an increasing primary-key cursor capture all changes? Answer: Not by itself. Updates and deletions to existing smaller keys can be missed.
  4. Question: Why does the batch receipt share a transaction with reporting rows? Answer: A receipt must not claim completion without the output, and a retry must recognise already committed output without duplicating it.
  5. Question: How many rows and paise should a replay of PIPE-01 add? Answer: Zero new rows and zero new amount when the stored batch identity agrees. The retained four rows still total 28,650.
  6. Question: Does a matching grand total prove record equality? Answer: No. Opposite 100-paise errors cancel. Compare complete keys and values as well as aggregates.
  7. Question: Does a source connector’s checkpoint prove a dashboard is current? Answer: No. Every downstream stage has its own applied progress and publication boundary.
  8. Question: What should change when the transformation definition changes? Answer: Its version and derived result identity, with a recorded comparison and a deliberate replacement or coexistence policy.

40.98 Common wrong ideas#

  1. Wrong: A pipeline is just a scheduled copy. Right: It also carries contracts, progress, failure policy and evidence about meaning.
  2. Wrong: A checksum authenticates the producer. Right: It identifies bytes under a hash; authenticity needs a separate trust mechanism.
  3. Wrong: Every valid input record belongs in the report. Right: Eligibility follows the declared population and purpose.
  4. Wrong: A snapshot followed by any recent log position is sufficient. Right: The boundary must prevent gaps and handle overlap.
  5. Wrong: Repeating a job is harmless because it reads the same file. Right: Destination effects need replay safety too.
  6. Wrong: CDC captures the full meaning of a business action. Right: Row changes require application interpretation.
  7. Wrong: Matching counts and totals prove all records are correct. Right: Keyed differences can remain hidden.
  8. Wrong: A successful local import proves production delivery guarantees. Right: It proves only the bounded behaviours actually exercised.

40.99 Chapter summary in 20 lines#

  1. A pipeline moves and transforms information under a contract.
  2. Define source grain before writing extraction code.
  3. Preserve units, identity and missing-value meaning.
  4. A digest binds bytes, not truth or origin.
  5. A batch needs an explicit observation boundary.
  6. Changing sources can produce inconsistent multi-query exports.
  7. An increasing ID alone does not capture every mutation.
  8. CDC follows inserts, updates and deletions.
  9. An initial snapshot must connect correctly to the change stream.
  10. Log retention limits how a consumer can resume.
  11. Redelivery is a normal recovery possibility.
  12. Transformations can discard detail and change grain.
  13. Historical prices are not today’s catalogue prices.
  14. Record all inputs that affect a derived answer.
  15. Progress must correspond to preserved output.
  16. One local transaction can bind output and its receipt.
  17. Remote effects need additional delivery and replay rules.
  18. Reconcile keys and values as well as counts and totals.
  19. Reprocessing with new rules creates a new result version.
  20. State exactly which pipeline guarantees the evidence demonstrates.

Return to contents