Schedule icon
Request icon
If icon
Fail icon
Read icon
Update icon
Insert icon
Log icon
Delete icon
SlackIncomingWebhook icon

Upsert a DocumentDB record from an external source

Sync one external record into DocumentDB with a serialized, race-safe emulated upsert and optional id-scoped retirement, with retry and Slack alerts.

Categories
Data

Keep a DocumentDB collection continuously in sync with an external HTTP source of truth, one record per run, using a race-safe emulated upsert. The DocumentDB plugin's Update task has no native upsert flag, so this blueprint reads by external_id, then conditionally updates or inserts. It solves the classic time-of-check / time-of-use race in read-then-write upserts by serializing executions, and it avoids the silent data loss of age-based collection prunes by making record retirement strictly opt-in and scoped to a single document. Use it for change-data-capture style mirroring, reference-data sync, and lightweight CDC into Amazon DocumentDB or any MongoDB-compatible Data REST API.

How it works

  1. fetch_source (io.kestra.plugin.core.http.Request) pulls the current record over HTTP with retry and options.allowFailed: true, so a 404 or 410 routes to the retirement branch instead of failing the run.
  2. route (io.kestra.plugin.core.flow.If) branches on outputs.fetch_source.code == 200.
  3. On a live record, guard_body (io.kestra.plugin.core.execution.Fail) validates the body has scalar id, name, and email before any write; read_existing (io.kestra.plugin.documentdb.Read, fetchType: FETCH_ONE) looks up the document by external_id.
  4. upsert (io.kestra.plugin.core.flow.If) chooses update_existing (io.kestra.plugin.documentdb.Update, updateMany: false, with $set and a $inc on sync_count) when the row exists, otherwise insert_new (io.kestra.plugin.documentdb.Insert).
  5. When the source returns 404 or 410 and retire_when_gone is true, delete_retired (io.kestra.plugin.documentdb.Delete, deleteMany: false) removes exactly that one document; otherwise the run logs and writes nothing.
  6. The hourly_sync Schedule trigger fires every hour, and concurrency.limit: 1 serializes runs so no two executions interleave.

What you get

  • A safe upsert that never double-inserts under concurrency.
  • Opt-in, id-scoped retirement instead of dangerous collection-wide deletes.
  • Source-body validation before any database write.
  • A Slack failure alert carrying the execution ID so a partial sync is never silent.
  • A sync_count counter and synced_at timestamp on every document.

Who it's for

  • Data engineers mirroring an external API into DocumentDB.
  • Platform teams running reference-data or master-data sync.
  • Anyone needing dependable CDC-style upserts without a heavy ETL stack.

Why orchestrate this with Kestra

The DocumentDB API has no scheduler, no retries, and no notion of a serialized read-then-write across runs. Kestra adds the event-driven Schedule trigger, per-task retry blocks, concurrency.limit: 1 to close the upsert race, full execution lineage and outputs for audit, and an errors block for alerting, all in declarative YAML you can version and review. Pair it with a unique index on external_id so even an external writer cannot create a duplicate.

Prerequisites

  • A reachable DocumentDB Data REST API endpoint with a catalog collection keyed by external_id (unique index recommended).
  • An external HTTP source returning JSON with scalar id, name, and email, and 404/410 when a record is retired.
  • Outbound network access from Kestra workers to both the source and DocumentDB.

Secrets

  • DOCUMENTDB_HOST: base HTTPS endpoint of the DocumentDB Data REST API.
  • DOCUMENTDB_USERNAME: basic-auth username for the DocumentDB API.
  • DOCUMENTDB_PASSWORD: basic-auth password for the DocumentDB API.
  • SLACK_WEBHOOK: Slack incoming webhook URL used by the failure alert.

Quick start

  1. Add the four secrets above to your Kestra instance.
  2. Point the source_url input at your record endpoint and set the database and collection variables.
  3. Run once to confirm the insert path, then run again to confirm the update path (collection size unchanged, sync_count incremented).
  4. Enable the hourly schedule, and turn on retire_when_gone only once you trust the source to return 404/410 for genuinely retired records.

How to extend

  • Swap the single-record fetch for a paged endpoint and loop with Loop to sync many records per run.
  • Replace the $set fields to mirror your real document shape.
  • Add a downstream notification or a metrics task keyed on outputs.upsert.evaluationResult.
  • Trigger on a webhook instead of a schedule for near-real-time sync.

Links

See How

New to Kestra?

Use blueprints to kickstart your first workflows.