Write Output & Auto-Sink
Write a finished run's output to a database, warehouse, bucket, cloud drive, topic, stream, endpoint or business system, with replace, append, fail and merge modes, and send it automatically.
What it does
Writes one output of a completed run straight to a connection: a database or warehouse table, a file in a bucket or cloud drive, a worksheet, a search index, a topic or stream, an HTTP endpoint, or records in a business or plant system. Every connector type except Google Analytics 4 can be a destination. In the dashboard, use the database button on a file's row or the Write to a destination panel under the results: pick a connection, the pipeline stage, the destination name, what happens if it already exists, and the destination's own file format and options. Partitioned writes and auto-sink are set through the API and the Python SDK.
Modes
if_exists | What it does |
|---|---|
replace | The target ends up holding exactly this output: a table is recreated, a file is overwritten. |
append | Adds the rows to the target, creating it when it does not exist. |
fail | Creates the target, and stops if it already exists. |
merge | Updates the rows that match on merge_keys and inserts the rest. |
Each destination offers the modes that apply to it. Leave if_exists out to use the destination's default: replace where it is offered, otherwise append. A request for a mode the destination does not offer is answered with the list of supported ones. A saved connection reports its own modes, file formats and options in the write field of GET /api/v1/credentials.
Destinations
| Types | table_name | Modes | Options |
|---|---|---|---|
sql, oracle, sap_hana, synapse, redshift, snowflake, bigquery, databricks, fabric | table, optionally qualified (schema.table, dataset.table) | replace, append, fail, merge | |
timescaledb | table | replace, append, fail, merge | time_column makes a new table a hypertable |
clickhouse | table | replace, append, fail | order_by: sort key columns |
delta | the connection's table path | replace, append, fail, merge | |
mongodb | collection or database.collection | replace, append, fail, merge | |
elasticsearch | index | replace, append, fail, merge | |
cassandra | table or keyspace.table | append, replace, fail | merge_keys name the primary key of a new table |
dynamodb | table | append, fail | merge_keys name the key of a new table |
influxdb | measurement | append, replace | time_column (required), tag_columns |
s3, gcs, azure_blob, adls_gen2, google_drive, dropbox, onedrive, box | file name or path | replace, append, fail | file_format; folder_id for Google Drive and Box |
google_sheets | worksheet | replace, append, fail, merge | |
kafka, kinesis | topic or stream | append | partition_key_column |
mqtt | topic | append | qos (0 or 1) |
http | a URL, or any label to use the connection's URL | append | batch_size, method (POST, PUT, PATCH) |
salesforce, hubspot | object type | append, merge | |
netsuite | record type | append, merge | |
sharepoint | list | append | |
stripe | customers, products or prices | append | |
pi_web_api | point name | append, merge | point_column, time_column, value_column |
opcua | any label; node ids come from a column | append | node_column, value_column |
Merge needs merge_keys, which must be columns of the output. Topics, streams and endpoints receive one JSON document per row. Row limits per write: Redshift and HTTP 1,000,000; MQTT 200,000.
File destinations
Amazon S3, Google Cloud Storage, Azure Blob Storage, ADLS Gen2, Google Drive, Dropbox, OneDrive and Box take a file_format: csv (default), parquet, xlsx, json, jsonl, ndjson, feather or orc. An extension on table_name (reports/customers.parquet) selects the format too.
- replace overwrites one fixed file.
- append writes a new timestamped file next to the earlier ones.
- fail stops if the file exists.
Destinations that change live records
Writing to Salesforce, HubSpot, NetSuite, SharePoint lists, Stripe, PI Web API or OPC UA creates or changes live records in that system. These destinations are written to only through a saved connection whose Allow writing switch is on (allow_write on the credential), never with inline secrets. They do not offer replace.
| Destination | Rows per write | Notes |
|---|---|---|
| Salesforce | 150,000 | Merge upserts on one external-id field. |
| HubSpot | 100,000 | Merge upserts on one unique property. |
| NetSuite | 1,000 | Merge upserts on the external id. |
| SharePoint lists | 5,000 | One list item per row. |
| Stripe | 10,000 | Creates customers, products and prices only. |
| PI Web API | 500,000 | Merge replaces the value at the same timestamp; no merge_keys needed. |
| OPC UA | 500 | Writes the last row for each node. |
Partitioned writes
partition_by names a column of the output and writes one target per distinct value, up to 1,000 partitions. Object stores get name/column=value/data; other destinations get name_value. A date column is formatted with partition_format (default %Y-%m-%d).
Request
POST /api/v1/sessions/{session_id}/write-output
{ "connector_type": "sql", "credential_id": "cred-uuid",
"table_name": "customers_prepared", "output_stage": "dsg",
"if_exists": "merge", "merge_keys": ["customer_id"] }
| Key | Default | Meaning |
|---|---|---|
table_name | required | |
connector_type | the saved connection's type, else sql | Any connector type except ga4. |
credential_id | Or inline_secrets. | |
output_stage | dsg | dsg, dsm, cds, mdh, dtc, anormaly_fixed, multimodal, compleated_dataset; falls back to dsg, cds, mdh, dtc. Inference and retraining runs write their scored rows with targets. |
if_exists | the destination's default | replace, append, fail, merge. |
merge_keys | List or comma-separated. | |
file_format | csv | File destinations. |
partition_by | One target per distinct value of this column. | |
partition_format | %Y-%m-%d | Date format for a date partition column. |
options | Destination options from the table above, as an object: {"time_column": "ts", "tag_columns": ["site"]}. |
Success: { success, rows_written, table, message }, with partitions for a partitioned write. When a destination accepts some rows and declines others, the answer reports rows_written, rows_failed and the first reasons in errors. Python SDK: write_output(session_id, connector_type, table_name, secrets=None, credential_id=None, output_stage="dsg", if_exists=None, merge_keys=None, file_format=None, partition_by=None, partition_format=None, options=None).
Auto-sink
The same body plus "enabled": true, stored on an automation. It is written once when the run completes successfully, for full-pipeline, inference and retraining runs alike. Set it through the API or the Python SDK. The outcome is reported in the run's result as auto_sink_result.
| Surface | Key | Applies to |
|---|---|---|
| Folder listener / SFTP inbox | auto_sink_config | Every run. |
| Scheduled run | pipeline_config.auto_sink | Every run. |
| Trigger | target_pipeline_config.auto_sink | Every run. |
| Connector run | auto_sink | That run. |
LM Readiness sessions
For an LM Readiness session, the replayed rows are written (a preparation writes its training partition). List columns such as token ids are written to a database as JSON text, and to an object store unchanged.