Skip to main content

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 TypeSourceOperators
Database (DB)SQL databases reached with SQLPostgreSQL, SQL Server, Oracle, MySQL
APIHTTP endpoints returning JSONGeneric REST, Heartpace, Salesforce Service Cloud, Talkdesk, SharePoint (Graph), Microsoft Teams, Google Drive / Sheets
FileFile storage systemsLocal disk, Amazon S3, Azure Data Lake Gen2, SFTP
CustomMessage queues and bespoke sourcesAWS SQS, Azure Service Bus, Redis, plus per-customer connectors
ManualNo automated extractionData uploaded manually
note

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:

SettingEffect
AddServiceColValWhen set, every exported record gets a __service column carrying this value. Empty disables the column.
ColumnDelimiterField separator when CDC writes CSV output. Ignored for JSON.
LineDelimiterRecord 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.

SettingDescription
databaseNameReference to the connected database. Read-only, taken from the connection.
fromClauseTable, view, function or stored procedure to query. Must be schema-qualified (dbo.Customers) — the metadata lookup splits on the dot.
whereClauseOptional SQL condition. Inserted verbatim, so it must include the WHERE keyword. An incremental filter is appended with AND.
sqlOverrideFull 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.
inclTablesExtra 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.
inclColumnsJoined with commas straight into the SELECT clause, so aliases and expressions work. ["*"] selects everything.
filterColumnEnables incremental loading. The comparison is numeric, so the column must be an increasing number.
lastFilterValueWatermark 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.

SettingDescription
apiUrlBase API URL. Read-only, taken from the connection.
apiPathEndpoint path, appended to the base URL.
filterConditionKey=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:

OperatorPagination model
Heartpacelimit/offset, driven by the response's meta block
Salesforce Service CloudFollows nextRecordsUrl until the result set reports itself done
TalkdeskAsynchronous report jobs, polled every 10 seconds for up to 50 minutes
SharePoint (Graph)Recursive folder traversal, ordered newest-first with early termination
Teams, Google DriveA 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.

SettingDescription
filePathSource 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__].
fileNameA substring match against the file name — not a glob.
filterConditionFile selection and reader options (below).
outputFormatjson or csv convert the files; binary uploads them untouched and skips CDC. Converted output is written gzipped.

File selection keys:

KeyDefaultDescription
deleteFalseRemove source files after a successful copy. There is no undo.
use_filename_timestampsFalseTrue 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_timestampnowISO datetime — start of the time window.
end_timestampnowISO datetime — end of the time window.
find_partitionsFalseAmazon 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:

KeyDefaultDescription
excel_sheet, excel_range–Sheet name and cell range for .xlsx files.
csv_has_headertrueFirst row is the header.
csv_sample_size10000Rows sampled for type detection.
json_sample_size10000Rows sampled for schema inference.
json_formatautoauto, newline_delimited, array, and the other DuckDB JSON formats.
json_union_by_nametrueUnion schemas across files instead of requiring identical ones.
file_encodingutf-8CSV 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:

KeyDefaultDescription
max_batch_bytes31457280 (30 MB)Close the batch file at this uncompressed size.
max_batch_seconds900 (15 min)Close the batch file after this much wall-clock time.
max_batches_per_run0 (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: 0 to 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 __parseError marker, 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:

SettingDescriptionExample
scheduleCycleFrequencynever, minute, hourly, daily, monthly. Intervals below 15 minutes are not recommended.
scheduleCycleIntervalMultiplier2 with hourly = every 2 hours
scheduleDoNotStartBeforeTimeEarliest allowed start02: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: -1 for 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:

  1. Set filterColumn to an increasing numeric column — the comparison is numeric, so a text or date column will not work
  2. Leave lastFilterValue alone; the agent posts the new watermark back after each run, starting from 0

How it works:

  • The generated query gains a filterColumn > lastFilterValue condition
  • The column's maximum is tracked per chunk and returned as the new watermark
  • A run truncated by stopAtRow does not advance the watermark, so nothing is skipped next time
  • sqlOverride bypasses 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:

  1. Set enableCDC: 1 on the export (0 exports everything every run)
  2. Mark the fields to track with includeInCDC: 1 in the sourcefile structure
  3. Use a validFrom field to identify which version of a record is current
  4. Set sortOrder on the versioning field (positive = descending, most recent first)

How it works:

  • The previous extract is kept as a baseline file named <alias>_latestversion.txt and compared in DuckDB
  • Records are emitted with a __cdc column of change, delete or equal
  • Where a __checksum is 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
note

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.

warning

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):

  1. Set use_filename_timestamps: True if the timestamp is in the file name, and file_timestamp_format to match it — otherwise the window is matched against last-modified time
  2. Express the window with tokens rather than dates:
