salesforce
Executes a Salesforce SOQL query once and emits a message for each record.
Introduced in version 4.90.3.
Runs a SOQL query against the Salesforce REST API, paginates through all result pages, and emits one message per record. When results are exhausted the input shuts down, letting the pipeline terminate gracefully (or the next input in a sequence to take over).
When to use this input
Use salesforce for:
-
One-shot extracts (for example, dump all Accounts into a warehouse).
-
Periodic full-table refreshes via a scheduled pipeline or sequence.
-
Backfills and ad-hoc queries.
-
Warming up a downstream pipeline before switching to CDC.
Use a different Salesforce input instead if:
-
You need continuous change events: use
salesforce_cdc. -
You need a GraphQL query (cross-object in one request): use
salesforce_graphql.
Metadata
This input adds the following metadata fields to each message:
-
sobject: The sObject API name (for example, "Account").
Authentication
Uses the Salesforce OAuth 2.0 Client Credentials flow. Create a Connected App in Salesforce, enable OAuth settings and the Client Credentials Flow, then supply the Consumer Key as client_id and Consumer Secret as client_secret.
-
Common
-
Advanced
input:
label: ""
salesforce:
org_url: "" # No default (required)
client_id: "" # No default (required)
client_secret: "" # No default (required)
api_version: v65.0
object: "" # No default (required)
columns: [] # No default (required)
where: "" # No default (optional)
args_mapping: "" # No default (optional)
auto_replay_nacks: true
input:
label: ""
salesforce:
org_url: "" # No default (required)
client_id: "" # No default (required)
client_secret: "" # No default (required)
api_version: v65.0
object: "" # No default (required)
columns: [] # No default (required)
where: "" # No default (optional)
args_mapping: "" # No default (optional)
prefix: "" # No default (optional)
suffix: "" # No default (optional)
auto_replay_nacks: true
http:
timeout: 5s
tls:
enabled: false
skip_cert_verify: false
enable_renegotiation: false
root_cas: ""
root_cas_file: ""
client_certs: []
proxy_url: ""
disable_http2: false
tps_limit: 0
tps_burst: 1
backoff:
initial_interval: 1s
max_interval: 30s
max_retries: 3
tcp:
connect_timeout: 0s
keep_alive:
idle: 15s
interval: 15s
count: 9
tcp_user_timeout: 0s
http:
max_idle_conns: 100
max_idle_conns_per_host: 0
max_conns_per_host: 64
idle_conn_timeout: 1m30s
tls_handshake_timeout: 10s
expect_continue_timeout: 1s
response_header_timeout: 0s
disable_keep_alives: false
disable_compression: false
max_response_header_bytes: 1048576
max_response_body_bytes: 10485760
write_buffer_size: 4096
read_buffer_size: 4096
h2:
strict_max_concurrent_requests: false
max_decoder_header_table_size: 4096
max_encoder_header_table_size: 4096
max_read_frame_size: 16384
max_receive_buffer_per_connection: 1048576
max_receive_buffer_per_stream: 1048576
send_ping_timeout: 0s
ping_timeout: 15s
write_byte_timeout: 0s
access_log_level: ""
access_log_body_limit: 0
Fields
api_version
Salesforce REST API version to target, prefixed with v. Affects endpoint paths (/services/data/{api_version}/…) and available fields/objects. Must be supported by your org; check Setup → Company Information. Older versions may lack recent fields.
Type: string
Default: v65.0
# Examples:
api_version: v65.0
# ---
api_version: v62.0
args_mapping
Optional Bloblang mapping whose result must be an array of values matching the count of ? placeholders in where. Values are SOQL-escaped: strings become quoted literals, timestamps become ISO-8601, booleans and numbers pass through. The mapping runs once at startup with no message context, so it can use literal values and functions that don’t read a message, such as now(), but can’t reference message content or metadata.
Type: string
# Examples:
args_mapping: root = [ (now() - "1h").ts_format("2006-01-02T15:04:05Z") ]
# ---
args_mapping: root = [ "Active", (now() - "24h").ts_format("2006-01-02T15:04:05Z") ]
auto_replay_nacks
Whether to automatically replay rejected messages (negative acknowledgements, or nacks) at the output level. If the cause of rejections persists, leaving this option enabled can result in back pressure.
Set auto_replay_nacks to false to delete rejected messages. Disabling auto replays can greatly improve memory efficiency of high throughput streams, as the original shape of the data is discarded immediately upon consumption and mutation.
Type: bool
Default: true
client_id
Consumer Key of the Salesforce Connected App authorized for the OAuth Client Credentials flow. Create the Connected App under Setup → App Manager → New Connected App, enable OAuth settings, enable the Client Credentials Flow under Flow Enablement, then copy the Consumer Key from Manage Consumer Details.
Type: string
client_secret
Consumer Secret of the Salesforce Connected App, paired with client_id. Sensitive; prefer environment variable interpolation (${SALESFORCE_CLIENT_SECRET}) over inlining.
|
This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets. |
Type: string
columns[]
Ordered list of field API names to retrieve. SOQL does not accept *; every field must be listed explicitly. Standard fields use their documented names; custom fields end with __c. Relationship fields traverse parents via dot notation (Account.Name, Owner.Manager.Email) up to 5 levels deep. Requesting a non-existent or non-queryable field fails at Connect time with a SOQL compile error.
Type: array<string>
# Examples:
columns:
- Id
- Name
- LastModifiedDate
# ---
columns:
- Id
- Account.Name
- Owner.Email
# ---
columns:
- Id
- MyCustom__c
http
HTTP client configuration for Salesforce REST calls (OAuth token endpoint and, where applicable, data queries).
Type: object
http.access_log_body_limit
Maximum bytes of request/response body to include in logs. 0 to skip body logging.
Type: int
Default: 0
http.access_log_level
Log level for HTTP request/response logging. Empty disables logging.
Type: string
Default: ""
Options: "", TRACE, DEBUG, INFO, WARN, ERROR
http.backoff
Adaptive backoff configuration for 429 (Too Many Requests) responses. Always active.
Type: object
http.backoff.initial_interval
Initial interval between retries on 429 responses.
Type: string
Default: 1s
http.backoff.max_interval
Maximum interval between retries on 429 responses.
Type: string
Default: 30s
http.http
HTTP transport settings controlling connection pooling, timeouts, and HTTP/2.
Type: object
http.http.disable_compression
Disable automatic decompression of gzip responses.
Type: bool
Default: false
http.http.disable_keep_alives
Disable HTTP keep-alive connections; each request uses a new connection.
Type: bool
Default: false
http.http.expect_continue_timeout
Maximum time to wait for a server’s 100-continue response before sending the body. 0 means the body is sent immediately.
Type: string
Default: 1s
http.http.h2.max_decoder_header_table_size
Upper limit in bytes for the HPACK header table used to decode headers from the peer. Must be less than 4 MiB.
Type: int
Default: 4096
http.http.h2.max_encoder_header_table_size
Upper limit in bytes for the HPACK header table used to encode headers sent to the peer. Must be less than 4 MiB.
Type: int
Default: 4096
http.http.h2.max_read_frame_size
Largest HTTP/2 frame this endpoint will read. Valid range: 16 KiB to 16 MiB.
Type: int
Default: 16384
http.http.h2.max_receive_buffer_per_connection
Maximum flow-control window size in bytes for data received on a connection. Must be at least 64 KiB and less than 4 MiB.
Type: int
Default: 1048576
http.http.h2.max_receive_buffer_per_stream
Maximum flow-control window size in bytes for data received on a single stream. Must be less than 4 MiB.
Type: int
Default: 1048576
http.http.h2.ping_timeout
Timeout waiting for a PING response before closing the connection.
Type: string
Default: 15s
http.http.h2.send_ping_timeout
Idle timeout after which a PING frame is sent to verify connection health. 0 disables health checks.
Type: string
Default: 0s
http.http.h2.strict_max_concurrent_requests
When true, new requests block when a connection’s concurrency limit is reached instead of opening a new connection.
Type: bool
Default: false
http.http.h2.write_byte_timeout
Timeout for writing data to a connection. The timer resets whenever bytes are written. 0 disables the timeout.
Type: string
Default: 0s
http.http.idle_conn_timeout
How long an idle connection remains in the pool before being closed. 0 disables the timeout.
Type: string
Default: 1m30s
http.http.max_conns_per_host
Maximum total connections (active + idle) per host. 0 means unlimited.
Type: int
Default: 64
http.http.max_idle_conns
Maximum total number of idle (keep-alive) connections across all hosts. 0 means unlimited.
Type: int
Default: 100
http.http.max_idle_conns_per_host
Maximum idle connections to keep per host. 0 (the default) uses GOMAXPROCS+1.
Type: int
Default: 0
http.http.max_response_body_bytes
Maximum bytes of response body the client will read. The response body is wrapped with a limit reader; reads beyond this cap return EOF. 0 disables the limit.
Type: int
Default: 10485760
http.http.max_response_header_bytes
Maximum bytes of response headers to allow.
Type: int
Default: 1048576
http.http.response_header_timeout
Maximum time to wait for response headers after writing the full request. 0 disables the timeout.
Type: string
Default: 0s
http.http.tls_handshake_timeout
Maximum time to wait for a TLS handshake to complete. 0 disables the timeout.
Type: string
Default: 10s
http.http.write_buffer_size
Size in bytes of the per-connection write buffer.
Type: int
Default: 4096
http.tcp
Configure TCP socket-level settings to optimize network performance and reliability. These low-level controls are useful for:
-
Unresponsive hosts: Set
connect_timeoutto limit how long a connection attempt can take (the default0ssets no limit) -
Long-lived connections: Configure
keep_alivesettings to detect and recover from stale connections -
Unstable networks: Tune keep-alive probes to balance between quick failure detection and avoiding false positives
-
Linux systems with specific requirements: Use
tcp_user_timeout(Linux 2.6.37+) to control data acknowledgment timeouts
Most users should keep the default values. Only modify these settings if you’re experiencing connection stability issues or have specific network requirements.
Type: object
http.tcp.connect_timeout
Maximum amount of time a dial will wait for a connect to complete. Zero disables.
Type: string
Default: 0s
http.tcp.keep_alive.count
Maximum unanswered keep-alive probes before dropping the connection. Zero defaults to 9.
Type: int
Default: 9
http.tcp.keep_alive.idle
Duration the connection must be idle before sending the first keep-alive probe. Zero defaults to 15s. Negative values disable keep-alive probes.
Type: string
Default: 15s
http.tcp.keep_alive.interval
Duration between keep-alive probes. Zero defaults to 15s.
Type: string
Default: 15s
http.tcp.tcp_user_timeout
Maximum time to wait for acknowledgment of transmitted data before killing the connection. Linux-only (kernel 2.6.37+), ignored on other platforms. When enabled, keep_alive.idle must be greater than this value per RFC 5482. Zero disables.
Type: string
Default: 0s
http.tls
Configure Transport Layer Security (TLS) settings to secure network connections. This includes options for standard TLS as well as mutual TLS (mTLS) authentication where both client and server authenticate each other using certificates. Key configuration options include enabled to enable TLS, client_certs for mTLS authentication, root_cas/root_cas_file for custom certificate authorities, and skip_cert_verify for development environments.
Type: object
http.tls.client_certs[]
A list of client certificates for mutual TLS (mTLS) authentication. Configure this field to enable mTLS, authenticating the client to the server with these certificates.
You must set tls.enabled: true for the client certificates to take effect.
Certificate pairing rules: For each certificate item, provide either:
-
Inline PEM data using both
certandkeyor -
File paths using both
cert_fileandkey_file.
Mixing inline and file-based values within the same item is not supported.
Type: array<object>
Default: []
# Examples:
client_certs:
- cert: foo
key: bar
# ---
client_certs:
- cert_file: ./example.pem
key_file: ./example.key
http.tls.client_certs[].cert
The plaintext certificate to use for TLS authentication. Must be paired with the corresponding private key in the key field when using inline PEM data for mTLS client certificates.
Type: string
Default: ""
http.tls.client_certs[].cert_file
The path to a file containing the certificate to use for TLS authentication. Must be paired with the corresponding private key file in the key_file field when using file-based configuration for mTLS client certificates.
Type: string
Default: ""
http.tls.client_certs[].key
Private key for mTLS client certificate as inline PEM data. Must correspond to the client certificate specified in the cert field. Use this field together with cert when providing certificate data inline rather than through files.
|
This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets. |
Type: string
Default: ""
http.tls.client_certs[].key_file
Path to private key file for mTLS client certificate in PEM format. Must correspond to the client certificate specified in the cert_file field. Use this field together with cert_file when loading certificate data from files.
Type: string
Default: ""
http.tls.client_certs[].password
The password to use for the private key (specified in the key or key_file fields), if it is password-protected. The PKCS#1 and PKCS#8 formats are supported. Supports environment variable interpolation for secure password management.
The pbeWithMD5AndDES-CBC algorithm is obsolete and not supported for the PKCS#8 format. This algorithm does not authenticate the ciphertext, making it vulnerable to padding oracle attacks that can let an attacker recover the plaintext.
|
This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets. |
Type: string
Default: ""
# Examples:
password: foo
# ---
password: ${KEY_PASSWORD}
http.tls.enable_renegotiation
Whether to allow the remote server to repeatedly request renegotiation. Enable this option if you’re seeing the error message local error: tls: no renegotiation.
Type: bool
Default: false
http.tls.enabled
Whether to enable TLS for secure connections. Set to true to enable TLS encryption. Required to be true for other TLS options (like client_certs, root_cas, etc.) to take effect.
Type: bool
Default: false
http.tls.root_cas
Specify a root certificate authority to use (optional). This is a string that represents a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for inline certificate data or root_cas_file for file-based certificate loading.
|
This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets. |
Type: string
Default: ""
# Examples:
root_cas: |-
-----BEGIN CERTIFICATE-----
...
-----END CERTIFICATE-----
http.tls.root_cas_file
Specify the path to a root certificate authority file (optional). This is a file, often with a .pem extension, which contains a certificate chain from the parent-trusted root certificate, through possible intermediate signing certificates, to the host certificate. Use either this field for file-based certificate loading or root_cas for inline certificate data.
Type: string
Default: ""
# Examples:
root_cas_file: ./root_cas.pem
http.tls.skip_cert_verify
Whether to skip server-side certificate verification. Set to true only for testing environments as this reduces security by disabling certificate validation. When using self-signed certificates or in development, this may be necessary, but should never be used in production. Consider using root_cas or root_cas_file to specify trusted certificates instead of disabling verification entirely.
Type: bool
Default: false
object
The sObject API name to SELECT from. Case-sensitive; uses the API name, not the display label. Standard objects use the noun (Account, Opportunity); custom objects end with __c; Big Objects end with __b; External Objects end with __x. Confirm the exact API name in Setup → Object Manager.
Type: string
# Examples:
object: Account
# ---
object: Contact
# ---
object: MyCustom__c
org_url
Salesforce instance base URL for your org, protocol included and no trailing slash. Used as the base for both the OAuth token endpoint and REST queries. Production orgs use https://{my-domain}.my.salesforce.com; sandboxes use https://{my-domain}.sandbox.my.salesforce.com. Legacy instance URLs (https://na123.salesforce.com) still work but My Domain URLs are strongly recommended by Salesforce.
Type: string
# Examples:
org_url: https://acme.my.salesforce.com
# ---
org_url: https://acme--staging.sandbox.my.salesforce.com
prefix
Optional SOQL fragment inserted before the SELECT keyword. Rarely needed; provided for forward compatibility with future SOQL extensions or Bulk API framing.
Type: string
suffix
Optional SOQL fragment appended after the WHERE clause. Typical uses: ORDER BY for deterministic pagination, LIMIT to cap result size, FOR REFERENCE / FOR VIEW to mark records for Chatter tracking.
Type: string
# Examples:
suffix: ORDER BY LastModifiedDate DESC
# ---
suffix: ORDER BY Id LIMIT 1000
# ---
suffix: ORDER BY CreatedDate DESC LIMIT 10000
where
Optional SOQL WHERE body, without the WHERE keyword. ? placeholders are substituted client-side from args_mapping with SOQL literal escaping (quoted strings, ISO-8601 datetimes). Supports the full WHERE grammar: AND/OR/NOT, LIKE, IN, date literals (TODAY, LAST_N_DAYS:7), subqueries. Date/datetime comparisons require ISO-8601 with explicit timezone.
Type: string
# Examples:
where: LastModifiedDate > ?
# ---
where: Status__c = ? AND CreatedDate > ?
# ---
where: OwnerId IN (?, ?)
Examples
Recent Account Changes
Query Accounts modified in the last hour and emit one message per record:
input:
salesforce:
org_url: https://acme.my.salesforce.com
client_id: ${SALESFORCE_CLIENT_ID}
client_secret: ${SALESFORCE_CLIENT_SECRET}
object: Account
columns: [ Id, Name, LastModifiedDate ]
where: LastModifiedDate > ?
args_mapping: |
root = [ (now() - "1h").ts_format("2006-01-02T15:04:05Z") ]