Skip to content
AITroveRead. Build. Understand.
Make this comfortable

DAG intervals and idempotent task outputs

Last updated: 6 Oct 20265 min read
tutorial
AdvancedBy AITrove Editorial

A pipeline DAG declares dependencies among bounded tasks; each run must identify the exact input interval and produce the same committed result on retry.

Use interval identity, not wall-clock now

A daily billing task for event day 47 should read the declared day-47 source snapshot even when it is retried on day 52. Calling the current time inside the transformation changes its input and makes the retry a different calculation. Put interval start, interval end, input snapshot IDs and code version in the run context. A replay manifest makes those choices inspectable.

Give each task a publication boundary

A task may build a candidate output in temporary storage, but downstream tasks should consume only a validated committed artifact. Writing 70% of a partition and then failing must not leave a visible partial table. An upsert keyed by grain or an atomic partition replacement can make a retry converge. An append without a duplicate key usually cannot. Idempotent loads cover the database side of this rule.

Model dependencies at the data level

A transform should wait for the source interval and its quality checks, not merely for a clock tick. If source A is ready and source B is delayed, do not build a silently partial join. Mark the run blocked with the missing dataset and cutoff. Readiness checks distinguish absent data from empty but complete data.

Bound parallelism

A backfill of 240 daily intervals can overwhelm the warehouse if every interval starts at once. Limit active runs and expensive tasks, then prioritize recent consumer-critical intervals if that matches the recovery policy. Keep each interval independently rerunnable, but document dependencies when a cumulative model needs sequential checkpoints.

Test retry equivalence

Run the same interval twice from the same pinned inputs. Compare output row count, keys and a canonical checksum; the second run must not add duplicate bills or change amounts. Then inject a failure after writing candidate data and verify the published table still points to the previous complete result.

Implementation

python
from hashlib import sha256

def run_key(pipeline, interval_start, interval_end, input_snapshot):
    identity = f"{pipeline}|{interval_start}|{interval_end}|{input_snapshot}"
    return sha256(identity.encode()).hexdigest()[:16]

first = run_key("daily-billing", 47, 48, "snap-41")
retry = run_key("daily-billing", 47, 48, "snap-41")
different_input = run_key("daily-billing", 47, 48, "snap-42")
assert first == retry
assert first != different_input

Performance and operating cost

Run-key calculation is O(L) in the identity length. The expensive part is staging, validating, and publishing each interval. More parallel tasks reduce elapsed time until warehouse slots, source rate limits, or commit conflicts become the bottleneck.

Common Mistakes

  • Do not use the retry time as the event interval.
  • Do not let a failed task expose a partial partition.
  • Do not equate a scheduler success state with validated downstream data.

Read next

ai-data
data-engineering
Storage details