DEV Community

LunarBreeze4173085
LunarBreeze4173085

Posted on

Nightly Pipeline Reconstruction — Serverless Polling Windows for Error Tracking API Timeouts

TL;DR: Treat a serverless timeout while polling an error-tracking API as a query-planning failure, not as a reason to raise the function timeout. For a nightly education-data pipeline, persist a closed time range, traverse it with deterministic pagination, and split any range whose query budget is nearly exhausted. Alert from durable run state after collection finishes; do not infer pipeline health from whether one long polling invocation returned.

The decision rule is strict: an invocation may stop, but an incident-reconstruction interval may never become ambiguous. That distinction matters when a 02:00 enrollment import fails and the on-call engineer must determine which tenant, course, and stage were affected. A larger search window appears to provide more context, yet it also increases scan work, response size, page count, and the chance that the collector dies between pages. More context can produce less evidence.

This architecture decision record chooses bounded, replayable retrieval over a single long query. It assumes an API that can filter by event time and paginate results, but it does not assume a particular vendor, SDK, or serverless runtime.

How should serverless polling handle error tracking API timeouts?

The first invariant is temporal closure. A poller records both ends of its target interval before sending the first request: [start, end). The end is fixed, rather than continually set to “now,” so late pages cannot pull newly arriving events into a result set whose earlier pages were calculated against a different population. Half-open intervals also give adjacent partitions one unambiguous boundary: an event at 03:00:00Z belongs to the interval beginning at that instant, not the one ending there.

The second invariant is restartability. A checkpoint contains the interval, the next pagination token, and the identity of the pipeline run being reconstructed. The token is opaque. Code must store and return it unchanged rather than parse it, increment it, or assume it remains useful outside the query that produced it. If an API documents token expiry, the recovery path restarts that bounded interval and relies on idempotent writes; it must not quietly advance the watermark.

No guessed cursors.

Third, persistence precedes acknowledgement. Normalize each event into a durable incident store using a stable source event identifier, then commit the next-page checkpoint. Reversing those operations creates a loss window: a crash after checkpointing but before storing the page makes the missing events invisible on retry. Reprocessing is acceptable when the destination enforces uniqueness. Loss is not.

The failure boundaries are now visible. Consider one interval assigned to a worker near the end of its runtime: the remote API may answer slowly, the page write may succeed, and the worker may then disappear before its continuation acknowledgement is observed. A retry receives the old checkpoint and writes that page again. With a unique key on (run_id, source_id), this is routine replay; without one, the incident view can show duplicate failures and inflate the apparent blast radius. The inverse ordering is worse. If the worker saves the next token first and disappears before the page write, the retry begins after evidence that was never stored. The runtime may also approach its deadline, a response may be lost, or events may arrive after their event-time interval has closed. None of these conditions is equivalent to “the nightly pipeline succeeded,” and none should erase collection progress. This is why I would accept duplicate processing and reject any checkpoint scheme that can create a silent gap.

Severity deserves similar skepticism. RFC 5424 defines syslog severity values, including Error, Critical, Alert, and Emergency, but an application label alone does not establish the impact of an edtech batch failure. A rejected optional analytics export and a failed enrollment import can both be logged as errors. The alerting decision therefore combines severity with pipeline stage, run identifier, terminal status, and affected tenant or dataset.

Decision and option boundary

The chosen design is a coordinator plus short-lived page workers. The coordinator creates immutable interval jobs and owns the high-water mark. A worker reads pages until either pagination ends or its remaining execution budget reaches a safety threshold, then durably schedules continuation state. If the first page of an interval is too expensive, the coordinator bisects that interval and retries the two child intervals. The minimum interval width is a declared limit; below it, the system records an explicit collection failure and alerts rather than splitting forever.

Option Incident reconstruction Main failure mode Valid boundary
One invocation, one large time range Simple only when the full scan completes Runtime ends after partial pagination, leaving an uncertain gap Small, predictably sparse datasets
Fixed time buckets Easy to replay and reason about A dense bucket can still exceed the deadline Stable event density with measured headroom
Adaptive interval splitting with checkpoints Preserves explicit coverage under uneven density More state transitions and duplicate processing on retry Bursty nightly workloads where evidence completeness matters
Continuous ingestion into owned storage Fast retrospective queries after ingestion Ingestion lag or outage becomes the evidence boundary Teams prepared to operate a durable ingestion path

Adaptive splitting wins here because nightly education jobs are uneven: one tenant may contribute a few validation messages while another imports a large roster. Fixed five-minute buckets merely move the guess into configuration.

Density wins the argument.

Still, adaptive splitting does not make the upstream API a durable archive. Retention, deletion, and delayed arrival remain external constraints, and the incident store should record the coverage it actually observed.

There is a cost trade-off. Smaller windows increase request and checkpoint overhead; larger windows risk wasted work when an invocation expires. Tune from page latency, page count, throttling responses, and remaining-runtime telemetry, not from a fashionable bucket size. The useful limit is the largest interval that completes with credible deadline headroom at a high but realistic event density.

The critical path in Python

The following code sketches the control path, not a vendor client. Its interface requires event-time bounds, an optional opaque token, and a bounded page size. The worker reserves 15 seconds for persistence, retry scheduling, and runtime shutdown; that number is an example policy, not a universal constant.

from dataclasses import dataclass
from datetime import datetime
from typing import Optional, Protocol, Sequence


