Skip to content

Dataflows

Prerequisites · Reusable connections are available, or you are intentionally using an inline connection.
End state · A dataflow that describes one source-to-destination unit and can be validated before execution.

A dataflow is the unit DataCoolie reads, transforms, and writes. Its three runtime phases are siblings in the metadata model:

dataflows[]
├── source       read and filter input
├── transform    optionally shape the rows
└── destination  write rows with a load strategy

At runtime the phases run as Source -> Transform -> Destination. In the authoring process, choose the destination contract early enough to know whether the transform needs merge keys, SCD2 columns, partition expressions, or watermark columns. JSON property order does not change execution order.

Dataflow at a glance

{
  "dataflows": [
    {
      "name": "orders_to_bronze",
      "stage": "ingest",
      "source": {
        "connection_name": "orders_input",
        "table": "orders"
      },
      "transform": {
        "schema_hints": [
          { "column_name": "order_id", "data_type": "long" }
        ]
      },
      "destination": {
        "connection_name": "bronze",
        "table": "orders",
        "load_type": "append"
      }
    }
  ]
}

The Dataflow contract requires source and destination. transform is optional; when it is omitted, DataCoolie still applies the driver-managed pipeline behavior that belongs to every run.

Define the dataflow envelope

Give every dataflow a stable name or an explicit dataflow_id, plus a stage. Names are the normal authoring identity; add scheduling fields only when the run needs them.

Field Use it for
name Recommended human-readable identity for name-based authoring
dataflow_id Explicit stable identity for ID-based integrations; it can replace name
stage Select related dataflows with driver.run(stage=...)
description Business purpose or operational context
group_number Put dependent or related flows in the same scheduling group
execution_order Order buckets within a non-null group
processing_mode Select the configured processing mode; the normal built-in ETL path is batch-oriented
is_active Keep a definition in metadata while preventing its execution
configure Pass dataflow-level settings for supported extensions
{
  "name": "orders_to_silver",
  "description": "Load the current order state",
  "stage": "transform",
  "group_number": 1,
  "execution_order": 20,
  "processing_mode": "batch",
  "is_active": true,
  "source": { "connection_name": "bronze", "table": "orders" },
  "destination": {
    "connection_name": "silver",
    "table": "orders",
    "load_type": "merge_upsert",
    "merge_keys": ["order_id"]
  }
}

Use the same non-null group_number and increasing execution_order when a combined stage selection has an explicit dependency. Independent dataflows can omit both fields. See the run and orchestration guides for execution selection and failure behavior.

Keep name as the normal dataflow identity. See Dataflow identity when another system requires an explicit stable ID.

Dataflow identity

Use a unique name when the dataflow is authored or selected by name. DataCoolie derives dataflow_id from the name when an explicit ID is absent. An explicit ID may be used without a name for ID-based integrations; keep that ID unique and stable across deployments. The document mapper enforces dataflow ID uniqueness and permits distinct explicit IDs with the same display name, but any name-based consumer must still have one unambiguous match. When workspace_id is supplied, name scope is that workspace. Watermark and provider records use dataflow identity, so changing an explicit ID can select a different state record. For most documents, omit it and use the name. See the Dataflow contract.

For an externally managed state identity, add the ID to the named dataflow (dataflow fragment):

{
  "name": "orders_to_bronze",
  "dataflow_id": "flow-orders-bronze-v1",
  "source": {"connection_name": "orders_source", "table": "orders"},
  "destination": {
    "connection_name": "bronze", "table": "orders", "load_type": "append"
  }
}

Activation and selection

A dataflow can run only when its own is_active, its source connection's is_active, and its destination connection's is_active are all true. All three flags default to true. Normal metadata selection omits inactive dataflows; get_dataflows(active_only=False) returns them for inspection. Connections remain resolvable in the metadata provider even when inactive.

If a selected dataflow references an inactive connection, it finishes as SKIPPED with an activation reason. The same execution check applies when passing dataflows directly to driver.run(), replay or maintenance. A dataflow filtered out before execution has no run record. Disabling an upstream dataflow does not automatically disable downstream dataflows; check stage dependencies yourself. See Connection activation and the dataflow is_active field.

Configure the three phases

Source

The Source block chooses the read mode: table, query, or python_function. A query source may also set table as a human-readable logical label. A function source may use table and other Source fields as inputs to the custom function. Add watermarks and source filters when the source is incremental or needs a bounded read. A relative or artifact:/ SQL path still belongs in source.query.

Continue with Source for file, database, SQL file, API, function, and watermark cases.

Transform

Use Transform for casts, value rules, computed columns, deduplication, hashing, masking, projection, renaming, and row filters. Transform behavior may depend on the destination strategy: merge and SCD2 need key or audit columns, and partition expressions must survive the transform pipeline.

See Transform patterns for the execution order and feature-specific examples. Use Datatypes and schema hints when the source datatype or reusable schema metadata needs explicit control.

Destination

The Destination block chooses the target connection, table, load_type, and any write modifiers. merge_keys, partition_columns, SCD2 settings, and replace_by_watermark are conditional on the load strategy and source coverage.

See Destination & load patterns for append, overwrite, merge, SCD2, and window-replacement cases.

Add shared schema hints when a table is reused

Top-level schema_hints[] is an optional reusable layer defined by Shared schema hint. Match it by source connection_name, optional schema_name, and table_name; the provider can attach the matching hints[] items, shaped as Schema hint, to transform.schema_hints for a table source.

{
  "schema_hints": [
    {
      "connection_name": "warehouse",
      "schema_name": "sales",
      "table_name": "orders",
      "hints": [
        { "column_name": "order_id", "data_type": "long" },
        { "column_name": "created_at", "data_type": "timestamp" }
      ]
    }
  ]
}

This layer is shared by matching table sources; it is not an all-database datatype override. The source connection must allow schema hints, and inline dataflows[].transform.schema_hints takes precedence. Query and function sources can also match shared hints when source.table identifies their output. Without it, put their hints inline. See Datatypes and schema hints for matching and precedence.

If the source is incremental, choose the watermark and look-back together; see Incremental windows and look-back. If the destination replaces a watermark window, the source must cover the complete delete scope and the destination must use the matching load strategy.

Validate and run

Validate the complete metadata document after composing its connections and dataflows, and repeat the check after changing a source, transform, or load strategy. Use the Validation checklist before the first run. Keep SQL-file roots, secrets, merge keys, source coverage, and destination support in the same review.