filterCondition: ["start_timestamp=[__SCHEDULEDTIME_MINUS_1_DAYS__]", "end_timestamp=[__SCHEDULEDTIME__]"]
  1. For a partitioned source, address the partition from filePath with [__SCHEDULEDTIME_AS_HIVE__]

How to configure (API exports): Use the same tokens in the request parameters:

filterCondition: ["date_from=[__LASTEXECUTION__]", "date_to=[__SCHEDULEDTIME__]"]
note

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 inclColumns contains "*", all fields are included
  • exclColumns acts 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
warning

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 }:

PropertyDescription
columnThe source column the substitution reads
aliasOptional target column. Without it, the source column is overwritten
replaceThe expression or literal used as the replacement
typetransform — 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.

warning

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):

FormatDescription
jsonNewline-delimited JSON. The default, and the only format the plain API operators write
csvDelimited text, using the connection's ColumnDelimiter
binaryFiles 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:

TokenValue
[__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:

SourceUnit
API, customRows
DatabaseChunks of 25,000 rows
File, convertedRows written — every matching file is still downloaded
File, binaryFiles — the only case where the unit is not rows
QueueMessages, 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:

ColumnAdded whenValue
__serviceAddServiceColVal is set on the connectionThat value
__scheduledAtAlwaysThe run's scheduled time, falling back to export time
__exportedAtAlwaysWhen the export started writing
__checksumDatabase sourcesHash of the record, enabling the fast CDC path
__cdcCDC is enabledchange, delete or equal
__parseErrorA queue message could not be parsedThe parse error, with the raw body preserved

Destination paths​

Exported data lands in configured zone paths:

ZoneSourcePurpose
Landing ZoneFrom Settings (landingZoneName)Temporary staging
Raw ZoneFrom Settings (rawZoneName)Permanent archive
Trusted ZoneFrom 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
PropertyDescription
sourceFilenameUnique identifier
systemSource system reference
fileTypeCSV, JSON, or XML
fileEncodingUTF-8, UTF-16, Windows-1252
columnDelimiterDelimiter for CSV files
targetMethodTRANSACTION, APPEND, CHANGES ONLY, LATEST VERSION or OVERWRITE — see the DLS field reference
targetFormatTABLE, ICEBERG or JSON
targetNormalizationNONE, LISTS, or LISTS AND OBJECTS
fileCompressionTypeCompression type: gzip or none
enableEncryptionEncryption 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
PropertyDescription
fieldKeySource column name (case-sensitive)
fieldAliasOptional renamed column name
dataTypeVarchar, Integer, Decimal, Boolean, Time, Date, Timestamp
keyFieldIndicator1 = part of primary key
keyOrderPosition in composite key
fieldOrderColumn ordering in target
excludeField1 = skip this field entirely
excludeFromProfiling1 = don't run quality checks
sensitive1 = PII/sensitive data flag
includeInCDC1 = track changes for this field
validFrom1 = SCD Type 2 versioning field (max one per level)
fieldDomainGovernance classification category

Best practices​

Choosing an Export Strategy
ScenarioRecommended Strategy
Small dimension table (<100K rows)Full export with CDC, daily
Large transaction table with an increasing keyIncremental watermark on filterColumn
Large table with no change indicatorFull export with CDC — the checksum path keeps the comparison cheap
File drop with timestamped filenamesSliding window on start_timestamp / end_timestamp
Partitioned lake (year=/month=/day=)Hive token in filePath, one partition per run
API with pagination + date filtersSliding window through request parameters
Message queueBatch bounds on the connection, CDC off
Initial data loadFull export, then switch to incremental
Column Selection Tips
  • Start with inclColumns: ["*"] and use exclColumns to remove unwanted columns
  • On database exports, drop columns from inclColumns instead — exclColumns does 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 scheduleDoNotStartBeforeTime to avoid running during peak hours
  • For dependent exports, stagger schedules (e.g. dimension tables before fact tables)
  • Use monthly cycle for slowly-changing reference data
  • Use hourly with a small interval for near-real-time needs

Next steps​