@dataclass(frozen=True)
class Interval:
    run_id: str
    start: datetime
    end: datetime
    token: Optional[str] = None


@dataclass(frozen=True)
class Event:
    source_id: str
    occurred_at: datetime
    severity: int
    pipeline_stage: str
    tenant_id: str
    message: str


@dataclass(frozen=True)
class Page:
    events: Sequence[Event]
    next_token: Optional[str]


class ErrorSource(Protocol):
    def fetch(
        self,
        *,
        start: datetime,
        end: datetime,
        token: Optional[str],
        limit: int,
    ) -> Page: ...


class Store(Protocol):
    def put_events(self, run_id: str, events: Sequence[Event]) -> None: ...
    def continue_from(self, interval: Interval) -> None: ...
    def mark_complete(self, interval: Interval) -> None: ...


def collect_interval(
    interval: Interval,
    source: ErrorSource,
    store: Store,
    remaining_seconds,
    *,
    reserve_seconds: float = 15.0,
    page_size: int = 500,
) -> None:
    token = interval.token

    while remaining_seconds() > reserve_seconds:
        page = source.fetch(
            start=interval.start,
            end=interval.end,
            token=token,
            limit=page_size,
        )

        # put_events must be idempotent on (run_id, source_id).
        store.put_events(interval.run_id, page.events)
        token = page.next_token

        if token is None:
            store.mark_complete(interval)
            return

    store.continue_from(
        Interval(
            run_id=interval.run_id,
            start=interval.start,
            end=interval.end,
            token=token,
        )
    )
Enter fullscreen mode Exit fullscreen mode

One subtlety is deliberately absent: this loop does not update the run's global watermark. Only the coordinator may advance that watermark, and only after every interval in the closed range is marked complete. Otherwise, two workers finishing out of order can skip the unfinished gap between them.

Retries need classification. Rate limits and temporary upstream failures should use bounded exponential backoff with jitter, constrained by the remaining execution budget and any server-provided retry guidance. Authentication failures, invalid queries, and exhausted minimum-width intervals are terminal for that collection attempt. Retrying those until the runtime expires creates noise while concealing a configuration or contract error.

Stop early.

That short instruction prevents a common deadline failure: beginning one more remote request when too little time remains to store its result and checkpoint safely. The reserve should be derived from upper-tail persistence and scheduling latency, then tested under injected delay. A runtime's advertised maximum duration is a ceiling, not a processing budget.

Alert from durable run state, not a poller's exit code

Collection and alerting answer different questions. Collection asks, “Which events are durably available for this closed interval?” Alerting asks, “Does the evidence show that the nightly run failed or exceeded its completion deadline?” Joining them inside one invocation makes a transport timeout look like a pipeline incident and, worse, can suppress a real incident when the invocation disappears before sending the notification.

Those are separate alarms.

Model the run explicitly: expected, started, succeeded, failed, or overdue. Store stage transitions with the run identifier and event time. An alert evaluator reads that durable state after interval coverage is complete, or emits a separate “observability incomplete” alert when coverage cannot be completed. Those alerts must remain distinct. One concerns the education pipeline; the other concerns confidence in the evidence.

For incident reconstruction, the stored record should answer a bounded set of questions: which run failed, which stage last completed, which tenants or datasets were implicated, what source events support the conclusion, and which time intervals were fully searched. Preserve source timestamps and ingestion timestamps because late arrival makes their difference operationally meaningful. Do not place student content, email addresses, or unbounded payloads into alert text. Link the alert to access-controlled evidence by internal identifier.

Deletion requirements shape this design too. GDPR Article 17 describes a right to erasure and enumerates exceptions; it is not a blanket instruction to retain all logs because they may help debugging. Separate operational fields needed for reconstruction from personal data, define retention and erasure behavior with the responsible legal and security teams, and ensure derived incident records follow the same policy. Hashing an identifier is not automatically erasure if the organization can still relate it to a person.

Test the awkward transitions: a crash after writing events but before saving continuation; a duplicated page; an empty page with a continuation token if the API contract permits it; throttling near the deadline; child intervals completing out of order; and events arriving after a range closes. Deployment should begin with shadow collection that compares coverage and duplicate counts without paging anyone, followed by alerts routed to a test destination. These tests expose storage semantics that a successful happy-path query never touches.

Why reject a single scheduled scan?

A single scheduled scan is attractive because it has almost no control plane: compute the previous night's range, request errors, paginate, and alert. It is also a valid choice when measured worst-case volume fits comfortably inside the runtime, the upstream query has predictable latency, and losing one invocation merely delays a replay from an independently stored watermark.

Those conditions do not describe a reconstruction-critical pipeline with bursty tenants. Extending the function timeout changes the cliff but preserves it. Reducing the window without recording coverage creates a different problem: overlapping windows duplicate evidence, while adjacent windows can miss late events unless a deliberate lookback and idempotent destination are used. Pagination alone is not progress; durably recorded coverage is progress.

The rejected design can remain useful for low-volume operational summaries. It should not be allowed to advance an incident watermark after a partial response, and its alert should say when evidence is incomplete. That boundary turns a shortcut into an explicit engineering choice rather than an undocumented durability claim.

For the nightly pipeline, keep the architecture boring: immutable ranges, opaque tokens, idempotent event writes, deadline-aware continuation, and a coordinator that advances coverage only after all partitions finish.

Limits remain.

The design makes those limits visible, recoverable, and separate from the failure signal the on-call engineer is trying to trust.

References

Top comments (0)