Source contract
Define keys, schema, insert and update and delete semantics, timestamps, ordering, partitioning, expected volume, ownership, access, retention, and the source change path.
Source change to usable state
BluePi builds change-data-capture and event paths around a defined operating freshness target. The system keeps source position, schema version, processing state, publish watermark, consumer freshness, replay history, and reconciliation evidence visible.
Fresh operational data with visible recovery state
A real-time data service connects source position, event semantics, processing state, consumer watermark, replay history, and reconciliation evidence to the operating decision that needs current data.
The architecture begins with the required freshness and completeness. Capture technology, processing, storage, and serving follow that service contract.
Define keys, schema, insert and update and delete semantics, timestamps, ordering, partitioning, expected volume, ownership, access, retention, and the source change path.
Connect the initial snapshot to a known change position. Record checkpoints, idempotency, quarantine, replay range, backfill isolation, promotion, and safe consumer behavior.
Publish schema, source coverage, latest processed position, watermark, freshness, completeness, known gaps, correction behavior, replay status, and support ownership.
When this is the right starting point
Good fit
One operating workflow has a named owner, a measurable baseline, and users who can judge whether the result improves.
Poor fit
The request is capacity-only staffing, an unowned demonstration, or a broad transformation without a first decision and finish condition.
Evidence produced during delivery
Open the part you need. Each section expands into the full delivery scope for that step.
Real time has meaning only when the team names the decision, source event, consumer state, tolerated delay, completeness condition, recovery objective, owner, and consequence of stale or missing data.
Each source change or event records its source object, primary or business key, insert or update or delete semantics, correction behavior, event identifier, source position, event and commit and arrival times, schema version, partition and ordering rule, expected volume, access requirements, retention, source owner, and change-notification path.
A continuous path starts with an initial source state and then applies changes from a known position. The system records the snapshot boundary, source position, extraction range, validation result, and the point where continuous processing takes ownership so the handoff introduces no hidden gaps or repeated changes.
The ingestion layer persists a durable source position or checkpoint. Stable event and operation identifiers make retries safe.
Each consumer contract states the available object or topic, key, schema version, source coverage, latest processed position, event and publish watermark, freshness and completeness state, known gaps, quarantined records, correction and delete behavior, replay or backfill status, and owner. Consumers can distinguish current, delayed, incomplete, and unavailable data.
Replay records the source and event-time range, schema and transformation versions, destination scope, consumer impact, start and finish positions, approving owner, and reconciliation result. Backfill can run in an isolated path before promotion when consumers cannot safely receive historical or repeated events.
Reconciliation can compare counts by source and operation, key presence, duplicates, checksums, watermarks, source positions, field rules, sampled records, business totals, quarantine, and unresolved differences. Monitoring separates source, capture, processing, publish, consumer, schema, replay, reconciliation, throughput, and cost layers. See it in production in the logistics change-data-capture path on Google Cloud and the Delhivery reporting case study.
Define the source commit or event point, consumer state, capture and processing and publish timestamps, freshness target, completeness or watermark condition, tolerated delay, outage behavior, recovery objective, and owner. Measure the full source-to-consumer path.
Use stable source positions, event identifiers, schema and transformation versions, idempotent processing, bounded time and key ranges, isolated validation where needed, consumer-impact rules, reconciliation checks, an approving owner, and a recorded promotion result.
Start with one source, bounded objects or event types, and one consumer. Define freshness, completeness, snapshot-to-stream handoff, checkpoint ownership, replay range, schema behavior, reconciliation evidence, and the safe response to late or unavailable data.
We will review the owner, the baseline, the data path, the system boundary, and the route to go-live.
Discuss one workflow