Keyboard shortcuts

Press ← or → to navigate between chapters

Press S or / to search in the book

Press ? to show this help

Press Esc to hide this help

Connect plugin: outputs

The 68 Redpanda Connect outputs the Connect plugin links, generated from the component specs of Redpanda Connect v4.110.0. Field descriptions are Redpanda Connect’s own, shortened to one sentence; follow Upstream documentation for the rest. Install and configuration forms are on the plugin page.

Use one as a route output:

output:
  custom:
    name: connect
    config:
      connector: mqtt          # the component name below
      urls: ["tcp://localhost:1883"]
      topic: "orders"

amqp_0_9

Sends messages to an AMQP (0.91) exchange. AMQP is a messaging protocol used by various message brokers, including RabbitMQ.

connector: amqp_0_9 · URI connect+amqp-0-9://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
exchangestringrequiredAn AMQP exchange to publish to.
keystring""The binding key to set for each message.
typestring""The type property to set for each message.
metadataobjectSpecify criteria for which metadata values are attached to messages as headers.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: exchange_declare, content_type, content_encoding, correlation_id, reply_to, expiration, message_id, user_id, app_id, priority, persistent, mandatory, immediate, timeout, tls.

Upstream documentation

amqp_1

Sends messages to an AMQP (1.0) server.

connector: amqp_1 · URI connect+amqp-1://

FieldTypeDefaultDescription
urlslist of stringA list of URLs to connect to.
target_addressstring""The target address to write to.
max_in_flightint64The maximum number of messages to have in flight at a given time.
metadataobjectSpecify criteria for which metadata values are attached to messages as headers.

Advanced: tls, application_properties_map, sasl, content_type, persistent, target_capabilities, message_properties_to, message_properties_message_id, message_properties_correlation_id, message_properties_subject, message_properties_reply_to, message_properties_group_id, message_properties_group_sequence, message_properties_reply_to_group_id, message_properties_user_id, message_properties_content_type, message_properties_content_encoding.

Upstream documentation

arc

Writes data to an Arc database via the msgpack ingestion endpoint.

connector: arc · URI connect+arc://

