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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
exchange | string | required | An AMQP exchange to publish to. |
key | string | "" | The binding key to set for each message. |
type | string | "" | The type property to set for each message. |
metadata | object | Specify criteria for which metadata values are attached to messages as headers. | |
max_in_flight | int | 64 | The 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.
amqp_1
Sends messages to an AMQP (1.0) server.
connector: amqp_1 · URI connect+amqp-1://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | A list of URLs to connect to. | |
target_address | string | "" | The target address to write to. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
metadata | object | Specify 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.
arc
Writes data to an Arc database via the msgpack ingestion endpoint.
connector: arc · URI connect+arc://
| Field | Type | Default | Description |
|---|---|---|---|
base_url | string | required | Base URL of the target service (e.g., https://api.example.com). |
timeout | string | "5s" | HTTP request timeout. |
token | string | Bearer token for authentication. | |
database | string | "default" | The target database name. |
measurement | string | required | The measurement (table) name. |
format | string | "columnar" | The payload format. |
tags_mapping | string | An optional Bloblang mapping to extract tags from each message. | |
compression | string | "zstd" | Compression algorithm for the request body. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows 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.
aws_kinesis
Sends messages to a Kinesis stream.
connector: aws_kinesis · URI connect+aws-kinesis://
| Field | Type | Default | Description |
|---|---|---|---|
stream | string | required | The stream to publish messages to. |
partition_key | string | required | A required key for partitioning messages. |
max_in_flight | int | 64 | The maximum number of parallel message batches to have in flight at any given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: hash_key, region, endpoint, tcp, credentials, max_retries, backoff.
aws_kinesis_firehose
Sends messages to a Kinesis Firehose delivery stream.
connector: aws_kinesis_firehose · URI connect+aws-kinesis-firehose://
| Field | Type | Default | Description |
|---|---|---|---|
stream | string | required | The stream to publish messages to. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: region, endpoint, tcp, credentials, max_retries, backoff.
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://
| Field | Type | Default | Description |
|---|---|---|---|
bucket | string | required | The bucket to upload messages to. |
path | string | "${!counter()}-${!timestamp_unix_nano… | The path of each message to upload. |
tags | map of string | {} | Key/value pairs to store with the object as tags. |
content_type | string | "application/octet-stream" | The content type to set for each object. |
metadata | object | Specify criteria for which metadata values are attached to objects as headers. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows 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.
aws_sns
Sends messages to an AWS SNS topic.
connector: aws_sns · URI connect+aws-sns://
| Field | Type | Default | Description |
|---|---|---|---|
topic_arn | string | required | The topic to publish to. |
message_group_id | string | An optional group ID to set for messages. | |
message_deduplication_id | string | An optional deduplication ID to set for messages. | |
subject | string | An optional subject to set for messages. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
metadata | object | Specify criteria for which metadata values are sent as headers. |
Advanced: timeout, region, endpoint, tcp, credentials.
aws_sqs
Sends messages to an SQS queue.
connector: aws_sqs · URI connect+aws-sqs://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target SQS queue. |
message_group_id | string | An optional group ID to set for messages. | |
message_deduplication_id | string | An optional deduplication ID to set for messages. | |
delay_seconds | string | An optional delay time in seconds for message. | |
max_in_flight | int | 64 | The maximum number of parallel message batches to have in flight at any given time. |
metadata | object | Specify criteria for which metadata values are sent as headers. | |
batching | object | Allows you to configure a batching policy. |
Advanced: max_records_per_request, region, endpoint, tcp, credentials, max_retries, backoff.
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://
| Field | Type | Default | Description |
|---|---|---|---|
storage_account | string | "" | The storage account to access. |
storage_access_key | string | "" | The storage account access key. |
storage_connection_string | string | "" | A storage account connection string. |
storage_sas_token | string | "" | The storage account SAS token. |
container | string | required | The container for uploading the messages to. |
path | string | "${!counter()}-${!timestamp_unix_nano… | The path of each message to upload. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: blob_type, public_access_level.
azure_cosmosdb
Creates or updates messages as JSON documents in Azure CosmosDB.
connector: azure_cosmosdb · URI connect+azure-cosmosdb://
| Field | Type | Default | Description |
|---|---|---|---|
endpoint | string | CosmosDB endpoint. | |
account_key | string | Account key. | |
connection_string | string | Connection string. | |
database | string | required | Database. |
container | string | required | Container. |
partition_keys_map | string | required | A Bloblang mapping which should evaluate to a single partition key value or an array of partition key values of type string, integer or boolean. |
operation | string | "Create" | Operation. |
item_id | string | ID of item to replace or delete. | |
batching | object | Allows you to configure a batching policy. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: patch_operations, patch_condition, auto_id.
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://
| Field | Type | Default | Description |
|---|---|---|---|
storage_account | string | "" | The storage account to access. |
storage_access_key | string | "" | The storage account access key. |
storage_connection_string | string | "" | A storage account connection string. |
storage_sas_token | string | "" | The storage account SAS token. |
filesystem | string | required | The data lake storage filesystem name for uploading the messages to. |
path | string | "${!counter()}-${!timestamp_unix_nano… | The path of each message to upload within the filesystem. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
azure_queue_storage
Sends messages to an Azure Storage Queue.
connector: azure_queue_storage · URI connect+azure-queue-storage://
| Field | Type | Default | Description |
|---|---|---|---|
storage_account | string | "" | The storage account to access. |
storage_access_key | string | "" | The storage account access key. |
storage_connection_string | string | "" | A storage account connection string. |
storage_sas_token | string | "" | The storage account SAS token. |
queue_name | string | required | The name of the target Queue Storage queue. |
max_in_flight | int | 64 | The maximum number of parallel message batches to have in flight at any given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: ttl.
azure_table_storage
Stores messages in an Azure Table Storage table.
connector: azure_table_storage · URI connect+azure-table-storage://
| Field | Type | Default | Description |
|---|---|---|---|
storage_account | string | "" | The storage account to access. |
storage_access_key | string | "" | The storage account access key. |
storage_connection_string | string | "" | A storage account connection string. |
storage_sas_token | string | "" | The storage account SAS token. |
table_name | string | required | The table to store messages into. |
partition_key | string | "" | The partition key. |
row_key | string | "" | The row key. |
properties | map of string | {} | A map of properties to store into the table. |
max_in_flight | int | 64 | The maximum number of parallel message batches to have in flight at any given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: transaction_type, timeout.
beanstalkd
Write messages to a Beanstalkd queue.
connector: beanstalkd · URI connect+beanstalkd://
| Field | Type | Default | Description |
|---|---|---|---|
address | string | required | An address to connect to. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
broker
Allows you to route messages to multiple child outputs using a range of brokering patterns.
connector: broker · URI connect+broker://
| Field | Type | Default | Description |
|---|---|---|---|
pattern | string | "fan_out" | The brokering pattern to use. |
outputs | list of output | required | A list of child outputs to broker. |
batching | object | Allows you to configure a batching policy. |
Advanced: copies.
cache
Stores each message in a cache.
connector: cache · URI connect+cache://
| Field | Type | Default | Description |
|---|---|---|---|
target | string | required | The target cache to store messages in. |
key | string | "${!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_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: ttl.
cassandra
Runs a query against a Cassandra database for each message in order to insert data.
connector: cassandra · URI connect+cassandra://
| Field | Type | Default | Description |
|---|---|---|---|
addresses | list of string | required | A list of Cassandra nodes to connect to. |
timeout | string | "600ms" | The client connection timeout. |
reconnect_interval | string | "60s" | Attempts to reconnect known DOWN nodes in every ReconnectInterval. |
query | string | required | A query to execute for each message. |
args_mapping | string | A Bloblang mapping that can be used to provide arguments to Cassandra queries. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows 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.
couchbase
Performs operations against Couchbase for each message, allowing you to store or delete data.
connector: couchbase · URI connect+couchbase://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | Couchbase connection string. |
username | string | Username to connect to the cluster. | |
password | string | Password to connect to the cluster. | |
bucket | string | required | Couchbase bucket. |
id | string | required | Document id. |
content | string | Document content. | |
operation | string | "upsert" | Couchbase operation to perform. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: collection, scope, transcoder, timeout, ttl.
cyborgdb
Inserts items into a CyborgDB encrypted vector index.
connector: cyborgdb · URI connect+cyborgdb://
| Field | Type | Default | Description |
|---|---|---|---|
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. | |
host | string | required | The host for the CyborgDB instance. |
api_key | string | required | The CyborgDB API key for authentication. |
index_name | string | "redpanda-vectors" | The name of the index to write to. |
index_key | string | required | The base64-encoded encryption key for the index. |
operation | string | "upsert" | The operation to perform against the CyborgDB index. |
id | string | required | The ID for the vector entry in CyborgDB. |
vector_mapping | string | The mapping to extract out the vector from the document. | |
metadata_mapping | string | An optional mapping of message to metadata for the vector entry. |
Advanced: create_if_missing.
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://
| Field | Type | Default | Description |
|---|---|---|---|
uri | string | required | The connection URI to connect to. |
cypher | string | required | The cypher expression to execute against the graph database. |
database_name | string | "" | Set the target database for which expressions are evaluated against. |
args_mapping | string | The mapping from the message to the data that is passed in as parameters to the cypher expression. | |
basic_auth | object | Allows you to specify basic authentication. | |
batching | object | Allows you to configure a batching policy. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: tls.
discord
Writes messages to a Discord channel.
connector: discord · URI connect+discord://
| Field | Type | Default | Description |
|---|---|---|---|
channel_id | string | required | A discord channel ID to write messages to. |
bot_token | string | required | A bot token used for authentication. |
doris_stream_load
Beta. Writes batches of messages into Apache Doris using Stream Load.
connector: doris_stream_load · URI connect+doris-stream-load://
| Field | Type | Default | Description |
|---|---|---|---|
fe_urls | list of string | [] | A list of Doris FE HTTP URLs. |
database | string | required | Target Doris database. |
table | string | required | Target Doris table. |
username | string | required | Doris username. |
password | string | required | Doris password. |
format | string | "json" | Body format sent to Doris Stream Load. |
read_json_by_line | bool | true | Encode a JSON batch as newline-delimited JSON and set the Doris header read_json_by_line=true. |
strip_outer_array | bool | false | Encode a JSON batch as a single JSON array and set the Doris header strip_outer_array=true. |
columns | list of string | [] | Optional Doris columns header. |
label_prefix | string | "redpanda_connect" | Prefix used when generating Doris labels. |
timeout | string | "30s" | Timeout for each Doris HTTP request. |
max_in_flight | int | 8 | Maximum number of parallel in-flight Doris Stream Load requests. |
batching | object | Allows 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.
drop
Drops all messages.
connector: drop · URI connect+drop://
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://
| Field | Type | Default | Description |
|---|---|---|---|
error | bool | false | Whether messages should be dropped when the child output returns an error of any type. |
error_patterns | list of string | A 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_pressure | string | An 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. | |
output | output | required | A child output to wrap with this drop mechanism. |
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://
| Field | Type | Default | Description |
|---|---|---|---|
outputs | map of output | {} | A map of outputs to statically create. |
prefix | string | "" | A path prefix for HTTP endpoints that are registered. |
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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
index | string | required | The index to place messages. |
action | string | required | The action to take on the document. |
id | string | required | The ID for indexed messages. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
api_key | string | "" | An API key to authenticate with. |
batching | object | Allows you to configure a batching policy. |
Advanced: pipeline, routing, retry_on_conflict, tls, basic_auth.
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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
index | string | required | The index to place messages. |
action | string | required | The action to take on the document. |
id | string | required | The ID for indexed messages. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
api_key | string | "" | An API key to authenticate with. |
batching | object | Allows you to configure a batching policy. |
Advanced: pipeline, routing, retry_on_conflict, tls, basic_auth.
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.
file
Writes messages to files on disk based on a chosen codec.
connector: file · URI connect+file://
| Field | Type | Default | Description |
|---|---|---|---|
path | string | required | The file to write to, if the file does not yet exist it will be created. |
codec | string | "lines" | The way in which the bytes of messages should be written out into the output data stream. |
gcp_bigquery
Sends messages as new rows to a Google Cloud BigQuery table.
connector: gcp_bigquery · URI connect+gcp-bigquery://
| Field | Type | Default | Description |
|---|---|---|---|
project | string | "" | The project ID of the dataset to insert data to. |
job_project | string | "" | The project ID in which jobs will be executed. |
dataset | string | required | The BigQuery Dataset ID. |
table | string | required | The table to insert messages to. |
format | string | "NEWLINE_DELIMITED_JSON" | The format of each incoming message. |
max_in_flight | int | 64 | The maximum number of message batches to have in flight at a given time. |
job_labels | map of string | {} | A list of labels to add to the load job. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
csv | object | Specify how CSV data should be interpreted. | |
batching | object | Allows you to configure a batching policy. |
Advanced: write_disposition, create_disposition, ignore_unknown_values, max_bad_records, auto_detect.
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://
| Field | Type | Default | Description |
|---|---|---|---|
bucket | string | required | The bucket to upload messages to. |
path | string | "${!counter()}-${!timestamp_unix_nano… | The path of each message to upload. |
content_type | string | "application/octet-stream" | The content type to set for each object. |
collision_mode | string | "overwrite" | Determines how file path collisions should be dealt with. |
timeout | string | "3s" | The maximum period to wait on an upload before abandoning it and reattempting. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
max_in_flight | int | 64 | The maximum number of message batches to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: content_encoding, chunk_size.
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://
| Field | Type | Default | Description |
|---|---|---|---|
project | string | required | The project ID of the topic to publish to. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
topic | string | required | The topic to publish to. |
endpoint | string | "" | An optional endpoint to override the default of pubsub.googleapis.com:443. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
count_threshold | int | 100 | Publish a pubsub buffer when it has this many messages |
delay_threshold | string | "10ms" | Publish a non-empty pubsub buffer after this delay has passed. |
byte_threshold | int | 1000000 | Publish a batch when its size in bytes reaches this value. |
metadata | object | Specify criteria for which metadata values are sent as attributes, all are sent by default. | |
batching | object | Configures a batching policy on this output. |
Advanced: ordering_key, publish_timeout, validate_topic, flow_control.
hdfs
Sends message parts as files to a HDFS directory.
connector: hdfs · URI connect+hdfs://
| Field | Type | Default | Description |
|---|---|---|---|
hosts | list of string | required | A list of target host addresses to connect to. |
user | string | "" | A user ID to connect as. |
directory | string | required | A directory to store message files within. |
path | string | "${!counter()}-${!timestamp_unix_nano… | The path to upload messages as, interpolation functions should be used in order to generate unique file paths. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
http_client
Sends messages to an HTTP server.
connector: http_client · URI connect+http-client://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL to connect to. |
verb | string | "POST" | A verb to connect with |
headers | map of string | {} | A map of headers to add to the request. |
rate_limit | string | An optional rate limit to throttle requests by. | |
timeout | string | "5s" | A static timeout to apply to requests. |
max_in_flight | int | 64 | The maximum number of parallel message batches to have in flight at any given time. |
batching | object | Allows 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
address | string | "" | An alternative address to host from. |
path | string | "/get" | The path from which discrete messages can be consumed. |
stream_path | string | "/get/stream" | The path from which a continuous stream of messages can be consumed. |
ws_path | string | "/get/ws" | The path from which websocket connections can be established. |
allowed_verbs | list 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.
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.
mongodb
Inserts items into a MongoDB collection.
connector: mongodb · URI connect+mongodb://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target MongoDB server. |
database | string | required | The name of the target MongoDB database. |
username | string | "" | The username to connect to the database. |
password | string | "" | The password to connect to the database. |
collection | string | required | The name of the target collection. |
operation | string | "update-one" | The mongodb operation to perform. |
write_concern | object | The write concern settings for the mongo connection. | |
document_map | string | "" | A bloblang map representing a document to store within MongoDB, expressed as extended JSON in canonical form. |
filter_map | string | "" | A bloblang map representing a filter for a MongoDB command, expressed as extended JSON in canonical form. |
hint_map | string | "" | A bloblang map representing the hint for the MongoDB command, expressed as extended JSON in canonical form. |
upsert | bool | false | The upsert setting is optional and only applies for update-one and replace-one operations. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: app_name, aws.
mqtt
Pushes messages to an MQTT broker.
connector: mqtt · URI connect+mqtt://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
client_id | string | "" | An identifier for the client connection. |
connect_timeout | string | "30s" | The maximum amount of time to wait in order to establish a connection before the attempt is abandoned. |
topic | string | required | The topic to publish messages to. |
qos | int | 1 | The QoS value to set for each message. |
write_timeout | string | "3s" | The maximum amount of time to wait to write data before the attempt is abandoned. |
retained | bool | false | Set message as retained on the topic. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: dynamic_client_id_suffix, will, user, password, keepalive, tls, retained_interpolated.
nanomsg
Send messages over a Nanomsg socket.
connector: nanomsg · URI connect+nanomsg://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
bind | bool | false | Whether the URLs listed should be bind (otherwise they are connected to). |
socket_type | string | "PUSH" | The socket type to send with. |
poll_timeout | string | "5s" | The maximum period of time to wait for a message to send before the request is abandoned and reattempted. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
nats
Publish to an NATS subject.
connector: nats · URI connect+nats://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
subject | string | required | The subject to publish to. |
headers | map of string | {} | Explicit message headers to add to messages. |
metadata | object | Determine which (if any) metadata values should be added to messages as headers. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.
nats_jetstream
Write messages to a NATS JetStream subject.
connector: nats_jetstream · URI connect+nats-jetstream://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
subject | string | required | A subject to write to. |
headers | map of string | {} | Explicit message headers to add to messages. |
metadata | object | Determine which (if any) metadata values should be added to messages as headers. | |
max_in_flight | int | 1024 | The maximum number of messages to have in flight at a given time. |
Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.
nats_kv
Put messages in a NATS key-value bucket.
connector: nats_kv · URI connect+nats-kv://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
bucket | string | required | The name of the KV bucket. |
key | string | required | The key for each message. |
max_in_flight | int | 1024 | The maximum number of messages to have in flight at a given time. |
Advanced: max_reconnects, tls, tls_handshake_first, auth.
nats_stream
Publish to a NATS Stream subject.
connector: nats_stream · URI connect+nats-stream://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
cluster_id | string | required | The cluster ID to publish to. |
subject | string | required | The subject to publish to. |
client_id | string | "" | The client ID to connect with. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: max_reconnects, tls, tls_handshake_first, auth, inject_tracing_map.
nsq
Publish to an NSQ topic.
connector: nsq · URI connect+nsq://
| Field | Type | Default | Description |
|---|---|---|---|
nsqd_tcp_address | string | required | The address of the target NSQD server. |
topic | string | required | The topic to publish to. |
user_agent | string | A user agent to assume when connecting. | |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: tls.
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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
index | string | required | The index to place messages. |
action | string | required | The action to take on the document. |
id | string | required | The ID for indexed messages. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: pipeline, routing, tls, basic_auth, aws.
pinecone
Inserts items into a Pinecone index.
connector: pinecone · URI connect+pinecone://
| Field | Type | Default | Description |
|---|---|---|---|
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. | |
host | string | required | The host for the Pinecone index. |
api_key | string | required | The Pinecone api key. |
operation | string | "upsert-vectors" | The operation to perform against the Pinecone index. |
id | string | required | The ID for the index entry in Pinecone. |
vector_mapping | string | The mapping to extract out the vector from the document. | |
metadata_mapping | string | An optional mapping of message to metadata in the Pinecone index entry. |
Advanced: namespace.
pulsar
Write messages to an Apache Pulsar server.
connector: pulsar · URI connect+pulsar://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | A URL to connect to. |
topic | string | required | The topic to publish to. |
tls | object | Specify the path to a custom CA certificate to trust broker TLS service. | |
key | string | "" | The key to publish messages with. |
ordering_key | string | "" | The ordering key to publish messages with. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: auth.
pusher
Output for publishing messages to Pusher API (https://pusher.com)
connector: pusher · URI connect+pusher://
| Field | Type | Default | Description |
|---|---|---|---|
batching | object | maximum batch size is 10 (limit of the pusher library) | |
channel | string | required | Pusher channel to publish to. |
event | string | required | Event to publish to |
appId | string | required | Pusher app id |
key | string | required | Pusher key |
secret | string | required | Pusher secret |
cluster | string | required | Pusher cluster |
secure | bool | true | Enable SSL encryption |
max_in_flight | int | 1 | The maximum number of parallel message batches to have in flight at any given time. |
qdrant
Adds items to a Qdrant collection
connector: qdrant · URI connect+qdrant://
| Field | Type | Default | Description |
|---|---|---|---|
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. | |
grpc_host | string | required | The gRPC host of the Qdrant server. |
api_token | string | "" | The Qdrant API token for authentication. |
collection_name | string | required | The name of the collection in Qdrant. |
id | string | required | The ID of the point to insert. |
vector_mapping | string | required | The mapping to extract the vector from the document. |
payload_mapping | string | "root = {}" | An optional mapping of message to payload associated with the point. |
Advanced: tls.
questdb
Pushes messages to a QuestDB table
connector: questdb · URI connect+questdb://
| Field | Type | Default | Description |
|---|---|---|---|
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. | |
address | string | required | Address of the QuestDB server’s HTTP port (excluding protocol) |
username | string | Username for HTTP basic auth | |
password | string | Password for HTTP basic auth | |
token | string | Bearer token for HTTP auth (takes precedence over basic auth username & password) | |
table | string | required | Destination table |
designated_timestamp_field | string | Name of the designated timestamp field | |
designated_timestamp_unit | string | "auto" | Designated timestamp field units |
timestamp_string_fields | list of string | String fields with textual timestamps | |
timestamp_string_format | string | "Jan _2 15:04:05.000000Z0700" | Timestamp format, used when parsing timestamp string fields. |
symbols | list of string | Columns that should be the SYMBOL type (string values default to STRING) | |
doubles | list of string | Columns that should be double type, (int is default) | |
error_on_empty_messages | bool | false | Mark a message as errored if it is empty after field validation |
Advanced: tls, retry_timeout, request_timeout, request_min_throughput.
redis_hash
Sets Redis hash objects using the HSET command.
connector: redis_hash · URI connect+redis-hash://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
key | string | required | The key for each message, function interpolations should be used to create a unique key per message. |
walk_metadata | bool | false | Whether all metadata fields of messages should be walked and added to the list of hash fields to set. |
walk_json_object | bool | false | Whether to walk each message as a JSON object and add each key/value pair to the list of hash fields to set. |
fields | map of string | {} | A map of key/value pairs to set as hash fields. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: kind, master, client_name, tls.
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://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
key | string | required | The key for each message, function interpolations can be optionally used to create a unique key per message. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: kind, master, client_name, tls, command.
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://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
channel | string | required | The channel to publish messages to. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
batching | object | Allows you to configure a batching policy. |
Advanced: kind, master, client_name, tls.
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://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
stream | string | required | The stream to add messages to. |
id | string | "*" | The entry ID for the stream message. |
body_key | string | "body" | A key to set the raw body of the message to. |
max_length | int | 0 | When greater than zero enforces a rough cap on the length of the target stream. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
metadata | object | Specify criteria for which metadata values are included in the message body. | |
batching | object | Allows you to configure a batching policy. |
Advanced: kind, master, client_name, tls.
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.
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.
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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
output | output | required | A child output. |
Advanced: max_retries, backoff.
sftp
Writes files to an SFTP server.
connector: sftp · URI connect+sftp://
| Field | Type | Default | Description |
|---|---|---|---|
address | string | required | The address of the server to connect to. |
credentials | object | The credentials to use to log into the target server. | |
path | string | required | The file to save the messages to on the server. |
codec | string | "all-bytes" | The way in which the bytes of messages should be written out into the output data stream. |
max_in_flight | int | 64 | The maximum number of messages to have in flight at a given time. |
Advanced: connection_timeout.
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://
| Field | Type | Default | Description |
|---|---|---|---|
network | string | required | A network type to connect as. |
address | string | required | The address to connect to. |
codec | string | "lines" | The way in which the bytes of messages should be written out into the output data stream. |
Advanced: tls.
sql
Deprecated. Executes an arbitrary SQL query for each message.
connector: sql · URI connect+sql://
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
data_source_name | string | required | Data source name. |
query | string | required | The query to execute. |
args_mapping | string | An 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_flight | int | 64 | The maximum number of inserts to run in parallel. |
batching | object | Allows you to configure a batching policy. |
sql_insert
Inserts a row into an SQL database for each message.
connector: sql_insert · URI connect+sql-insert://
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
dsn | string | required | A Data Source Name to identify the target database. |
table | string | required | The table to insert to. |
columns | list of string | required | A list of columns to insert. |
args_mapping | string | required | A Bloblang mapping which should evaluate to an array of values matching in size to the number of columns specified. |
max_in_flight | int | 64 | The maximum number of inserts to run in parallel. |
batching | object | Allows 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.
sql_raw
Executes an arbitrary SQL query for each message.
connector: sql_raw · URI connect+sql-raw://
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
dsn | string | required | A Data Source Name to identify the target database. |
query | string | The query to execute. | |
args_mapping | string | An optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query. | |
queries | list of object | A list of query statements. | |
max_in_flight | int | 64 | The maximum number of batches to be sending in parallel at any given time. |
batching | object | Allows 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.
stdout
Prints messages to stdout as a continuous stream of data.
connector: stdout · URI connect+stdout://
| Field | Type | Default | Description |
|---|---|---|---|
codec | string | "lines" | The way in which the bytes of messages should be written out into the output data stream. |
subprocess
Beta. Executes a command, runs it as a subprocess, and writes messages to it over stdin.
connector: subprocess · URI connect+subprocess://
| Field | Type | Default | Description |
|---|---|---|---|
name | string | required | The command to execute as a subprocess. |
args | list of string | [] | A list of arguments to provide the command. |
codec | string | "lines" | The way in which messages should be written to the subprocess. |
switch
The switch output type allows you to route messages to different outputs based on their contents.
connector: switch · URI connect+switch://
| Field | Type | Default | Description |
|---|---|---|---|
retry_until_success | bool | false | If a selected output fails to send a message this field determines whether it is reattempted indefinitely. |
cases | list of object | A list of switch cases, outlining outputs that can be routed to. |
Advanced: strict_mode.
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://
websocket
Sends messages to an HTTP server via a websocket connection.
connector: websocket · URI connect+websocket://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL to connect to. |
Advanced: proxy_url, tls, oauth, basic_auth, jwt.