INGEST layer
The Ingest layer pulls data out of source systems and lands it in the data lake. It handles the what, when and how of extraction.
This page is the reference for what the layer does. For the screen-by-screen guide to configuring it, see Configuring ingestion.
The Ingest layer is responsible for pulling data out of source systems and landing it in the data lake. It handles the what, when, and how of data extraction.
Connection types
Before you can create exports, a connection must be configured for your source system. The connection type decides which fields the connection carries; within a type, the runtime runs exactly one operator class per source system, and the operator decides which of those fields are actually read.
| Connection Type | Source | Operators |
|---|---|---|
| Database (DB) | SQL databases reached with SQL | PostgreSQL, SQL Server, Oracle, MySQL |
| API | HTTP endpoints returning JSON | Generic REST, Heartpace, Salesforce Service Cloud, Talkdesk, SharePoint (Graph), Microsoft Teams, Google Drive / Sheets |
| File | File storage systems | Local disk, Amazon S3, Azure Data Lake Gen2, SFTP |
| Custom | Message queues and bespoke sources | AWS SQS, Azure Service Bus, Redis, plus per-customer connectors |
| Manual | No automated extraction | Data uploaded manually |
The queue operators ship under the Custom contract, so a queue's coordinates travel as key=value connection properties rather than as dedicated fields. Nothing about the queue behaviour changes because of it.
Three connection settings shape every export made from that connection:
| Setting | Effect |
|---|---|
AddServiceColVal | When set, every exported record gets a __service column carrying this value. Empty disables the column. |
ColumnDelimiter | Field separator when CDC writes CSV output. Ignored for JSON. |
LineDelimiter | Record separator passed to CDC. Every operator except the generic database one forces a newline regardless of this value. |
The full connection field reference — credentials, authentication methods, storage settings and queue coordinates — is in the Ingest developer guide.
Export types
Each source system can have one or more source exports. The export type matches the connection type, and fields an operator ignores are shown but marked as unused.
Database Exports
Extract data using SQL queries from connected databases. The query is streamed in 25,000-row chunks and written as newline-delimited JSON.
| Setting | Description |
|---|---|
databaseName | Reference to the connected database. Read-only, taken from the connection. |
fromClause | Table, view, function or stored procedure to query. Must be schema-qualified (dbo.Customers) — the metadata lookup splits on the dot. |
whereClause | Optional SQL condition. Inserted verbatim, so it must include the WHERE keyword. An incremental filter is appended with AND. |
sqlOverride | Full custom SQL query — replaces the generated query entirely. Time tokens are still substituted, but the incremental filter is not appended; build it into the override yourself. |
inclTables | Extra tables whose column metadata is gathered for type coercion. Defaults to the from clause. Use it to list every object behind a view or procedure, so column types are resolved correctly. |
inclColumns | Joined with commas straight into the SELECT clause, so aliases and expressions work. ["*"] selects everything. |
filterColumn | Enables incremental loading. The comparison is numeric, so the column must be an increasing number. |
lastFilterValue | Watermark from the previous run, posted back to the platform after each run. 0 on the first run. |
Two modes:
- Standard Mode — build the query from
fromClause+whereClause+ column selection - SQL Override Mode — write a complete custom SQL statement that bypasses the generated query
Database exports also compute a __checksum per record, which enables the fast CDC comparison path.
API Exports
Extract data from HTTP endpoints. Output is always newline-delimited JSON — outputFormat is ignored here.
| Setting | Description |
|---|---|
apiUrl | Base API URL. Read-only, taken from the connection. |
apiPath | Endpoint path, appended to the base URL. |
filterCondition | Key=value pairs sent as request parameters — a query string, or a JSON body when the connection's method header is post. |
Filter conditions are string arrays like ["user_status=active", "date_from=2024-01-01"]. Time tokens are substituted in the values, and a value containing braces is parsed as JSON, so a JSON request body can be configured as a filter condition.
Errors: a 401 triggers one token refresh and a single retry. Any other status of 400 or above fails the workload, as does a response with no detectable content type.
Operator specialisations handle pagination and job models the generic operator cannot:
| Operator | Pagination model |
|---|---|
| Heartpace | limit/offset, driven by the response's meta block |
| Salesforce Service Cloud | Follows nextRecordsUrl until the result set reports itself done |
| Talkdesk | Asynchronous report jobs, polled every 10 seconds for up to 50 minutes |
| SharePoint (Graph) | Recursive folder traversal, ordered newest-first with early termination |
| Teams, Google Drive | A single request returning a values array whose first row is the header |
File Exports
Extract data from file storage systems. Matching files are copied locally, then either converted with DuckDB or passed through untouched.
| Setting | Description |
|---|---|
filePath | Source directory prefix, appended to the connection's base path. Time tokens are substituted here, so a partitioned source can be addressed with [__SCHEDULEDTIME_AS_HIVE__]. |
fileName | A substring match against the file name — not a glob. |
filterCondition | File selection and reader options (below). |
outputFormat | json or csv convert the files; binary uploads them untouched and skips CDC. Converted output is written gzipped. |
File selection keys:
| Key | Default | Description |
|---|---|---|
delete | False | Remove source files after a successful copy. There is no undo. |
use_filename_timestamps | False | True matches the timestamp in the file name; False filters on last-modified time. |
file_timestamp_format | [YYYY][MM][DD] | Timestamp pattern inside file names. The finest part present also sets the search granularity. |
start_timestamp | now | ISO datetime — start of the time window. |
end_timestamp | now | ISO datetime — end of the time window. |
find_partitions | False | Amazon S3 only. Lists sub-prefixes under the base path and searches each. Not partition-aware — for a year=/month=/day= layout, put a Hive token in filePath instead. |
Reader keys:
| Key | Default | Description |
|---|---|---|
excel_sheet, excel_range | – | Sheet name and cell range for .xlsx files. |
csv_has_header | true | First row is the header. |
csv_sample_size | 10000 | Rows sampled for type detection. |
json_sample_size | 10000 | Rows sampled for schema inference. |
json_format | auto | auto, newline_delimited, array, and the other DuckDB JSON formats. |
json_union_by_name | true | Union schemas across files instead of requiring identical ones. |
file_encoding | utf-8 | CSV encoding. |
Format is detected in the order .xlsx → .parquet → JSON → CSV; files matching none of them fall back to a binary copy. All matched files are read in one pass, so they must share a schema unless json_union_by_name covers the difference.
Queue Exports
Turn a message stream into batch files. One workload run drains the queue and then ends, so the export stays on the normal schedule instead of running its own loop.
The queue itself is configured on the connection; the export adds only the batch bounds, through filterCondition:
| Key | Default | Description |
|---|---|---|
max_batch_bytes | 31457280 (30 MB) | Close the batch file at this uncompressed size. |
max_batch_seconds | 900 (15 min) | Close the batch file after this much wall-clock time. |
max_batches_per_run | 0 (unlimited) | Cap a hot queue so it cannot starve other workloads. |
Behaviour worth knowing:
- An empty queue writes nothing. No zero-byte file is created and no downstream processing runs.
- CDC is bypassed. Events are consumed once and never updated, so set
enableCDC: 0to avoid a pointless baseline file. - Delivery is at-least-once. Acknowledgement happens after the batch is flushed to disk, so a crash mid-run redelivers rather than loses.
- Malformed messages are forwarded, written with a
__parseErrormarker, rather than left on the queue as poison. - Batch files are named
<alias>_<suffix>_batch<NNNNNN>.txt.
This family targets high-volume streams of small events. Land large payloads in blob storage and ingest them as a file export, letting the queue carry only the pointer.
Custom Exports
Minimal configuration for proprietary integrations. The connection's properties are the only configuration channel, filterCondition is substituted and handed to the connector, and outputFormat becomes the output file extension.
Scheduling
All export types support scheduling:
| Setting | Description | Example |
|---|---|---|
scheduleCycle | Frequency | never, minute, hourly, daily, monthly. Intervals below 15 minutes are not recommended. |
scheduleCycleInterval | Multiplier | 2 with hourly = every 2 hours |
scheduleDoNotStartBeforeTime | Earliest allowed start | 02:00:00 (don't start before 2 AM) |
never disables an export without deleting it.
Export strategies
Choosing the right export strategy is critical for balancing data freshness against cost and volume.
Full Export
Extracts the entire dataset every time. Simple but expensive for large tables.
When to use:
- Small reference/dimension tables
- Source systems that don't support change tracking
- Initial loads or periodic full refreshes
How to configure:
- No special settings needed — this is the default behaviour
- Leave
stopAtRow: -1for no row limit
Incremental Export (Watermark)
Reads only the rows the source has added since the last run. Database sources only.
When to use:
- Large transaction tables where full exports are too costly
- Tables carrying a reliably increasing key or sequence
How to configure:
- Set
filterColumnto an increasing numeric column — the comparison is numeric, so a text or date column will not work - Leave
lastFilterValuealone; the agent posts the new watermark back after each run, starting from0
How it works:
- The generated query gains a
filterColumn > lastFilterValuecondition - The column's maximum is tracked per chunk and returned as the new watermark
- A run truncated by
stopAtRowdoes not advance the watermark, so nothing is skipped next time sqlOverridebypasses this entirely — the filter is not appended to a custom query
Change Data Capture (CDC)
Compares this run's extract against the previous one and marks each record as changed, deleted or unchanged. Available for every source type that produces rows.
When to use:
- Sources with no change indicator of their own
- Any source where downstream needs to know about deletions
How to configure:
- Set
enableCDC: 1on the export (0exports everything every run) - Mark the fields to track with
includeInCDC: 1in the sourcefile structure - Use a
validFromfield to identify which version of a record is current - Set
sortOrderon the versioning field (positive = descending, most recent first)
How it works:
- The previous extract is kept as a baseline file named
<alias>_latestversion.txtand compared in DuckDB - Records are emitted with a
__cdccolumn ofchange,deleteorequal - Where a
__checksumis present — database sources — an indexed anti-join replaces the full comparison, which is significantly faster on wide tables - If the column set changed between runs, the whole new load is returned and flagged as new
- Output is split at roughly 2 GB per file; every part is uploaded
No baseline file means a full load, so the first run after enabling CDC always exports everything. Each run also unions in a small sample of the new data as a deliberate re-seed, so CDC output is never a pure delta.
stopAtRow forces CDC off for that run and discards the truncated output rather than keeping it as the baseline — a capped extract compared against a full baseline would read every row it did not write as a delete. It is a smoke-test knob, not a throttle.
Binary file exports and queue exports bypass CDC entirely.
Sliding Window
A time-bounded extraction that moves forward with each run, because the window is expressed as time tokens that re-render every run rather than as fixed dates.
When to use:
- File ingestion where filenames or partition paths carry timestamps
- APIs with date-range parameters
- Scenarios where you need the last N days/hours
How to configure (file exports):
- Set
use_filename_timestamps: Trueif the timestamp is in the file name, andfile_timestamp_formatto match it — otherwise the window is matched against last-modified time - Express the window with tokens rather than dates:
filterCondition: ["start_timestamp=[__SCHEDULEDTIME_MINUS_1_DAYS__]", "end_timestamp=[__SCHEDULEDTIME__]"]
- For a partitioned source, address the partition from
filePathwith[__SCHEDULEDTIME_AS_HIVE__]
How to configure (API exports): Use the same tokens in the request parameters:
filterCondition: ["date_from=[__LASTEXECUTION__]", "date_to=[__SCHEDULEDTIME__]"]
One run reads one window. [__SCHEDULEDTIME__] is the slot the run belongs to rather than the clock time, so a retry reads the same window — but a missed day means replaying that schedule, not widening the window.
Column selection
Control which columns appear in the exported data.
Include & Exclude Modes
Include Mode (default):
By default, inclColumns: ["*"] — all columns are included. To include only specific columns:
inclColumns: ["CustomerID", "Name", "Email", "CreatedDate"]
exclColumns: []
Exclude Mode: Include everything except specific columns:
inclColumns: ["*"]
exclColumns: ["InternalNotes", "TempFlag", "DebugData"]
Rules:
- If
inclColumnscontains"*", all fields are included exclColumnsacts as a blacklist applied on top of the inclusion list — deny wins- The UI prevents leaving no included fields — it defaults back to
["*"] - Column names are case-sensitive and must match the source exactly
On database exports the two lists behave differently. inclColumns is joined into the SELECT clause, so it is the only way to drop a column; exclColumns merely suppresses type coercion and does not remove anything from the output.
Data transformations
Substitutions
Substitutions derive or default a column while the export is being written. Each one is defined as { column, alias, replace, type }:
| Property | Description |
|---|---|
column | The source column the substitution reads |
alias | Optional target column. Without it, the source column is overwritten |
replace | The expression or literal used as the replacement |
type | transform — replace is an expression evaluated against the cell value, applied only when the value is not null. null — replace is the literal default used when the value is null |
Typical uses are null replacement and masking sensitive values before they leave the source.
replace is evaluated as code by the agent. Treat export definitions as trusted configuration, not as user input.
Output Format & File Naming
Output Format (outputFormat):
| Format | Description |
|---|---|
json | Newline-delimited JSON. The default, and the only format the plain API operators write |
csv | Delimited text, using the connection's ColumnDelimiter |
binary | Files uploaded untouched. Skips conversion and CDC; file sources only |
TABLE, ICEBERG and JSON are target formats, chosen on the sourcefile in DLS — they are not export output formats.
Dynamic File Naming (Suffix):
The suffix appended to the output file name is either a literal string or a time token:
| Token | Value |
|---|---|
[__SCHEDULEDTIME__] | The slot this run belongs to — unchanged by a late start or a retry |
[__EXECUTIONTIME__] | Actual execution time, in the agent's timezone |
[__LASTEXECUTION__] | Previous successful run's scheduled time (1900-01-01T00:00:00 on the first run) |
The same tokens work in request parameters, file paths, where clauses and the ends of a file-selection window. They take an optional offset and an optional format:
[__<TIME>[_MINUS_<N>_<UNIT>][_AS_<FORMAT>]__]
<UNIT> is SECONDS, MINUTES, HOURS, DAYS or WEEKS; there is no _PLUS_, since an export reads a window that has already happened. <FORMAT> is TIMESTAMP, DATETIME, DATE, DATECOMPACT, TIME, ISO, EPOCH, EPOCHMS or HIVE, or a literal pattern built from [YYYY] [YY] [MM] [DD] [HH] [MI] [SS]. A value with no _AS_ segment inherits the export's dateTimeFormat, which defaults to TIMESTAMP — or DATETIME for database sources, which suits SQL literals.
HIVE renders a partition path such as year=2026/month=8/day=13 rather than a datetime, for reading sources laid out that way. It belongs in a file path and cannot be used as the export-wide format. See the Ingest developer guide for the grain and padding variants.
Row Limiting:
Set stopAtRow to limit the size of an export: -1 is no limit, and any positive number stops the run there. What it counts depends on the source:
| Source | Unit |
|---|---|
| API, custom | Rows |
| Database | Chunks of 25,000 rows |
| File, converted | Rows written — every matching file is still downloaded |
| File, binary | Files — the only case where the unit is not rows |
| Queue | Messages, counted across the whole run |
Because a capped run also disables CDC, stopAtRow is for testing a new export rather than for limiting a production one.
Technical Columns
Every export adds columns of its own alongside the source data:
| Column | Added when | Value |
|---|---|---|
__service | AddServiceColVal is set on the connection | That value |
__scheduledAt | Always | The run's scheduled time, falling back to export time |
__exportedAt | Always | When the export started writing |
__checksum | Database sources | Hash of the record, enabling the fast CDC path |
__cdc | CDC is enabled | change, delete or equal |
__parseError | A queue message could not be parsed | The parse error, with the raw body preserved |
Destination paths
Exported data lands in configured zone paths:
| Zone | Source | Purpose |
|---|---|---|
| Landing Zone | From Settings (landingZoneName) | Temporary staging |
| Raw Zone | From Settings (rawZoneName) | Permanent archive |
| Trusted Zone | From Settings (trustedZoneName) | Validated data |
Default path pattern: [system]/[filename]/[YYYY]/[MM]/[DD]/
Sourcefile structure
Once data is ingested, it's described by a sourcefile — the schema definition for the incoming data.
Key Properties
| Property | Description |
|---|---|
sourceFilename | Unique identifier |
system | Source system reference |
fileType | CSV, JSON, or XML |
fileEncoding | UTF-8, UTF-16, Windows-1252 |
columnDelimiter | Delimiter for CSV files |
targetMethod | TRANSACTION, APPEND, CHANGES ONLY, LATEST VERSION or OVERWRITE — see the DLS field reference |
targetFormat | TABLE, ICEBERG or JSON |
targetNormalization | NONE, LISTS, or LISTS AND OBJECTS |
fileCompressionType | Compression type: gzip or none |
enableEncryption | Encryption flag |
Hierarchical File Structure
For nested data (e.g. JSON with arrays), the sourcefile supports a hierarchical structure:
Level 0 (OBJECT): Root
├── customer.id (keyFieldIndicator=1)
├── customer.name
└── customer.loadDate (validFrom=1)
Level 1 (LIST): orders (useForSplittingRecords=1)
├── order.orderId
├── order.amount
└── order.status
Each level can be an OBJECT (single record) or LIST (array that gets split into rows).
Field-Level Properties
| Property | Description |
|---|---|
fieldKey | Source column name (case-sensitive) |
fieldAlias | Optional renamed column name |
dataType | Varchar, Integer, Decimal, Boolean, Time, Date, Timestamp |
keyFieldIndicator | 1 = part of primary key |
keyOrder | Position in composite key |
fieldOrder | Column ordering in target |
excludeField | 1 = skip this field entirely |
excludeFromProfiling | 1 = don't run quality checks |
sensitive | 1 = PII/sensitive data flag |
includeInCDC | 1 = track changes for this field |
validFrom | 1 = SCD Type 2 versioning field (max one per level) |
fieldDomain | Governance classification category |
Best practices
Choosing an Export Strategy
| Scenario | Recommended Strategy |
|---|---|
| Small dimension table (<100K rows) | Full export with CDC, daily |
| Large transaction table with an increasing key | Incremental watermark on filterColumn |
| Large table with no change indicator | Full export with CDC — the checksum path keeps the comparison cheap |
| File drop with timestamped filenames | Sliding window on start_timestamp / end_timestamp |
Partitioned lake (year=/month=/day=) | Hive token in filePath, one partition per run |
| API with pagination + date filters | Sliding window through request parameters |
| Message queue | Batch bounds on the connection, CDC off |
| Initial data load | Full export, then switch to incremental |
Column Selection Tips
- Start with
inclColumns: ["*"]and useexclColumnsto remove unwanted columns - On database exports, drop columns from
inclColumnsinstead —exclColumnsdoes not remove them there - Exclude large text/blob columns that aren't needed downstream
- Exclude internal/debug columns that add noise
- Remember: column names are case-sensitive
Scheduling Considerations
- Use
scheduleDoNotStartBeforeTimeto avoid running during peak hours - For dependent exports, stagger schedules (e.g. dimension tables before fact tables)
- Use
monthlycycle for slowly-changing reference data - Use
hourlywith a small interval for near-real-time needs
Next steps
- DLS layer — what happens to a delivery once it has landed
- Configuring ingestion — the field-by-field configuration guide
- Architecture — how the layers fit together