Records plus transaction/change metadata.
SQL databasePROJECT WALKTHROUGH / CDC
CDC Data Pipeline
Process inserts, updates, and deletes while keeping downstream records consistent and recoverable.
Illustrative change feed and ordered offsets.
CDC conceptAppend-only change records with operation type.
Python · filesDeduplicate, order, and apply upsert/delete rules.
SQLMaintain a queryable latest record per key.
SnowflakeExpose current views and change lineage.
SQL viewsUnderstand the system before building it.
A change data capture walkthrough covering ordered changes, replay safety, delete handling, checkpointing, and downstream current-state tables. It uses generic operational data and does not connect to a real CDC stream.
Business problem
Make the data need understandable before choosing tools or designing pipelines.
A generic subscription service updates customer and service records throughout the day. Downstream analysis needs a current state plus enough history to explain when changes occurred.
Data engineering provides a repeatable way to move, validate, transform, and make data available with clear ownership and operational expectations.
The design gives consumers a reconciled current-state view and an auditable change path while making replay and delete semantics explicit.
Requirements
Separate what the workflow must do from the reliability and operational qualities it needs.
- Capture inserts, updates, and deletes
- Apply changes in a deterministic order
- Support replay from a known checkpoint
- Use configuration rather than hard-coded environment-specific values.
- Add validation, structured run logging, and clear failure boundaries.
- Protect credentials and grant only the access each workload needs.
- Measure freshness, duration, volume, and quality outcomes.
Architecture
A layered view of how data moves from source systems to a useful consumer-facing output.
Records plus transaction/change metadata.
SQL databaseIllustrative change feed and ordered offsets.
CDC conceptAppend-only change records with operation type.
Python · filesDeduplicate, order, and apply upsert/delete rules.
SQLMaintain a queryable latest record per key.
SnowflakeExpose current views and change lineage.
SQL viewsSource systems
Map each source to its data shape, arrival pattern, ingestion choice, and likely failure modes.
Subscription service database
Subscriber, plan, and service-status changes
Continuous or micro-batch changes
Consume ordered change records from a supported source mechanism
Implementation flow
Build the pipeline in observable stages so each boundary can be tested and recovered independently.
Data ingestion
- Define the source change mechanism and its ordering key.
- Persist the source offset with each landed batch.
- Treat at-least-once delivery as possible and design for replay.
- Advance checkpoints only after target changes are committed.
Transformation
- Normalize operation type and source timestamps.
- Deduplicate repeated event IDs or source offsets.
- Apply updates by business key and deterministic sequence.
- Represent deletes explicitly rather than treating them as missing updates.
# Illustrative transformation outline
valid = [row for row in records if is_valid(row)]
rejected = [row for row in records if not is_valid(row)]
write_curated(valid)
write_quarantine(rejected, run_id=run_id)Storage / warehouse
- Retain append-only changes for audit within a defined retention plan.
- Maintain a current-state table keyed by business identity.
- Consider history tables where consumers need point-in-time state.
- Keep checkpoint and run metadata separate from business entities.
Data quality
Make quality expectations explicit and decide what happens when a record or batch does not pass.
- Check event ordering and uniqueness by source offset.
- Validate required business keys and operation values.
- Reconcile change counts against applied inserts, updates, and deletes.
- Monitor checkpoint lag and unexpected gaps.
Operations, performance, security & monitoring
A realistic project also explains how the system behaves when inputs change, work slows down, or something fails.
09 / Error handling
- Retry transient transport errors without committing a checkpoint.
- Quarantine malformed events with their source offset.
- Make replay idempotent by using event identity and merge keys.
- Define a recovery plan if source retention no longer covers the checkpoint.
10 / Performance
- Process bounded micro-batches and avoid repeated full scans.
- Cluster or partition based on measured access patterns.
- Reduce unnecessary state rewrites and small-file generation.
- Measure lag, throughput, and merge cost together.
11 / Security
- Protect customer identifiers and sensitive attributes.
- Separate capture, transform, and consumer roles.
- Use managed secret storage in an actual deployment.
- Audit who can read raw changes and history.
Monitoring and observability
Monitor system health and data health together: pipeline status alone does not tell you whether the delivered data is fresh and complete.
- Track source offset lag, event volume, and apply duration.
- Alert on gaps, dead-letter growth, and stale current-state tables.
- Compare change operations received with operations applied.
- This is an observability plan, not a running CDC service.
CDC Data Pipeline monitoring signals shown as a design concept. No live pipeline or alert integration is connected.
Deployment flow
Treat infrastructure, SQL, configuration, and validation as reviewed changes that move through separate environments.
Deployment workflow is a learning design. No CI/CD pipeline is implemented by this walkthrough.
Challenges and solutions
Strong project explanations show the problem-solving process, not only the happy path.
The same change arrives twice after a retry.
Use a stable event ID or source offset and idempotent merge semantics.
An update arrives after a later event for the same key.
Use source ordering metadata and define how late events affect current state.
A delete is represented by a tombstone rather than a full row.
Model delete operations explicitly and preserve the key and event metadata needed downstream.
How to explain this project in an interview
Use this outline to structure a truthful explanation. Adapt it to work you personally completed; this sample is not a claim about your experience.
EXAMPLE INTERVIEW EXPLANATIONExample interview explanation: “This CDC design consumes ordered source changes, lands them append-only with operation type and offset, and applies them idempotently to a current-state table. The checkpoint advances only after changes are committed, so retries can safely replay a bounded interval. I would explicitly handle deletes, duplicates, and late events, and reconcile received versus applied operation counts. Monitoring would focus on lag, gaps, failures, and downstream freshness. The actual capture mechanism depends on the source database and its retention guarantees.”
Practice the follow-up.
How do you distinguish event time from processing time?
How do you handle duplicate or out-of-order changes?
When should a checkpoint advance?
How would you represent deletes and preserve history?
What signals show CDC lag or data loss?
Work through the build.
Tick items as you explore. This checklist is session-only and is not saved.