Transaction Data Ingestion Pipeline

AboutSeptember 2025

05· Data EngineeringInternship project, PT Telkom Indonesia

The nightly pipeline I rebuilt at Telkom: it moves roughly 10,000 transaction records a day, asks the source a cheap question before starting an expensive sync, and records what happened on every run, including the runs where the answer was nothing.

ETL
How a run is triggered and carried outOrchestration
A Dagster sensor polls the source database for rows newer than the last one it saw. No new rows and it returns without doing anything; new rows and it triggers the acquisition job, which calls the loader and polls until the sync finishes. Either outcome (success or failure) is written to a log table with its metadata, so a failed run is a recorded event rather than a silent no-op.
Write-up

Moving data between two databases sounds like it should be a single step, and the reason it usually is not is that most of the work goes into deciding whether to move anything at all. This pipeline carries roughly 10,000 transaction records a day from one PostgreSQL database into another, and keeps that decision separate from the transfer itself: Dagster decides whether a sync should happen and watches it once it does, while Airbyte performs the actual movement. Neither has to know how the other works inside, so a new source can be connected without touching the trigger, and the trigger can be tightened without touching a connector.

The trigger is a Dagster sensor that asks the source database one cheap question on an interval: what is the newest created_at timestamp in the table? If that value has not moved since the last look, the sensor returns without doing anything and records that it found nothing. If it has moved, it fires the acquisition job. The reason to build it this way is cost asymmetry. Running that query is nearly free and running a full sync is not, so making the expensive operation conditional on the cheap one means the pipeline can be checked often without being run often. That is what brings ingestion latency down without paying for it in wasted transfers.

When the job does fire, it calls Airbyte through Airbyte's API rather than shelling out or reimplementing the transfer, then polls the connection until it reports a finished status. That polling is deliberate and easy to leave out. An API call that starts a sync returns almost immediately, long before any data has landed, so a job that treats that response as success reports green while rows are still in flight. Polling to completion means the run's duration reflects how long the sync really took, and anything scheduled to wait on this job waits for data that is actually there.

Both outcomes are written to a log table alongside their metadata, and a failure raises rather than being swallowed. This matters more than it looks, because without it three very different situations look identical from outside: the sensor found nothing to do, the sync ran and moved rows, or the sync failed. Recording the quiet outcome as explicitly as the loud one turns an empty period into a fact that can be checked rather than an absence that has to be interpreted. Raising on failure instead of logging and carrying on hands the problem back to the orchestrator, so Dagster's own retry and alerting apply rather than the error settling into a log nobody reads.

The approach has edges worth naming. A watermark based on the newest timestamp sees inserts and nothing else, so updates and deletes in the source pass unnoticed unless they happen to touch that column. It also assumes rows arrive roughly in the order their timestamps suggest, since a record written late but dated below the current watermark is skipped, which is an accepted trade of watermark-based incremental loading rather than a defect in any one implementation. And the sensor's interval sets a floor on how fresh the destination can ever be.

Things to underline
  • Moved roughly 10,000 transaction records a day between PostgreSQL databases, scoped to one use case rather than the full business process
  • Split orchestration from transport, with Dagster deciding whether a sync runs and Airbyte performing it, so either side can change without touching the other
  • Gates an expensive sync behind a cheap newest-timestamp check, letting the pipeline be polled frequently without being run frequently
  • Polls the Airbyte connection to a finished status rather than trusting the API's immediate response, so a run reports success only once rows have actually landed
  • Writes every outcome to a log table with its metadata, so a period with no new data is a recorded event rather than an ambiguous silence
  • Raises on sync failure instead of swallowing it, keeping the orchestrator's retry and alerting behaviour in play
  • A newest-timestamp watermark only sees inserts, so updates and deletes at the source pass unnoticed unless they touch that column
  • A row written late but dated below the current watermark is skipped, an inherent trade of watermark-based incremental loading
  • The sensor's polling interval sets a floor on how fresh the destination can ever be
Built with
PythonDagsterAirbytePostgreSQL