Source¶
Prerequisites · A reusable connection, an inline
connection, and a dataflow envelope from Dataflows.
End state · A working connections entry and source block for any DataCoolie source type.
A source connection describes where DataCoolie reads data from. Two things work together: the Connection (shared, reusable endpoint definition) and the Source block inside each dataflow (table/query/function selector and incremental options).
connections[].configure
carries reusable settings; source.configure carries
source-specific settings. Precedence depends on the field: read_options merge
key by key (source wins on matching keys), while a source look-back replaces
the entire connection look-back. API settings have their own placement and
precedence rules; there is no universal configure merge. See
connection defaults
and look-back overrides.
For the read lifecycle, bounded ranges and reader-specific option routing, see Sources & destinations.
Source at a glance¶
This is a valid query-source fragment. Here, table is a short, human-readable
alias for the query result, not the table DataCoolie reads. A table source uses
table without query; a function source uses python_function and may also
set table. The framework calls the function named by python_function, while
the function may use any fields on the Source it receives, including table.
"source": {
"connection_name": "postgres_src",
"table": "active_orders",
"query": "SELECT updated_at, status, order_id FROM sales.orders",
"watermark_columns": ["updated_at"],
"filter_expression": "status = 'active'"
}
| Field | When you use it |
|---|---|
connection_name or connection |
Select a named reusable connection with connection_name (or a string connection reference), or supply an inline connection object; use only one of the two fields |
schema_name |
Table mode: folder / schema namespace. Query mode: optional metadata namespace, such as for shared schema-hint matching; not needed merely because SQL names a schema |
table |
File folder or lakehouse/database table; optional logical alias alongside query, or metadata that a Python function can use |
query |
Inline SQL or a relative/artifact .sql file; takes precedence over table when reading |
python_function |
Function-source mode instead of a physical table |
watermark_columns |
Incremental loading |
filter_expression |
SQL predicate evaluated against the reader result before transforms; query mode uses returned columns/aliases |
configure |
Source-specific settings such as read_options or API endpoint/params; precedence varies by key |
In practice, choose the read mode first:
tablefor direct file, lakehouse, or database table reads.queryfor database or lakehouse SQL reads when one executable SQL statement expresses the source, including joins, filters, and CTEs.python_functionwhen reading data needs more complex Python logic or a source case the built-in readers do not cover.
These are choices about how data is read, not strict complexity levels. A CTE
still belongs in query if the result can be produced by one SQL statement;
use python_function when SQL alone is not the right read boundary.
For a query source, optional table is a readable logical alias; the SQL
controls the read. For a Python function source, table is also optional, but
it may be a real input to the user function because the function receives the
whole Source object. The same applies to other Source fields: configure them
according to that function's contract. The framework still selects the
function by source.python_function.
An uncommon but supported case is attaching top-level schema_hints to a
query or function source. The provider matches those hints using the source
connection and source.table, plus source.schema_name when supplied. Only
use a matching group if its hints describe the resulting columns. Usually SQL
or custom function code already handles types; a few casts can be defined in
transform.schema_hints. See Datatypes and schema hints.
watermark_columns enables incremental loading. The saved watermark and filter
depend on the reader: SQL may push bounds into the query, while files and
lakehouse sources can filter in the engine. API request push-down is optional;
without it, the API reader can still filter fetched records locally. See
Incremental windows and look-back for
bounds and look-back.
Decide source type¶
| Source family | Connection shape | Typical source fields |
|---|---|---|
| File | connection_type: file, file format |
schema_name, table, watermark_columns, configure.read_options |
| Lakehouse | connection_type: lakehouse, delta or iceberg |
table or query, optional schema_name, watermark_columns |
| Database | connection_type: database, format: sql |
table or query, schema_name, watermark_columns |
| REST API | connection_type: api, format: api |
configure.endpoint, configure.params/body, pagination, watermark push-down; source.table is optional |
| Python function | connection_type: function, format: function |
python_function, optional table and other Source inputs, watermark_columns, custom source.configure |
File source (CSV, Parquet, JSON, JSONL, Avro, Excel)¶
Use when your data is in flat files on local disk, cloud object storage (S3, ADLS, GCS), or a Fabric/Databricks lakehouse path.
{
"name": "raw_files",
"connection_type": "file",
"format": "csv",
"configure": {
"base_path": "data/input",
"read_options": { "separator": ";" }
}
}
| Option | Required | Notes |
|---|---|---|
base_path |
yes | Root folder. For S3: s3://my-bucket/raw. For ADLS: abfss://container@account.dfs.core.windows.net/raw |
read_options |
no | Reader defaults for this connection. source.configure.read_options can override them per dataflow |
use_hive_partitioning |
no | Enables partition-folder discovery like country=VN/year=2026/ |
date_folder_partitions |
no | Date-folder pattern such as {year}/{month}/{day} for folder pruning |
backward_days / backward |
no | Re-read a historical window behind the last watermark |
use_schema_hint |
no | Defaults to true; disable to ignore schema_hints for this connection |
format |
yes | One of csv, parquet, json, jsonl, avro, excel |
Dataflow source block:
"source": {
"connection_name": "raw_files",
"schema_name": "sales",
"table": "orders",
"watermark_columns": ["__file_modification_time"]
}
Reads the folder at data/input/sales/orders. For incremental ingestion of
new or modified files, __file_modification_time is a common watermark: the
reader lists files and selects those newer than the saved modification time.
It is not an automatic default; omit watermark_columns for a full read, or
use a reliable row column such as updated_at when the requirement is to
filter records rather than files. A modified file is read again as a whole,
so the destination load strategy must handle reprocessed rows appropriately.
The normal lower bound is strict (>): a file arriving later with the same or
an older modification time can be missed. If this is possible, use a deliberate
look-back window and an idempotent
destination load.
When a source does not use __file_modification_time as a watermark, a file
whose platform listing has no modification time is still read; its
__file_modification_time lineage value is null. When the field is configured
as a watermark, every discovered file must provide a modification time. A
missing value fails the source before file reading or watermark advancement,
including on the first run. Fix the platform metadata or source listing and
retry from the unchanged watermark.
Per-dataflow file overrides¶
If one dataset needs special reader settings, override the connection-level
defaults in source.configure:
"source": {
"connection_name": "raw_files",
"schema_name": "sales",
"table": "orders",
"configure": {
"read_options": {
"separator": "|",
"encoding": "utf8-lossy"
}
}
}
File-source incremental cases¶
File readers support more than simple row-level watermarks:
| Pattern | How it works |
|---|---|
| Row-column watermark | Use a trustworthy column in the file data; the reader filters rows after reading |
| File modification watermark | Use __file_modification_time in watermark_columns to select new/modified files before reading; requires file listing with modification times |
| Date-folder pruning | date_folder_partitions lets the reader prune old folders before reading |
| Historical replay window | backward_days / backward re-opens a look-back window |
| Independent bounded read | Replay may use chunk_column="__file_modification_time" even when mtime is not a persisted watermark; the reader still selects files by actual mtime |
For a partitioned dataset under folders such as
data/input/sales/orders/year=2026/month=09/day=23/, configure both options
on its connection:
{
"name": "partitioned_files",
"connection_type": "file",
"format": "parquet",
"configure": {
"base_path": "data/input",
"use_hive_partitioning": true,
"date_folder_partitions": "year={year}/month={month}/day={day}"
}
}
Then use a source with connection_name: "partitioned_files",
schema_name: "sales", and table: "orders". If you also use
__file_modification_time, folder pruning happens first: a file modified in
an older, pruned date folder will not be discovered unless the folder window
is reopened.
File reads also inject file lineage columns such as __file_name,
__file_path, and __file_modification_time.
Internal folder watermark
When you use date_folder_partitions, DataCoolie stores the folder-level
watermark internally as __date_folder_partition__. You normally do not
need to author that field yourself; the reader maintains it. It is used
only for conservative folder discovery. Folder boundaries are inclusive;
the file's modification time decides whether the file is read. The
internal folder key cannot be a replay chunk_column. The reader also
reloads this internal frontier on later ordinary runs when the source has
no authored watermark_columns, so folders older than the saved frontier
are not scanned again.
Excel is read-only
Excel (format: "excel") can only be a source. You cannot write Excel
as a destination.
Lakehouse source (Delta or Iceberg)¶
Use when reading from a Delta table on a lakehouse (Databricks, Fabric, local Delta folder, or S3 Delta Lake), or an Iceberg table registered in a catalog.
{
"name": "bronze_lake",
"connection_type": "lakehouse",
"format": "delta",
"configure": {
"base_path": "data/output/bronze"
}
}
Dataflow source block:
"source": {
"connection_name": "bronze_lake",
"schema_name": "sales",
"table": "orders",
"watermark_columns": ["updated_at"]
}
Choose path or catalog addressing based on the table format, engine, and
governance rules—not simply on whether the platform is Databricks or Fabric.
For a named Delta table such as a Databricks Unity Catalog managed table,
configure its catalog/database scope instead of a storage base_path:
{
"name": "unity_bronze",
"connection_type": "lakehouse",
"format": "delta",
"catalog": "my_catalog",
"database": "bronze",
"configure": {}
}
Table is then referenced by name: schema_name maps to the schema, table to
the table. Leave schema_name empty when the Unity Catalog layout only uses
three parts (catalog.database.table).
For Iceberg, prefer catalog addressing for an end-to-end read/write pipeline:
{
"name": "iceberg_lake",
"connection_type": "lakehouse",
"format": "iceberg",
"catalog": "glue_catalog",
"database": "raw",
"configure": {}
}
| Addressing mode | What to configure |
|---|---|
| Delta by path | configure.base_path, table, optional schema_name; a good default for local/cloud Delta and Fabric when direct path access is appropriate |
| Delta by registered name | catalog, database, table, optional schema_name; use for governed named tables, including Unity Catalog managed tables |
| Iceberg by catalog | catalog, database, table, optional schema_name; preferred for portable reads and writes |
Fabric Delta can use either a table name or a OneLake path when the runtime and access policy permit it. Databricks Unity Catalog managed tables must be accessed by name; external tables have different path-access rules. With OneLake security enabled on a Fabric table, direct path access can be blocked for non-privileged users, so use the named table. See the platform guidance for Unity Catalog paths and Fabric OneLake security.
Engine support also matters: Polars Delta uses paths, not named Delta tables. For Iceberg, use catalog addressing for a full read/write pipeline: Polars can read an Iceberg path at the engine level but requires a catalog table name for writes. Path-read fallbacks should not be treated as a portable Iceberg write contract.
Delta and Iceberg readers apply watermark bounds through the engine's DataFrame filter after the table read.
Lakehouse SQL query¶
Delta and Iceberg sources can also use source.query when one SQL statement
expresses the read, including a WITH CTE. For example:
"source": {
"connection_name": "unity_bronze",
"table": "recent_orders",
"query": "WITH recent AS (SELECT * FROM my_catalog.bronze.orders WHERE status = 'OPEN') SELECT * FROM recent"
}
Here table is an optional readable alias; the SQL determines which relations
are read. Lakehouse readers execute the query through the selected engine and
apply watermark bounds to the resulting DataFrame. SQL-file references use
the same query-file resolution rules described below.
With Polars, the SQL relations must be registered in the engine before the
query runs (for example, Delta or Iceberg table discovery in the runner);
metadata alone does not register them. Install the SQL resolver dependency and
follow Qualified SQL relations in Polars.
Database source (SQL via SQLAlchemy)¶
Use when reading from PostgreSQL, MySQL, MSSQL, Oracle, SQLite, or any SQLAlchemy-supported database.
{
"name": "postgres_src",
"connection_type": "database",
"format": "sql",
"configure": {
"database_type": "postgresql",
"host": "warehouse.internal",
"port": 5432,
"database": "analytics",
"username": "DC_DB_USER",
"password": "DC_DB_PASSWORD"
},
"secrets_ref": {
"env:": ["username", "password"]
}
}
Never hardcode passwords
Use secrets_ref instead of putting credentials in configure. See
Concepts · Secrets · secrets_ref schema and the credential section
below.
You can connect with either of these shapes:
| Pattern | When to use |
|---|---|
configure.url |
You already have one connection string |
database_type + host + port + database (+ credentials) |
You want DataCoolie to assemble database options more explicitly |
Resolving a full URL from environment variables:
{
"name": "postgres_src",
"connection_type": "database",
"format": "sql",
"configure": {
"url": "DC_POSTGRES_URL"
},
"secrets_ref": {
"env:": ["url"]
}
}
Set DC_POSTGRES_URL=postgresql+psycopg2://realuser:realpass@host:5432/mydb
in your environment. DataCoolie replaces configure.url at runtime.
Database authentication types¶
By default, database connections use username/password auth. The optional
auth_type
field enables alternative authentication methods:
auth_type |
Required fields | Use case |
|---|---|---|
password (default) |
username, password, unless the connection URL supplies credentials |
All databases — standard SQL auth |
service_principal |
username (= client ID), password (= client secret), tenant_id |
Azure SQL, Fabric SQL via Azure AD/Entra |
managed_identity |
none (or username = client ID for user-assigned MI) |
Azure-hosted runtimes (AKS, App Service, Fabric) |
access_token |
token (+ optional username for non-MSSQL) |
Pre-fetched token from any provider (Azure, AWS IAM, GCP) |
Service principal example (Azure SQL):
{
"name": "azure_sql_spn",
"connection_type": "database",
"format": "sql",
"configure": {
"database_type": "mssql",
"auth_type": "service_principal",
"host": "myserver.database.windows.net",
"port": 1433,
"database": "mydb",
"username": "AZURE_CLIENT_ID",
"password": "AZURE_CLIENT_SECRET",
"tenant_id": "AZURE_TENANT_ID"
},
"secrets_ref": { "env:": ["username", "password", "tenant_id"] }
}
Managed identity example (zero-credential, Fabric):
{
"name": "fabric_sql_mi",
"connection_type": "database",
"format": "sql",
"configure": {
"database_type": "mssql",
"auth_type": "managed_identity",
"host": "xyz.datawarehouse.fabric.microsoft.com",
"port": 1433,
"database": "mydb"
}
}
Pre-fetched access token example (AWS RDS IAM):
{
"name": "rds_postgres_iam",
"connection_type": "database",
"format": "sql",
"configure": {
"database_type": "postgresql",
"auth_type": "access_token",
"host": "mydb.xxx.us-east-1.rds.amazonaws.com",
"port": 5432,
"database": "analytics",
"username": "iam_db_user",
"token": "RDS_IAM_TOKEN"
},
"secrets_ref": { "env:": ["token"] }
}
Fabric SQL endpoint
Fabric SQL endpoints (*.datawarehouse.fabric.microsoft.com) only accept
Entra ID auth. DataCoolie rejects auth_type: "password" for these hosts
at validation time.
Engine notes
Spark (JDBC): SPN/MI auth uses native MSSQL JDBC driver properties —
no extra Python packages needed.
Polars: Non-password MSSQL auth routes through pyodbc + ODBC Driver
18 instead of connectorx. Other databases use token-as-password via
connectorx.
Database transport options¶
The database reader passes unhandled connection configuration keys through to the selected database transport. This is useful for provider-specific options such as MSSQL TLS settings:
{
"configure": {
"database_type": "mssql",
"url": "MSSQL_DATABASE_URL",
"encrypt": "yes",
"trustServerCertificate": "false"
},
"secrets_ref": {"env:": ["url"]}
}
Pass-through behavior depends on the selected engine, driver, and installed database capability. Use the exact spelling expected by that transport and do not assume a SQLAlchemy/ODBC/JDBC option works on every engine. The generated Connection reference lists DataCoolie-defined keys; provider-specific keys remain an open configure map.
Polars database result typing¶
The Polars reader preserves the result values from the selected source driver
before schema_hints are applied. Install
datacoolie[source-db-native-polars] for the default precision-preserving
MySQL and MSSQL paths, and datacoolie[source-db-oracle-polars] for Oracle.
These paths keep unsigned integers, DECIMAL/NUMBER, and temporal values in
their native Python representation so a source-aware hint can make the
intentional engine-owned cast.
configure.database_read_engine is an explicit transport override:
Resolve DC_MYSQL_URL through secrets_ref on the surrounding connection,
as in the database URL example; do not
put credentials in the metadata file.
"native" is the default for MySQL, MSSQL and Oracle. Set
"connectorx" only when that transport is required and its result typing is
acceptable; ConnectorX may project unsigned or high-precision numeric values
to floating point, and a later cast cannot restore precision that was already
lost. PostgreSQL and SQLite continue to use their existing URI readers unless
an explicit source-specific reader is selected.
The native MySQL/MSSQL readers apply only the connection coordinates and the documented read-size options. Unsupported URI query parameters (for example driver-specific TLS flags) fail explicitly instead of being ignored; select ConnectorX or a dedicated ODBC configuration when those transport options are required.
Dataflow source block:
"source": {
"connection_name": "postgres_src",
"schema_name": "public",
"table": "orders",
"watermark_columns": ["updated_at"]
}
schema_name maps to the SQL schema, table to the SQL table name.
Inline SQL or SQL-file query sources¶
When the source is a SQL query, put the executable SQL in source.query.
Optionally add source.table as a concise logical label:
"source": {
"connection_name": "postgres_src",
"table": "open_orders",
"query": "SELECT * FROM sales.orders WHERE status = 'OPEN'",
"watermark_columns": ["updated_at"]
}
For database sources, if watermark_columns are set, DataCoolie wraps the
query as a subquery and applies the watermark filter outside it. Lakehouse
readers filter the resulting DataFrame instead.
The same query field accepts a conservative file shorthand during Driver
preparation. A single relative token ending in .sql is read below the
provider's sql_base_path when the metadata provider declares SQL roots, or
below the Driver's SQL roots when it supplies the session fallback. Otherwise
the complete relative path is read directly below artifact_base_path; no
sql/ folder is assumed:
Read a SQL file¶
orders/incremental.sql is also valid and is not prefixed with sql/.
Use artifact:/sql/orders.sql when the artifact root must be selected
explicitly or the filename does not match the shorthand. File resolution is a
Driver preparation step; direct reader calls still require executable SQL. The
path is metadata, while the resolved SQL text is the value passed to the
database reader. The provider keeps the path configuration but does not need
the Driver platform to store it.
For multiple SQL roots, pass a sequence of roots and qualify the reference by
the final folder name of the selected root (for example, sql1/orders.sql or
queries/customers.sql). Prefixes must be unique; a single configured root
also accepts the root-relative form without its folder prefix.
The query file still belongs in source.query; an optional source.table
is a logical label (and possible shared-hint lookup key), not the SQL-file
selector. It does not affect which SQL file is resolved or executed. Validate
the resolved path with the same artifact and SQL-root arguments used by the
runner; a path that works only from the repository checkout is not a portable
metadata contract.
REST API source¶
Use when reading JSON records from a REST API, with or without pagination.
For request-shape fields, see Connection settings by endpoint type
and source.configure
in the Metadata reference.
{
"name": "orders_api",
"connection_type": "api",
"format": "api",
"configure": {
"base_url": "https://api.example.com/v1",
"auth_type": "bearer",
"auth_token": "DC_ORDERS_API_TOKEN",
"timeout": 30,
"default_headers": { "Accept": "application/json" }
},
"secrets_ref": {
"env:": ["auth_token"]
}
}
Put request shape and pagination on the source, not on the connection:
"source": {
"connection_name": "orders_api",
"table": "orders",
"watermark_columns": ["updated_at"],
"configure": {
"endpoint": "/orders",
"method": "GET",
"params": { "status": "open" },
"pagination_type": "offset",
"page_size": 200,
"max_pages": 1000,
"total_path": "meta.total",
"data_path": "data.items"
}
}
This example expects a response shaped like
{"meta":{"total":2},"data":{"items":[{"id":1,"updated_at":"2026-01-01T00:00:00Z"},{"id":2,"updated_at":"2026-01-02T00:00:00Z"}]}}.
Check data_path against a real response: a missing path produces zero records,
not a path-validation error. Omit pagination_type and total_path for a
single-response endpoint. watermark_columns alone filters rows after
fetching; it does not make the API return only changed records. Use
watermark push-down when the endpoint accepts
incremental parameters. The API source also applies source.filter_expression
locally to the returned DataFrame before transforms; it does not add the
predicate to request parameters or body.
If pagination_type is present, it must be offset, cursor, or next_link.
An omitted or null value means one response. A typo or unsupported value fails
before credentials are resolved or an API request is sent, so an incremental
run cannot advance its watermark without a known completion contract.
With total_path, offset pages are fetched concurrently (four workers by
default), and rate_limit_delay is not applied. max_pages is a safety
budget: if the response still requires another page at that limit, the read
fails instead of returning a partial result. Match page_size, max_pages, and
offset_max_workers to the provider's contract and rate limit. For sequential
offset requests, omit total_path; then rate_limit_delay applies between
pages. See pagination contracts.
| Connection-level key | Purpose |
|---|---|
base_url |
Required root URL |
auth_type |
bearer, basic, api_key, oauth2_client_credentials, aws_sigv4 |
auth_token / username / password / api_key_* |
Auth credentials |
token_url, client_id, client_secret |
OAuth2 client-credentials flow |
default_headers, timeout |
Shared request defaults |
| Source-level key | Purpose |
|---|---|
endpoint |
Path appended to base_url |
method, params, body |
Request shape |
pagination_type |
offset, cursor, or next_link |
page_size, max_pages, data_path |
Response traversal and size |
total_path, offset_max_workers |
Parallel offset pagination |
rate_limit_delay, max_retries |
Delay between sequential pages; retry HTTP 429 responses |
source.filter_expression is deliberately absent from the request table. Put
server-side filters in the endpoint-specific params or body when the API
supports them, and use the source predicate for a framework-side safety filter
over the returned columns.
The table above is the basic API shape. Use Connections for authentication and the API source configuration section below for all request, pagination, timezone, and split-range variants. For a combined example see Incremental API with pagination.
For API incremental loading, choose between local filtering, request watermark push-down, and split ranges. The complete field combinations, timezone handling, and inclusive-boundary rules are documented in API source configuration, especially Push down and split watermark ranges.
API source configuration¶
This section owns the request, response, pagination, and API watermark fields
that belong to source.configure. Connections
owns the base URL, authentication, and secrets. The exact source field types
and allowed values are in the Source reference
and source.configure.
An authenticated source can use any supported pagination and watermark mode. The examples below are source fragments, not complete metadata documents.
| Decision | Choose | Where to configure |
|---|---|---|
| Response | Root JSON list, nested list, or one nested object | source.configure.data_path when not a root list |
| Pagination | One response, record offset, cursor, or next-link | source.configure.pagination_type and its matching path/parameter keys |
| Incremental read | Fetch then filter, API push-down, or bounded split ranges | source.watermark_columns plus the applicable source.configure keys |
Match the request and response shape¶
source.configure.method is an HTTP method string (GET by default), not a
DataCoolie enum. params becomes URL query parameters and body becomes a
JSON request body. For example, a POST endpoint that returns one JSON object
under result can use:
{
"source": {
"connection_name": "orders_api_oauth",
"table": "order_summary",
"configure": {
"endpoint": "/reports/orders",
"method": "POST",
"params": {"region": "west"},
"body": {"status": "open"},
"data_path": "result"
}
}
}
The response {"result":{"count":12}} becomes one record. With no
data_path, the response must be a root-level JSON list. A wrong path yields
zero records. These keys belong to the source because different endpoints on
one connection may use different request and response shapes.
Choose the pagination contract¶
Pagination settings belong to source.configure because different endpoints
on the same API may return different response shapes. For every paginated
mode, choose max_pages large enough for the intended result. Reaching the cap
while another page is required raises a source error instead of returning
partial data. Use a provider-defined stable sort or snapshot when available,
especially for concurrent offset requests against changing data.
One response (no pagination)¶
Omit pagination_type for an endpoint that returns all records in one
response. For a root-level list, omit data_path too; for a nested list, set
data_path to its dot-separated path. The reader makes one request and does
not follow a cursor, link, or offset unless a pagination mode is configured.
Offset pagination with provider-specific parameter names¶
{
"source": {
"connection_name": "orders_api_key",
"table": "orders",
"configure": {
"endpoint": "/orders",
"pagination_type": "offset",
"page_size": 200,
"offset_param": "skip",
"limit_param": "per_page",
"total_path": "meta.total",
"data_path": "data.items",
"offset_max_workers": 4
}
}
}
With total_path, the first response provides the total and remaining pages
can be fetched concurrently. Without it, offset pages are fetched
sequentially until a short page or empty response is returned. The reader
sends record offsets (0, 200, 400, ... in this example), not page
numbers. A provider whose page parameter expects 1, 2, 3, ... does not
fit this built-in offset mode; do not merely rename offset_param to page.
The first response above must contain a numeric meta.total and an array at
data.items. A wrong data_path is treated as an empty result, so test it
against a real response. If total_path is missing or non-numeric, the
concurrent path fails. Setting total_path also enables up to four concurrent
offset workers by default; rate_limit_delay only delays sequential page
requests. max_pages defaults to 1,000 and is a safety budget: the reader
fails when the declared total needs more pages than that budget or when the
aggregate fetched count does not equal the declared total. max_retries
retries HTTP 429 responses, not arbitrary HTTP failures.
Cursor pagination¶
{
"source": {
"connection_name": "orders_api_oauth",
"table": "orders",
"configure": {
"endpoint": "/orders",
"pagination_type": "cursor",
"page_size": 100,
"cursor_path": "paging.next_cursor",
"cursor_param": "after",
"data_path": "results"
}
}
}
The reader sends the returned cursor under cursor_param. The default response
path is next_cursor and the default request parameter is cursor. The
authored request options are grouped under
source.configure.
Next-link pagination¶
{
"source": {
"connection_name": "orders_api_basic",
"table": "orders",
"configure": {
"endpoint": "/orders",
"pagination_type": "next_link",
"next_link_path": "paging.next",
"data_path": "data"
}
}
}
When the response contains an absolute next-link URL, the reader follows it
verbatim and uses the link's parameters. The continuation token is opaque: the
reader does not decode it or append the original query bounds by default. Set
next_link_bound_mode: "repeat_query_bounds" only for an endpoint whose
contract explicitly requires the original range parameters on every
continuation; matching values are checked and conflicting/duplicate values
fail before the request. It stops when the configured path is empty.
Relative links are resolved against the current page URL. A continuation must
stay on the configured HTTP(S) origin (scheme, host, and effective port); a
foreign host or port, scheme downgrade, userinfo, or non-HTTP(S) link fails
before the next request. Connection credentials therefore remain scoped to the
configured origin. Redirect following is disabled for the API client, so a
redirect cannot bypass the continuation check. Providers that require
cross-origin continuation need an explicit credential-scoping policy outside
this built-in mode.
Push down and split watermark ranges¶
For a bounded replay or another source-owned range, use the explicit
range_param_mapping contract. Each field declares its lower and upper API
parameter, the wire format, the response field used for residual filtering,
and how a successful read advances state:
{
"source": {
"connection_name": "orders_api",
"table": "orders",
"watermark_columns": ["updated_at"],
"configure": {
"endpoint": "/orders",
"range_param_mapping": {
"updated_at": {
"lower": {"name": "created_from", "operator": ">="},
"upper": {"name": "created_before", "operator": "<"},
"format": "iso",
"response_column": "updated_at",
"watermark_value": "observed_max"
}
},
"pagination_type": "next_link",
"next_link_path": "paging.next"
}
}
}
watermark_value: "observed_max" saves the maximum response value after a
successful read. watermark_value: "request_end" saves the exclusive upper
bound that the endpoint confirmed. Omitted incremental lower operators resolve
to > for observed_max and >= for request_end; active fields with mixed
semantics are rejected before HTTP. A bounded replay uses [start, end) and
requires endpoint operators and a response_column that can enforce any
boundary the endpoint cannot enforce itself. Legacy watermark_param_mapping
does not provide this independent bounded-range contract.
A tracked observed_max binding must have a non-null response_column; the
reader rejects a missing column mapping before HTTP, even when watermark saving
is disabled. A request_end binding can omit it when the endpoint enforces the
requested operators exactly and guarantees complete interval pagination.
New bindings encode explicit bounds without rounding. iso retains the
datetime offset and microseconds; nonzero ISO fractional digits beyond six
cannot be represented and are rejected. ISO timezone offsets that Python
cannot parse without precision loss are also rejected. date requires midnight, datetime
requires whole seconds, and datetime_ms / timestamp_ms require millisecond
alignment. timestamp preserves fractional Unix seconds using exact decimal
encoding. For example, 12:30:00.123456 cannot use datetime_ms; use iso or
supply a bound that is already aligned. Validation happens before HTTP for
both incremental and bounded reads. Legacy watermark_param_mapping retains
its existing formatting behavior.
When the reader generates an incremental upper bound from the current time,
it chooses a ceiling compatible with the binding's precision and logical
date/datetime type before creating requests. Split requests and request_end
state use that same covered ceiling. Explicit bounds and stored lower values
are never rounded to make them fit.
The mapped field does not have to be one of source.watermark_columns when it
is used only for bounded selection. For example, set chunk_column="created_at"
and define range_param_mapping.created_at while keeping
watermark_columns: ["updated_at"]. Replay sends the created_at bounds and
the API reader continues to calculate/persist state only for the configured
watermark columns. A maximum observed from that slice must not be treated as
proof that the complete updated_at history was read.
Set a binding's location to params for query parameters or body for a
top-level JSON body field. Nested body paths are outside this generic mapping
contract and require a source-specific adapter.
Use watermark_param_mapping to map stored watermark columns to API parameter
names. Add watermark_to_param when the API accepts an upper bound:
{
"source": {
"connection_name": "orders_api_oauth",
"watermark_columns": ["updated_at"],
"configure": {
"endpoint": "/orders",
"watermark_param_mapping": {"updated_at": "updated_since"},
"watermark_to_param": "updated_before",
"watermark_param_location": "params",
"watermark_param_format": "iso",
"watermark_to_param_timezone": "Asia/Ho_Chi_Minh",
"watermark_range_interval_unit": "day",
"watermark_range_interval_amount": 1,
"watermark_range_start": "2026-01-01T00:00:00Z",
"watermark_range_max_workers": 4,
"watermark_range_to_exclusive_offset": "1ms"
}
}
}
The legacy incremental split requires
watermark_range_interval_unit and watermark_to_param. On the first run it
also requires watermark_range_start; later runs can start from the saved
watermark when watermark_param_mapping is present. The API must accept both
mapped lower and configured upper bounds; omitting the mapping causes every
range request to lack a lower bound. Supported interval units are hour,
day, month, and year.
Canonical bounded reads and replay use range_param_mapping instead. Each
selected field declares both endpoint bindings and their operators, so a
canonical request can enforce the exact [start, end) interval without
watermark_to_param, watermark_range_start, or legacy interval fields. A
canonical numeric field is read as one finite range; use replay's integer
chunk_interval: {"step": ...} when several numeric chunks are needed.
The legacy upper-bound behavior is unchanged. Set
watermark_range_to_exclusive_offset to 1ms, 1s, or 1day only when that
API treats its upper bound as inclusive, for example a BETWEEN from AND to
query. The offset changes the value sent to the API; adjacent internal range
boundaries remain contiguous. This adjustment does not replace the explicit
operators on a canonical range_param_mapping binding.
watermark_to_param_timezone accepts an IANA name such as
Asia/Ho_Chi_Minh or a UTC offset such as +07:00. A source-level value wins
over the same connection-level default. The range endpoint and the stored
watermark must use a precision and timezone that the API understands.
watermark_param_format controls how a stored date/datetime becomes the upper
bound request parameter:
| Format | Example output | Behavior |
|---|---|---|
iso |
2026-01-02T03:04:05+00:00 |
ISO-8601 output; preserves the datetime's timezone state |
date |
2026-01-02 |
Calendar date only |
timestamp |
1767323045.0 |
Unix seconds; naive datetimes are treated as UTC |
timestamp_ms |
1767323045000 |
Unix milliseconds; naive datetimes are treated as UTC |
datetime |
2026-01-02T03:04:05 |
Naive ISO datetime truncated to whole seconds |
datetime_ms |
2026-01-02T03:04:05.123 |
Naive ISO datetime retaining millisecond precision |
Unparseable string watermarks pass through unchanged; typed date/datetime values use the conversion rules above. Choose a format and timezone accepted by the API rather than relying on a provider-specific default.
Avoid invalid combinations¶
- Do not put pagination keys on the connection; they belong to the source.
- Do not enable the legacy incremental split without
watermark_to_paramand a first-runwatermark_range_start. Canonical bounded reads and replay userange_param_mappinginstead. - Do not use
cursor_paramornext_link_pathwith the wrongpagination_type. - Do not use a custom or misspelled
pagination_type; onlyoffset,cursor, andnext_linkare supported, and an omitted value means one response. - Put auth settings and secrets on the connection.
- Do not let the API response omit a pushed-down watermark column without
understanding the reader's
nowadvancement behavior.
table is optional for APIs
The API reader uses connection.configure.base_url plus
source.configure.endpoint. source.table is best treated as a logical
label for your metadata or logging, not as the actual HTTP path. Keep a
stable table when shared schema hints
must match this source by connection/table identity.
Python function source¶
Use when reading data needs custom Python logic—for example, combining several
steps or using an SDK—or when the built-in table/query readers do not cover the
source. If one database or lakehouse SQL statement (including CTEs) is enough,
prefer source.query.
The actual function path goes on source.python_function:
Dataflow source block:
"source": {
"connection_name": "custom_src",
"table": "partner_orders",
"python_function": "mypkg.sources.load_orders",
"watermark_columns": ["updated_at"],
"configure": {
"api_base": "https://partner.example.com"
}
}
Define load_orders(engine, source, watermark_start, watermark_end) for
ordinary incremental reads (the framework passes all four as keyword
arguments). A bounded replay passes a separate read_range keyword and sets
watermark_start=None and watermark_end=None; a bounded function must accept
that keyword explicitly or through **kwargs. read_range carries the exact
source-owned column, start/end values, and comparison operators. Use it for
push-down when possible, and return rows that honor the range; the built-in
function reader also applies the exact residual range filter before it observes
the candidate watermark. It rejects a legacy function that cannot accept the
keyword. Return a DataFrame compatible with the active engine, or None when
there is no data. The function receives the full Source object, so it can use
source.table, source.schema_name, source.configure, and other fields as
meaningful inputs. In this example table may name the dataset that
load_orders reads; it is not necessarily just a display alias.
The framework uses source.python_function to choose the function.
To narrow which dotted function paths may be imported from metadata, pass
allowed_function_prefixes in DataCoolieRunConfig. This is a string-prefix
check, not a sandbox; only run trusted metadata and choose prefixes carefully.
Incremental reads with watermark_columns¶
For any source type, add watermark_columns to the source block to enable
incremental loading:
"source": {
"connection_name": "postgres_src",
"schema_name": "public",
"table": "orders",
"watermark_columns": ["updated_at"]
}
The first run has no stored lower bound, unless a configured initial range or
replay window supplies one. Later reads use the stored watermark and any
configured look-back. A database reader adds a SQL watermark condition; other
readers use different mechanisms. Checkpoint advancement also differs: most
readers use the maximum observed watermark, while canonical API bindings use
their configured watermark_value (observed_max or request_end) and the
legacy split path advances from its covered request boundary.
Behavior differs by source family:
| Source family | Watermark behavior |
|---|---|
| Parquet / Delta / Iceberg | Engine DataFrame filter after read |
| Database | WHERE clause pushed to SQL; query mode filters its returned columns |
| API | Mapped watermark fields are pushed into request params/body; unmapped fields are filtered in the engine after fetching |
| File formats | __file_modification_time can select files during listing; row-column watermarks filter in the engine after read |
| Python function | Incremental watermarks are passed into the function first; bounded read_range is a separate source-owned contract, then the framework filters again after the function returns |
API request mapping can reduce remote fetches; without it, local filtering does not prevent a full API fetch. See API watermark push-down and Concepts · Watermarks · Storage ownership and path binding for watermark storage and provider interaction.
Incremental windows and look-back¶
An incremental source normally starts at the saved watermark. A look-back
reopens part of that range so late-arriving or corrected records can be read
again. The runtime-computed property is called date_backward; do not author
date_backward in metadata. Author one of the supported backward_* fields or
the nested backward object in source.configure.
Choose a fixed look-back¶
Put a reusable default in connections[].configure:
{
"connections": [
{
"name": "orders_source",
"connection_type": "database",
"format": "sql",
"configure": {
"database_type": "postgresql",
"url": "ORDERS_DATABASE_URL",
"backward_days": 3
},
"secrets_ref": {"env:": ["url"]}
}
]
}
The supported top-level shorthand fields are:
| Authored field | Effective meaning |
|---|---|
backward_hours |
Subtract hours from the saved datetime watermark |
backward_days |
Subtract days from the saved watermark |
backward_months |
Subtract calendar months, with calendar-safe clamping |
backward_years |
Subtract calendar years, with calendar-safe clamping |
backward_closing_day |
Use the closing-day strategy with that day of the month |
You can express the same values in one nested object:
{
"source": {
"connection_name": "orders_source",
"table": "orders",
"configure": {
"backward": {
"days": 3,
"hours": 6
}
}
}
}
The nested keys are years, months, days, hours, and closing_day.
Use one clear strategy per connection unless you intentionally need combined
fixed offsets such as three days plus six hours. If the nested object repeats
a unit also supplied by a shorthand field, the nested value wins during
parsing.
On the first run, there is no saved watermark to adjust, so the look-back has
no effect until a saved watermark exists. Non-datetime watermark values pass through
unchanged. When a file reader has a date-folder watermark, the offset applies
only to that folder key; __file_modification_time remains anchored to its
saved value. Other readers adjust datetime values in the stored watermark. See
Late-arriving and updated files for the two-stage
folder and mtime selection.
Override the connection for one source¶
A look-back in source.configure
overrides the entire connection-level look-back when the source has a
non-empty backward configuration:
{
"source": {
"connection_name": "orders_source",
"table": "orders",
"watermark_columns": ["updated_at"],
"configure": {
"backward_hours": 12
}
}
}
This source uses twelve hours, not “the connection's three days plus twelve hours”. If a source should inherit the connection default, omit its backward fields. If it should have a combined offset, author the complete combination at the source level:
{
"source": {
"connection_name": "orders_source",
"table": "orders",
"watermark_columns": ["updated_at"],
"configure": {
"backward": {"days": 2, "hours": 6}
}
}
}
Use a closing-day strategy for monthly corrections¶
closing_day computes an absolute start boundary from the current date rather
than subtracting a fixed number of days. It is useful when upstream closes a
period on a known day of each month:
{
"source": {
"connection_name": "orders_source",
"table": "orders",
"watermark_columns": ["updated_at"],
"configure": {
"backward": {"closing_day": 10}
}
}
}
When closing_day is present, it takes priority over fixed offset keys. The
optional months and years values can move the closing-day boundary farther
back. Use a day valid for the upstream business calendar and test a boundary
around month/year changes before production use.
Combine look-back with watermark-window replacement¶
Window replacement combines a source watermark and authored look-back with
the Destination setting
destination.configure.replace_by_watermark
under load_type: merge_overwrite. The effective source window must cover the
destination delete scope. The source must return every row to retain in that
scope and preserve or deterministically map the watermark output column. With
a usable replacement window, merge_keys are not required; the key-based
fallback needs them.
The Cross-boundary combinations section collects these source and destination requirements.
An ordinary empty incremental read does not replace a window. An explicitly bounded empty replay can delete its window. Follow the complete window recipe for its assumptions, first run, output mapping, and delete/append failure boundary.
Relate look-back to API ranges and replay¶
Look-back changes the lower bound derived from the saved watermark. API range
splitting separately divides a [from, to) interval into requests; configure
that on the API source with watermark_to_param and a range interval. See
Push down and split watermark ranges
for timezone and inclusive-upper-bound handling.
Replay supplies an explicit bounded range and can cap an API range's upper bound. It is an operational choice, not a replacement for ordinary look-back. Keep replay configuration in the operations guide and use this section only to understand how the source's effective lower bound is formed.
Validate an incremental configuration¶
- Set
watermark_columnsand confirm every selected column is returned by the table/query/API response. - Choose either connection inheritance or a complete source override.
- Do not author computed
date_backward. - For
closing_day, test month-end, leap-year, and upstream period-boundary behavior. - For window replacement, pair
merge_overwritewithreplace_by_watermark: trueand a usable look-back. - Confirm the source returns complete window coverage and preserves watermark columns.
- Test first run, normal empty increment, and explicit replay separately.
Filter rows at read time (source.filter_expression)¶
filter_expression is a SQL predicate combined with or applied after the
watermark condition by the source reader, before transforms. It references
columns in the reader's result. For a database source.query, use columns or
aliases returned by that query, not columns hidden inside it.
"source": {
"connection_name": "postgres_src",
"schema_name": "public",
"table": "orders",
"watermark_columns": ["updated_at"],
"filter_expression": "status = 'active' AND region = 'US'"
}
This is particularly useful when the source table holds multiple logical datasets and you only want one segment, or when you want to exclude known bad data before it enters the pipeline at all.
Database sources¶
For SQL table/query sources, filter_expression is combined with the generated
WHERE clause alongside the watermark condition. In query mode, the reader
wraps your SQL as a derived table and filters its output columns, so the query
must be valid when nested and expose the watermark and predicate columns:
-- generated SQL (conceptual)
SELECT * FROM orders
WHERE updated_at > '2024-01-01'
AND status = 'active'
AND region = 'US'
File, Delta, Iceberg, API, and function sources¶
For file, lakehouse, API, and Python function readers, the predicate is applied by the engine to the returned DataFrame after the read. For API reads this happens after all pages/ranges have been fetched and before the transform pipeline starts; it does not reduce HTTP volume.
When to use source.filter_expression vs transform.filter_expression¶
source.filter_expression |
transform.filter_expression |
|
|---|---|---|
| Stage | Read time (earliest possible) | Transformer order 35 (after ColumnAdder) |
| Available columns | Reader output columns (including query aliases and available watermark columns), but not columns added by transforms | Reader output + columns from additional_columns |
| Best for | Excluding rows that should never enter the pipeline | Filtering on computed/derived columns |
For an incremental source, the pipeline selects the reader's new watermark
before transform filters run and saves it after a successful write. Rows removed
only by transform.filter_expression do not automatically return if you widen
that filter later. Plan a bounded replay
when previously filtered history must be restored.
Common mistakes¶
| Symptom | Likely cause | Fix |
|---|---|---|
ValidationError: format not allowed |
connection_type and format don't match |
Use lakehouse + delta, not file + delta |
| Path not found | base_path / schema_name / table resolves wrong |
Check base_path exists; schema_name is optional |
APIReader requires 'base_url' |
API connection used url instead of base_url |
Put the root URL in connection.configure.base_url |
PythonFunctionReader requires source.python_function |
Function path was put on the connection or omitted | Put a dotted path like mypkg.sources.load_orders on source.python_function |
| No incremental filtering | watermark_columns missing from source block |
Add "watermark_columns": [...] and verify the selected reader's watermark behavior |
| API fetches all pages despite incremental filtering | Watermark columns are set, but request push-down mapping is absent | Configure API watermark mapping when the provider supports it; local filtering alone does not reduce fetch volume |
| Credentials exposed in logs | URL/token/password is hardcoded | Move the field to configure, store the secret name there, and resolve with secrets_ref |
| API pages repeat or skip records | Pagination/range boundary does not match the provider contract | Use the exact pagination paths/parameter names and review API source configuration |
Next¶
→ Transform patterns · Destination & load patterns · Datatypes and schema hints