FieldTypeDefaultDescription
base_urlstringrequiredBase URL of the target service (e.g., https://api.example.com).
timeoutstring"5s"HTTP request timeout.
tokenstringBearer token for authentication.
databasestring"default"The target database name.
measurementstringrequiredThe measurement (table) name.
formatstring"columnar"The payload format.
tags_mappingstringAn optional Bloblang mapping to extract tags from each message.
compressionstring"zstd"Compression algorithm for the request body.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: tls, proxy_url, disable_http2, tps_limit, tps_burst, backoff, tcp, http, access_log_level, access_log_body_limit, timestamp_field, timestamp_unit.

Upstream documentation

aws_kinesis

Sends messages to a Kinesis stream.

connector: aws_kinesis · URI connect+aws-kinesis://

FieldTypeDefaultDescription
streamstringrequiredThe stream to publish messages to.
partition_keystringrequiredA required key for partitioning messages.
max_in_flightint64The maximum number of parallel message batches to have in flight at any given time.
batchingobjectAllows you to configure a batching policy.

Advanced: hash_key, region, endpoint, tcp, credentials, max_retries, backoff.

Upstream documentation

aws_kinesis_firehose

Sends messages to a Kinesis Firehose delivery stream.

connector: aws_kinesis_firehose · URI connect+aws-kinesis-firehose://

FieldTypeDefaultDescription
streamstringrequiredThe stream to publish messages to.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: region, endpoint, tcp, credentials, max_retries, backoff.

Upstream documentation

aws_s3

Sends message parts as objects to an Amazon S3 bucket. Each object is uploaded with the path specified with the path field.

connector: aws_s3 · URI connect+aws-s3://

FieldTypeDefaultDescription
bucketstringrequiredThe bucket to upload messages to.
pathstring"${!counter()}-${!timestamp_unix_nano…The path of each message to upload.
tagsmap of string{}Key/value pairs to store with the object as tags.
content_typestring"application/octet-stream"The content type to set for each object.
metadataobjectSpecify criteria for which metadata values are attached to objects as headers.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: content_encoding, cache_control, content_disposition, content_language, website_redirect_location, storage_class, kms_key_id, checksum_algorithm, server_side_encryption, force_path_style_urls, timeout, object_canned_acl, region, endpoint, tcp, credentials.

Upstream documentation

aws_sns

Sends messages to an AWS SNS topic.

connector: aws_sns · URI connect+aws-sns://

FieldTypeDefaultDescription
topic_arnstringrequiredThe topic to publish to.
message_group_idstringAn optional group ID to set for messages.
message_deduplication_idstringAn optional deduplication ID to set for messages.
subjectstringAn optional subject to set for messages.
max_in_flightint64The maximum number of messages to have in flight at a given time.
metadataobjectSpecify criteria for which metadata values are sent as headers.

Advanced: timeout, region, endpoint, tcp, credentials.

Upstream documentation

aws_sqs

Sends messages to an SQS queue.

connector: aws_sqs · URI connect+aws-sqs://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target SQS queue.
message_group_idstringAn optional group ID to set for messages.
message_deduplication_idstringAn optional deduplication ID to set for messages.
delay_secondsstringAn optional delay time in seconds for message.
max_in_flightint64The maximum number of parallel message batches to have in flight at any given time.
metadataobjectSpecify criteria for which metadata values are sent as headers.
batchingobjectAllows you to configure a batching policy.

Advanced: max_records_per_request, region, endpoint, tcp, credentials, max_retries, backoff.

Upstream documentation

azure_blob_storage

Sends message parts as objects to an Azure Blob Storage Account container. Each object is uploaded with the filename specified with the container field.

connector: azure_blob_storage · URI connect+azure-blob-storage://

FieldTypeDefaultDescription
storage_accountstring""The storage account to access.
storage_access_keystring""The storage account access key.
storage_connection_stringstring""A storage account connection string.
storage_sas_tokenstring""The storage account SAS token.
containerstringrequiredThe container for uploading the messages to.
pathstring"${!counter()}-${!timestamp_unix_nano…The path of each message to upload.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: blob_type, public_access_level.

Upstream documentation

azure_cosmosdb

Creates or updates messages as JSON documents in Azure CosmosDB.

connector: azure_cosmosdb · URI connect+azure-cosmosdb://

FieldTypeDefaultDescription
endpointstringCosmosDB endpoint.
account_keystringAccount key.
connection_stringstringConnection string.
databasestringrequiredDatabase.
containerstringrequiredContainer.
partition_keys_mapstringrequiredA Bloblang mapping which should evaluate to a single partition key value or an array of partition key values of type string, integer or boolean.
operationstring"Create"Operation.
item_idstringID of item to replace or delete.
batchingobjectAllows you to configure a batching policy.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: patch_operations, patch_condition, auto_id.

Upstream documentation

azure_data_lake_gen2

Sends message parts as files to an Azure Data Lake Gen2 filesystem. Each file is uploaded with the filename specified with the path field.

connector: azure_data_lake_gen2 · URI connect+azure-data-lake-gen2://

FieldTypeDefaultDescription
storage_accountstring""The storage account to access.
storage_access_keystring""The storage account access key.
storage_connection_stringstring""A storage account connection string.
storage_sas_tokenstring""The storage account SAS token.
filesystemstringrequiredThe data lake storage filesystem name for uploading the messages to.
pathstring"${!counter()}-${!timestamp_unix_nano…The path of each message to upload within the filesystem.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Upstream documentation

azure_queue_storage

Sends messages to an Azure Storage Queue.

connector: azure_queue_storage · URI connect+azure-queue-storage://

FieldTypeDefaultDescription
storage_accountstring""The storage account to access.
storage_access_keystring""The storage account access key.
storage_connection_stringstring""A storage account connection string.
storage_sas_tokenstring""The storage account SAS token.
queue_namestringrequiredThe name of the target Queue Storage queue.
max_in_flightint64The maximum number of parallel message batches to have in flight at any given time.
batchingobjectAllows you to configure a batching policy.

Advanced: ttl.

Upstream documentation

azure_table_storage

Stores messages in an Azure Table Storage table.

connector: azure_table_storage · URI connect+azure-table-storage://

FieldTypeDefaultDescription
storage_accountstring""The storage account to access.
storage_access_keystring""The storage account access key.
storage_connection_stringstring""A storage account connection string.
storage_sas_tokenstring""The storage account SAS token.
table_namestringrequiredThe table to store messages into.
partition_keystring""The partition key.
row_keystring""The row key.
propertiesmap of string{}A map of properties to store into the table.
max_in_flightint64The maximum number of parallel message batches to have in flight at any given time.
batchingobjectAllows you to configure a batching policy.

Advanced: transaction_type, timeout.

Upstream documentation

beanstalkd

Write messages to a Beanstalkd queue.

connector: beanstalkd · URI connect+beanstalkd://

FieldTypeDefaultDescription
addressstringrequiredAn address to connect to.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Upstream documentation

broker

Allows you to route messages to multiple child outputs using a range of brokering patterns.

connector: broker · URI connect+broker://

FieldTypeDefaultDescription
patternstring"fan_out"The brokering pattern to use.
outputslist of outputrequiredA list of child outputs to broker.
batchingobjectAllows you to configure a batching policy.

Advanced: copies.

Upstream documentation

cache

Stores each message in a cache.

connector: cache · URI connect+cache://

FieldTypeDefaultDescription
targetstringrequiredThe target cache to store messages in.
keystring"${!count(\"items\")}-${!timestamp_un…The key to store messages by, function interpolation should be used in order to derive a unique key for each message.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: ttl.

Upstream documentation

cassandra

Runs a query against a Cassandra database for each message in order to insert data.

connector: cassandra · URI connect+cassandra://

FieldTypeDefaultDescription
addresseslist of stringrequiredA list of Cassandra nodes to connect to.
timeoutstring"600ms"The client connection timeout.
reconnect_intervalstring"60s"Attempts to reconnect known DOWN nodes in every ReconnectInterval.
querystringrequiredA query to execute for each message.
args_mappingstringA Bloblang mapping that can be used to provide arguments to Cassandra queries.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: tls, password_authenticator, disable_initial_host_lookup, max_retries, backoff, host_selection_policy, exponential_reconnection, consistency, logged_batch.

Upstream documentation

couchbase

Performs operations against Couchbase for each message, allowing you to store or delete data.

connector: couchbase · URI connect+couchbase://

FieldTypeDefaultDescription
urlstringrequiredCouchbase connection string.
usernamestringUsername to connect to the cluster.
passwordstringPassword to connect to the cluster.
bucketstringrequiredCouchbase bucket.
idstringrequiredDocument id.
contentstringDocument content.
operationstring"upsert"Couchbase operation to perform.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: collection, scope, transcoder, timeout, ttl.

Upstream documentation

cyborgdb

Inserts items into a CyborgDB encrypted vector index.

connector: cyborgdb · URI connect+cyborgdb://

FieldTypeDefaultDescription
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.
hoststringrequiredThe host for the CyborgDB instance.
api_keystringrequiredThe CyborgDB API key for authentication.
index_namestring"redpanda-vectors"The name of the index to write to.
index_keystringrequiredThe base64-encoded encryption key for the index.
operationstring"upsert"The operation to perform against the CyborgDB index.
idstringrequiredThe ID for the vector entry in CyborgDB.
vector_mappingstringThe mapping to extract out the vector from the document.
metadata_mappingstringAn optional mapping of message to metadata for the vector entry.

Advanced: create_if_missing.

Upstream documentation

cypher

The cypher output type writes a batch of messages to any graph database that supports the Neo4j or Bolt protocols.

connector: cypher · URI connect+cypher://

FieldTypeDefaultDescription
uristringrequiredThe connection URI to connect to.
cypherstringrequiredThe cypher expression to execute against the graph database.
database_namestring""Set the target database for which expressions are evaluated against.
args_mappingstringThe mapping from the message to the data that is passed in as parameters to the cypher expression.
basic_authobjectAllows you to specify basic authentication.
batchingobjectAllows you to configure a batching policy.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: tls.

Upstream documentation

discord

Writes messages to a Discord channel.

connector: discord · URI connect+discord://

FieldTypeDefaultDescription
channel_idstringrequiredA discord channel ID to write messages to.
bot_tokenstringrequiredA bot token used for authentication.

Upstream documentation

doris_stream_load

Beta. Writes batches of messages into Apache Doris using Stream Load.

connector: doris_stream_load · URI connect+doris-stream-load://

FieldTypeDefaultDescription
fe_urlslist of string[]A list of Doris FE HTTP URLs.
databasestringrequiredTarget Doris database.
tablestringrequiredTarget Doris table.
usernamestringrequiredDoris username.
passwordstringrequiredDoris password.
formatstring"json"Body format sent to Doris Stream Load.
read_json_by_linebooltrueEncode a JSON batch as newline-delimited JSON and set the Doris header read_json_by_line=true.
strip_outer_arrayboolfalseEncode a JSON batch as a single JSON array and set the Doris header strip_outer_array=true.
columnslist of string[]Optional Doris columns header.
label_prefixstring"redpanda_connect"Prefix used when generating Doris labels.
timeoutstring"30s"Timeout for each Doris HTTP request.
max_in_flightint8Maximum number of parallel in-flight Doris Stream Load requests.
batchingobjectAllows you to configure a batching policy.

Advanced: url, query_port, jsonpaths, json_root, where, column_separator, line_delimiter, group_commit, max_filter_ratio, partitions, temporary_partitions, skip_lines, empty_field_as_null, trim_double_quotes, num_as_string, fuzzy_parse, headers, strict_mode.

Upstream documentation

drop

Drops all messages.

connector: drop · URI connect+drop://

Upstream documentation

drop_on

Attempts to write messages to a child output and if the write fails for one of a list of configurable reasons the message is dropped (acked) instead of being reattempted (or nacked).

connector: drop_on · URI connect+drop-on://

FieldTypeDefaultDescription
errorboolfalseWhether messages should be dropped when the child output returns an error of any type.
error_patternslist of stringA list of regular expressions (re2) where if the child output returns an error that matches any part of any of these patterns the message will be dropped.
back_pressurestringAn optional duration string that determines the maximum length of time to wait for a given message to be accepted by the child output before the message should be dropped instead.
outputoutputrequiredA child output to wrap with this drop mechanism.

Upstream documentation

dynamic

A special broker type where the outputs are identified by unique labels and can be created, changed and removed during runtime via a REST API.

connector: dynamic · URI connect+dynamic://

FieldTypeDefaultDescription
outputsmap of output{}A map of outputs to statically create.
prefixstring""A path prefix for HTTP endpoints that are registered.

Upstream documentation

elasticsearch_v8

Publishes messages into an Elasticsearch index. If the index does not exist then it is created with a dynamic mapping.

connector: elasticsearch_v8 · URI connect+elasticsearch-v8://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
indexstringrequiredThe index to place messages.
actionstringrequiredThe action to take on the document.
idstringrequiredThe ID for indexed messages.
max_in_flightint64The maximum number of messages to have in flight at a given time.
api_keystring""An API key to authenticate with.
batchingobjectAllows you to configure a batching policy.

Advanced: pipeline, routing, retry_on_conflict, tls, basic_auth.

Upstream documentation

elasticsearch_v9

Publishes messages into an Elasticsearch index. If the index does not exist then it is created with a dynamic mapping.

connector: elasticsearch_v9 · URI connect+elasticsearch-v9://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
indexstringrequiredThe index to place messages.
actionstringrequiredThe action to take on the document.
idstringrequiredThe ID for indexed messages.
max_in_flightint64The maximum number of messages to have in flight at a given time.
api_keystring""An API key to authenticate with.
batchingobjectAllows you to configure a batching policy.

Advanced: pipeline, routing, retry_on_conflict, tls, basic_auth.

Upstream documentation

fallback

Attempts to send each message to a child output, starting from the first output on the list. If an output attempt fails then the next output in the list is attempted, and so on.

connector: fallback · URI connect+fallback://

Takes a single list of output value, not a map of fields.

Upstream documentation

file

Writes messages to files on disk based on a chosen codec.

connector: file · URI connect+file://

FieldTypeDefaultDescription
pathstringrequiredThe file to write to, if the file does not yet exist it will be created.
codecstring"lines"The way in which the bytes of messages should be written out into the output data stream.

Upstream documentation

gcp_bigquery

Sends messages as new rows to a Google Cloud BigQuery table.

connector: gcp_bigquery · URI connect+gcp-bigquery://

FieldTypeDefaultDescription
projectstring""The project ID of the dataset to insert data to.
job_projectstring""The project ID in which jobs will be executed.
datasetstringrequiredThe BigQuery Dataset ID.
tablestringrequiredThe table to insert messages to.
formatstring"NEWLINE_DELIMITED_JSON"The format of each incoming message.
max_in_flightint64The maximum number of message batches to have in flight at a given time.
job_labelsmap of string{}A list of labels to add to the load job.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
csvobjectSpecify how CSV data should be interpreted.
batchingobjectAllows you to configure a batching policy.

Advanced: write_disposition, create_disposition, ignore_unknown_values, max_bad_records, auto_detect.

Upstream documentation

gcp_cloud_storage

Sends message parts as objects to a Google Cloud Storage bucket. Each object is uploaded with the path specified with the path field.

connector: gcp_cloud_storage · URI connect+gcp-cloud-storage://

FieldTypeDefaultDescription
bucketstringrequiredThe bucket to upload messages to.
pathstring"${!counter()}-${!timestamp_unix_nano…The path of each message to upload.
content_typestring"application/octet-stream"The content type to set for each object.
collision_modestring"overwrite"Determines how file path collisions should be dealt with.
timeoutstring"3s"The maximum period to wait on an upload before abandoning it and reattempting.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
max_in_flightint64The maximum number of message batches to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: content_encoding, chunk_size.

Upstream documentation

gcp_pubsub

Sends messages to a GCP Cloud Pub/Sub topic. Metadata from messages are sent as attributes.

connector: gcp_pubsub · URI connect+gcp-pubsub://

FieldTypeDefaultDescription
projectstringrequiredThe project ID of the topic to publish to.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
topicstringrequiredThe topic to publish to.
endpointstring""An optional endpoint to override the default of pubsub.googleapis.com:443.
max_in_flightint64The maximum number of messages to have in flight at a given time.
count_thresholdint100Publish a pubsub buffer when it has this many messages
delay_thresholdstring"10ms"Publish a non-empty pubsub buffer after this delay has passed.
byte_thresholdint1000000Publish a batch when its size in bytes reaches this value.
metadataobjectSpecify criteria for which metadata values are sent as attributes, all are sent by default.
batchingobjectConfigures a batching policy on this output.

Advanced: ordering_key, publish_timeout, validate_topic, flow_control.

Upstream documentation

hdfs

Sends message parts as files to a HDFS directory.

connector: hdfs · URI connect+hdfs://

FieldTypeDefaultDescription
hostslist of stringrequiredA list of target host addresses to connect to.
userstring""A user ID to connect as.
directorystringrequiredA directory to store message files within.
pathstring"${!counter()}-${!timestamp_unix_nano…The path to upload messages as, interpolation functions should be used in order to generate unique file paths.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Upstream documentation

http_client

Sends messages to an HTTP server.

connector: http_client · URI connect+http-client://

FieldTypeDefaultDescription
urlstringrequiredThe URL to connect to.
verbstring"POST"A verb to connect with
headersmap of string{}A map of headers to add to the request.
rate_limitstringAn optional rate limit to throttle requests by.
timeoutstring"5s"A static timeout to apply to requests.
max_in_flightint64The maximum number of parallel message batches to have in flight at any given time.
batchingobjectAllows you to configure a batching policy.

Advanced: metadata, dump_request_log_level, oauth, oauth2, basic_auth, jwt, tls, extract_headers, retry_period, max_retry_backoff, retries, follow_redirects, backoff_on, drop_on, successful_on, proxy_url, disable_http2, batch_as_multipart, propagate_response, multipart.

Upstream documentation

http_server

Sets up an HTTP server that will send messages over HTTP(S) GET requests. HTTP 2.0 is supported when using TLS, which is enabled when key and cert files are specified.

connector: http_server · URI connect+http-server://

FieldTypeDefaultDescription
addressstring""An alternative address to host from.
pathstring"/get"The path from which discrete messages can be consumed.
stream_pathstring"/get/stream"The path from which a continuous stream of messages can be consumed.
ws_pathstring"/get/ws"The path from which websocket connections can be established.
allowed_verbslist of string["GET"]An array of verbs that are allowed for the path and stream_path HTTP endpoint.

Advanced: timeout, cert_file, key_file, cors.

Upstream documentation

inproc

Sends data directly to Redpanda Connect inputs by connecting to a unique ID.

connector: inproc · URI connect+inproc://

Takes a single string value, not a map of fields.

Upstream documentation

mongodb

Inserts items into a MongoDB collection.

connector: mongodb · URI connect+mongodb://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target MongoDB server.
databasestringrequiredThe name of the target MongoDB database.
usernamestring""The username to connect to the database.
passwordstring""The password to connect to the database.
collectionstringrequiredThe name of the target collection.
operationstring"update-one"The mongodb operation to perform.
write_concernobjectThe write concern settings for the mongo connection.
document_mapstring""A bloblang map representing a document to store within MongoDB, expressed as extended JSON in canonical form.
filter_mapstring""A bloblang map representing a filter for a MongoDB command, expressed as extended JSON in canonical form.
hint_mapstring""A bloblang map representing the hint for the MongoDB command, expressed as extended JSON in canonical form.
upsertboolfalseThe upsert setting is optional and only applies for update-one and replace-one operations.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: app_name, aws.

Upstream documentation

mqtt

Pushes messages to an MQTT broker.

connector: mqtt · URI connect+mqtt://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
client_idstring""An identifier for the client connection.
connect_timeoutstring"30s"The maximum amount of time to wait in order to establish a connection before the attempt is abandoned.
topicstringrequiredThe topic to publish messages to.
qosint1The QoS value to set for each message.
write_timeoutstring"3s"The maximum amount of time to wait to write data before the attempt is abandoned.
retainedboolfalseSet message as retained on the topic.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: dynamic_client_id_suffix, will, user, password, keepalive, tls, retained_interpolated.

Upstream documentation

nanomsg

Send messages over a Nanomsg socket.

connector: nanomsg · URI connect+nanomsg://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
bindboolfalseWhether the URLs listed should be bind (otherwise they are connected to).
socket_typestring"PUSH"The socket type to send with.
poll_timeoutstring"5s"The maximum period of time to wait for a message to send before the request is abandoned and reattempted.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Upstream documentation

nats

Publish to an NATS subject.

connector: nats · URI connect+nats://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
subjectstringrequiredThe subject to publish to.
headersmap of string{}Explicit message headers to add to messages.
metadataobjectDetermine which (if any) metadata values should be added to messages as headers.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.

Upstream documentation

nats_jetstream

Write messages to a NATS JetStream subject.

connector: nats_jetstream · URI connect+nats-jetstream://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
subjectstringrequiredA subject to write to.
headersmap of string{}Explicit message headers to add to messages.
metadataobjectDetermine which (if any) metadata values should be added to messages as headers.
max_in_flightint1024The maximum number of messages to have in flight at a given time.

Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.

Upstream documentation

nats_kv

Put messages in a NATS key-value bucket.

connector: nats_kv · URI connect+nats-kv://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
bucketstringrequiredThe name of the KV bucket.
keystringrequiredThe key for each message.
max_in_flightint1024The maximum number of messages to have in flight at a given time.

Advanced: max_reconnects, tls, tls_handshake_first, auth.

Upstream documentation

nats_stream

Publish to a NATS Stream subject.

connector: nats_stream · URI connect+nats-stream://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
cluster_idstringrequiredThe cluster ID to publish to.
subjectstringrequiredThe subject to publish to.
client_idstring""The client ID to connect with.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.

Upstream documentation

nsq

Publish to an NSQ topic.

connector: nsq · URI connect+nsq://

FieldTypeDefaultDescription
nsqd_tcp_addressstringrequiredThe address of the target NSQD server.
topicstringrequiredThe topic to publish to.
user_agentstringA user agent to assume when connecting.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: tls.

Upstream documentation

opensearch

Publishes messages into an Elasticsearch index. If the index does not exist then it is created with a dynamic mapping.

connector: opensearch · URI connect+opensearch://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
indexstringrequiredThe index to place messages.
actionstringrequiredThe action to take on the document.
idstringrequiredThe ID for indexed messages.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: pipeline, routing, tls, basic_auth, aws.

Upstream documentation

pinecone

Inserts items into a Pinecone index.

connector: pinecone · URI connect+pinecone://

FieldTypeDefaultDescription
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.
hoststringrequiredThe host for the Pinecone index.
api_keystringrequiredThe Pinecone api key.
operationstring"upsert-vectors"The operation to perform against the Pinecone index.
idstringrequiredThe ID for the index entry in Pinecone.
vector_mappingstringThe mapping to extract out the vector from the document.
metadata_mappingstringAn optional mapping of message to metadata in the Pinecone index entry.

Advanced: namespace.

Upstream documentation

pulsar

Write messages to an Apache Pulsar server.

connector: pulsar · URI connect+pulsar://

FieldTypeDefaultDescription
urlstringrequiredA URL to connect to.
topicstringrequiredThe topic to publish to.
tlsobjectSpecify the path to a custom CA certificate to trust broker TLS service.
keystring""The key to publish messages with.
ordering_keystring""The ordering key to publish messages with.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: auth.

Upstream documentation

pusher

Output for publishing messages to Pusher API (https://pusher.com)

connector: pusher · URI connect+pusher://

FieldTypeDefaultDescription
batchingobjectmaximum batch size is 10 (limit of the pusher library)
channelstringrequiredPusher channel to publish to.
eventstringrequiredEvent to publish to
appIdstringrequiredPusher app id
keystringrequiredPusher key
secretstringrequiredPusher secret
clusterstringrequiredPusher cluster
securebooltrueEnable SSL encryption
max_in_flightint1The maximum number of parallel message batches to have in flight at any given time.

Upstream documentation

qdrant

Adds items to a Qdrant collection

connector: qdrant · URI connect+qdrant://

FieldTypeDefaultDescription
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.
grpc_hoststringrequiredThe gRPC host of the Qdrant server.
api_tokenstring""The Qdrant API token for authentication.
collection_namestringrequiredThe name of the collection in Qdrant.
idstringrequiredThe ID of the point to insert.
vector_mappingstringrequiredThe mapping to extract the vector from the document.
payload_mappingstring"root = {}"An optional mapping of message to payload associated with the point.

Advanced: tls.

Upstream documentation

questdb

Pushes messages to a QuestDB table

connector: questdb · URI connect+questdb://

FieldTypeDefaultDescription
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.
addressstringrequiredAddress of the QuestDB server’s HTTP port (excluding protocol)
usernamestringUsername for HTTP basic auth
passwordstringPassword for HTTP basic auth
tokenstringBearer token for HTTP auth (takes precedence over basic auth username & password)
tablestringrequiredDestination table
designated_timestamp_fieldstringName of the designated timestamp field
designated_timestamp_unitstring"auto"Designated timestamp field units
timestamp_string_fieldslist of stringString fields with textual timestamps
timestamp_string_formatstring"Jan _2 15:04:05.000000Z0700"Timestamp format, used when parsing timestamp string fields.
symbolslist of stringColumns that should be the SYMBOL type (string values default to STRING)
doubleslist of stringColumns that should be double type, (int is default)
error_on_empty_messagesboolfalseMark a message as errored if it is empty after field validation

Advanced: tls, retry_timeout, request_timeout, request_min_throughput.

Upstream documentation

redis_hash

Sets Redis hash objects using the HSET command.

connector: redis_hash · URI connect+redis-hash://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
keystringrequiredThe key for each message, function interpolations should be used to create a unique key per message.
walk_metadataboolfalseWhether all metadata fields of messages should be walked and added to the list of hash fields to set.
walk_json_objectboolfalseWhether to walk each message as a JSON object and add each key/value pair to the list of hash fields to set.
fieldsmap of string{}A map of key/value pairs to set as hash fields.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: kind, master, client_name, tls.

Upstream documentation

redis_list

Pushes messages onto the end of a Redis list (which is created if it doesn’t already exist) using the RPUSH command.

connector: redis_list · URI connect+redis-list://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
keystringrequiredThe key for each message, function interpolations can be optionally used to create a unique key per message.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: kind, master, client_name, tls, command.

Upstream documentation

redis_pubsub

Publishes messages through the Redis PubSub model. It is not possible to guarantee that messages have been received.

connector: redis_pubsub · URI connect+redis-pubsub://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
channelstringrequiredThe channel to publish messages to.
max_in_flightint64The maximum number of messages to have in flight at a given time.
batchingobjectAllows you to configure a batching policy.

Advanced: kind, master, client_name, tls.

Upstream documentation

redis_streams

Pushes messages to a Redis (v5.0+) Stream (which is created if it doesn’t already exist) using the XADD command.

connector: redis_streams · URI connect+redis-streams://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
streamstringrequiredThe stream to add messages to.
idstring"*"The entry ID for the stream message.
body_keystring"body"A key to set the raw body of the message to.
max_lengthint0When greater than zero enforces a rough cap on the length of the target stream.
max_in_flightint64The maximum number of messages to have in flight at a given time.
metadataobjectSpecify criteria for which metadata values are included in the message body.
batchingobjectAllows you to configure a batching policy.

Advanced: kind, master, client_name, tls.

Upstream documentation

reject

Rejects all messages, treating them as though the output destination failed to publish them.

connector: reject · URI connect+reject://

Takes a single string value, not a map of fields.

Upstream documentation

reject_errored

Rejects messages that have failed their processing steps, resulting in nack behavior at the input level, otherwise sends them to a child output.

connector: reject_errored · URI connect+reject-errored://

Takes a single output value, not a map of fields.

Upstream documentation

resource

Resource is an output type that channels messages to a resource output, identified by its name.

connector: resource · URI connect+resource://

Takes a single string value, not a map of fields.

Upstream documentation

retry

Attempts to write messages to a child output and if the write fails for any reason the message is retried either until success or, if the retries or max elapsed time fields are non-zero, either is reached.

connector: retry · URI connect+retry://

FieldTypeDefaultDescription
outputoutputrequiredA child output.

Advanced: max_retries, backoff.

Upstream documentation

sftp

Writes files to an SFTP server.

connector: sftp · URI connect+sftp://

FieldTypeDefaultDescription
addressstringrequiredThe address of the server to connect to.
credentialsobjectThe credentials to use to log into the target server.
pathstringrequiredThe file to save the messages to on the server.
codecstring"all-bytes"The way in which the bytes of messages should be written out into the output data stream.
max_in_flightint64The maximum number of messages to have in flight at a given time.

Advanced: connection_timeout.

Upstream documentation

socket

Connects to a (tcp/udp/unix) server and sends a continuous stream of data, dividing messages according to the specified codec.

connector: socket · URI connect+socket://

FieldTypeDefaultDescription
networkstringrequiredA network type to connect as.
addressstringrequiredThe address to connect to.
codecstring"lines"The way in which the bytes of messages should be written out into the output data stream.

Advanced: tls.

Upstream documentation

sql

Deprecated. Executes an arbitrary SQL query for each message.

connector: sql · URI connect+sql://

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
data_source_namestringrequiredData source name.
querystringrequiredThe query to execute.
args_mappingstringAn optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query.
max_in_flightint64The maximum number of inserts to run in parallel.
batchingobjectAllows you to configure a batching policy.

Upstream documentation

sql_insert

Inserts a row into an SQL database for each message.

connector: sql_insert · URI connect+sql-insert://

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
tablestringrequiredThe table to insert to.
columnslist of stringrequiredA list of columns to insert.
args_mappingstringrequiredA Bloblang mapping which should evaluate to an array of values matching in size to the number of columns specified.
max_in_flightint64The maximum number of inserts to run in parallel.
batchingobjectAllows you to configure a batching policy.

Advanced: prefix, suffix, options, init_files, init_statement, conn_max_idle_time, conn_max_life_time, conn_max_idle, conn_max_open.

Upstream documentation

sql_raw

Executes an arbitrary SQL query for each message.

connector: sql_raw · URI connect+sql-raw://

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
querystringThe query to execute.
args_mappingstringAn optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query.
querieslist of objectA list of query statements.
max_in_flightint64The maximum number of batches to be sending in parallel at any given time.
batchingobjectAllows you to configure a batching policy.

Advanced: unsafe_dynamic_query, init_files, init_statement, conn_max_idle_time, conn_max_life_time, conn_max_idle, conn_max_open.

Upstream documentation

stdout

Prints messages to stdout as a continuous stream of data.

connector: stdout · URI connect+stdout://

FieldTypeDefaultDescription
codecstring"lines"The way in which the bytes of messages should be written out into the output data stream.

Upstream documentation

subprocess

Beta. Executes a command, runs it as a subprocess, and writes messages to it over stdin.

connector: subprocess · URI connect+subprocess://

FieldTypeDefaultDescription
namestringrequiredThe command to execute as a subprocess.
argslist of string[]A list of arguments to provide the command.
codecstring"lines"The way in which messages should be written to the subprocess.

Upstream documentation

switch

The switch output type allows you to route messages to different outputs based on their contents.

connector: switch · URI connect+switch://

FieldTypeDefaultDescription
retry_until_successboolfalseIf a selected output fails to send a message this field determines whether it is reattempted indefinitely.
caseslist of objectA list of switch cases, outlining outputs that can be routed to.

Advanced: strict_mode.

Upstream documentation

sync_response

Returns the final message payload back to the input origin of the message, where it is dealt with according to that specific input type.

connector: sync_response · URI connect+sync-response://

Upstream documentation

websocket

Sends messages to an HTTP server via a websocket connection.

connector: websocket · URI connect+websocket://

FieldTypeDefaultDescription
urlstringrequiredThe URL to connect to.

Advanced: proxy_url, tls, oauth, basic_auth, jwt.

Upstream documentation