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

The 56 Redpanda Connect inputs 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 input:

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

amqp_0_9

Connects to an AMQP (0.91) queue. 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.
queuestringrequiredAn AMQP queue to consume from.
consumer_tagstring""A consumer tag.
prefetch_countint10The maximum number of pending messages to have consumed at a time.

Advanced: queue_declare, bindings_declare, auto_ack, nack_reject_patterns, prefetch_size, tls.

Upstream documentation

amqp_1

Reads messages from an AMQP (1.0) server.

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

FieldTypeDefaultDescription
urlslist of stringA list of URLs to connect to.
source_addressstringrequiredThe source address to consume from.

Advanced: azure_renew_lock, read_header, credit, tls, sasl.

Upstream documentation

aws_cloudwatch_logs

Consumes log events from AWS CloudWatch Logs.

connector: aws_cloudwatch_logs · URI connect+aws-cloudwatch-logs://

FieldTypeDefaultDescription
log_group_namestringrequiredThe name of the CloudWatch Log Group to consume from.
log_stream_nameslist of stringAn optional list of log stream names to consume from.
log_stream_prefixstringAn optional log stream name prefix to filter streams.
filter_patternstringAn optional CloudWatch Logs filter pattern to apply when querying log events.
start_timestringThe time to start consuming log events from.
poll_intervalstring"5s"The interval at which to poll for new log events.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: limit, structured_log, api_timeout, region, endpoint, tcp, credentials.

Upstream documentation

aws_kinesis

Receive messages from one or more Kinesis streams.

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

FieldTypeDefaultDescription
streamslist of stringrequiredOne or more Kinesis data streams to consume from.
dynamodbobjectDetermines the table used for storing and accessing the latest consumed sequence for shards, and for coordinating balanced consumers of streams.
checkpoint_limitint1024The maximum gap between the in flight sequence versus the latest acknowledged sequence at a given time.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
commit_periodstring"5s"The period of time between each update to the checkpoint table.
steal_grace_periodstring"2s"Determines how long beyond the next commit period a client will wait when stealing a shard for the current owner to store a checkpoint.
start_from_oldestbooltrueWhether to consume from the oldest message when a sequence does not yet exist for the stream.
batchingobjectAllows you to configure a batching policy.

Advanced: poll_period, enhanced_fan_out, rebalance_period, lease_period, region, endpoint, tcp, credentials.

Upstream documentation

aws_s3

Downloads objects within an Amazon S3 bucket, optionally filtered by a prefix, either by walking the items in the bucket or by streaming upload notifications in realtime.

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

FieldTypeDefaultDescription
bucketstring""The bucket to consume from.
prefixstring""An optional path prefix, if set only objects with the prefix are consumed when walking a bucket.
scannerscanner{"to_the_end": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.
sqsobjectConsume SQS messages in order to trigger key downloads.

Advanced: region, endpoint, tcp, credentials, force_path_style_urls, delete_objects.

Upstream documentation

aws_sqs

Consume messages from an AWS SQS URL.

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

FieldTypeDefaultDescription
urlstringrequiredThe SQS URL to consume from.
max_outstanding_messagesint1000The maximum number of outstanding pending messages to be consumed at a given time.

Advanced: delete_message, reset_visibility, max_number_of_messages, wait_time_seconds, message_timeout, region, endpoint, tcp, credentials.

Upstream documentation

azure_blob_storage

Downloads objects within an Azure Blob Storage container, optionally filtered by a prefix.

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 name of the container from which to download blobs.
prefixstring""An optional path prefix, if set only objects with the prefix are consumed.
scannerscanner{"to_the_end": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.
targets_inputinputEXPERIMENTAL: An optional source of download targets, configured as a regular Redpanda Connect input.

Advanced: delete_objects.

Upstream documentation

azure_cosmosdb

Executes a SQL query against Azure CosmosDB and creates a batch of messages from each page of items.

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.
querystringrequiredThe query to execute
args_mappingstringA Bloblang mapping that, for each message, creates a list of arguments to use with the query.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: batch_count.

Upstream documentation

azure_queue_storage

Dequeue objects from 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.
queue_namestringrequiredThe name of the source storage queue.

Advanced: dequeue_visibility_timeout, max_in_flight, track_properties.

Upstream documentation

azure_table_storage

Queries an Azure Storage Account Table, optionally with multiple filters.

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 read messages from.

Advanced: filter, select, page_size.

Upstream documentation

batched

Consumes data from a child input and applies a batching policy to the stream.

connector: batched · URI connect+batched://

FieldTypeDefaultDescription
childinputrequiredThe child input.
policyobjectAllows you to configure a batching policy.

Upstream documentation

beanstalkd

Reads messages from a Beanstalkd queue.

connector: beanstalkd · URI connect+beanstalkd://

FieldTypeDefaultDescription
addressstringrequiredAn address to connect to.

Upstream documentation

broker

Allows you to combine multiple inputs into a single stream of data, where each input will be read in parallel.

connector: broker · URI connect+broker://

FieldTypeDefaultDescription
inputslist of inputrequiredA list of inputs to create.
batchingobjectAllows you to configure a batching policy.

Advanced: copies.

Upstream documentation

cassandra

Executes a find query and creates a message for each row received.

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.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

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

Upstream documentation

cockroachdb_changefeed

Listens to a CockroachDB Core Changefeed and creates a message for each row received.

connector: cockroachdb_changefeed · URI connect+cockroachdb-changefeed://

FieldTypeDefaultDescription
dsnstringrequiredA Data Source Name to identify the target database.
tableslist of stringrequiredCSV of tables to be included in the changefeed
cursor_cachestringA cache resource to use for storing the current latest cursor that has been successfully delivered, this allows Redpanda Connect to continue from that cursor upon restart, rather than consume the entire state of the table.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: tls, options.

Upstream documentation

csv

Reads one or more CSV files as structured records following the format described in RFC 4180.

connector: csv · URI connect+csv://

FieldTypeDefaultDescription
pathslist of stringrequiredA list of file paths to read from.
parse_header_rowbooltrueWhether to reference the first row as a header row.
delimiterstring","The delimiter to use for splitting values in each record.
lazy_quotesboolfalseIf set to true, a quote may appear in an unquoted field and a non-doubled quote may appear in a quoted field.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: delete_on_finish, batch_count.

Upstream documentation

discord

Consumes messages posted in a Discord channel.

connector: discord · URI connect+discord://

FieldTypeDefaultDescription
channel_idstringrequiredA discord channel ID to consume messages from.
bot_tokenstringrequiredA bot token used for authentication.
cachestringrequiredA cache resource to use for performing unread message backfills, the ID of the last message received will be stored in this cache and used for subsequent requests.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: cache_key.

Upstream documentation

dynamic

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

connector: dynamic · URI connect+dynamic://

FieldTypeDefaultDescription
inputsmap of input{}A map of inputs to statically create.
prefixstring""A path prefix for HTTP endpoints that are registered.

Upstream documentation

file

Consumes data from files on disk, emitting messages according to a chosen codec.

connector: file · URI connect+file://

FieldTypeDefaultDescription
pathslist of stringrequiredA list of paths to consume sequentially.
scannerscanner{"lines": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: delete_on_finish.

Upstream documentation

gcp_bigquery_select

Executes a SELECT query against BigQuery and creates a message for each row received.

connector: gcp_bigquery_select · URI connect+gcp-bigquery-select://

FieldTypeDefaultDescription
projectstringrequiredGCP project where the query job will execute.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
tablestringrequiredFully-qualified BigQuery table name to query.
columnslist of stringrequiredA list of columns to query.
wherestringAn optional where clause to add.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
job_labelsmap of string{}A list of labels to add to the query job.
prioritystring""The priority with which to schedule the query.
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 where.
prefixstringAn optional prefix to prepend to the select query (before SELECT).
suffixstringAn optional suffix to append to the select query.

Upstream documentation

gcp_cloud_storage

Downloads objects within a Google Cloud Storage bucket, optionally filtered by a prefix.

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

FieldTypeDefaultDescription
bucketstringrequiredThe name of the bucket from which to download objects.
prefixstring""An optional path prefix, if set only objects with the prefix are consumed.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
scannerscanner{"to_the_end": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.

Advanced: delete_objects.

Upstream documentation

gcp_pubsub

Consumes messages from a GCP Cloud Pub/Sub subscription.

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

FieldTypeDefaultDescription
projectstringrequiredThe project ID of the target subscription.
credentials_jsonstring""An optional field to set Google Service Account Credentials json.
subscriptionstringrequiredThe target subscription ID.
endpointstring""An optional endpoint to override the default of pubsub.googleapis.com:443.
syncboolfalseEnable synchronous pull mode.
max_outstanding_messagesint1000The maximum number of outstanding pending messages to be consumed at a given time.
max_outstanding_bytesint1000000000The maximum number of outstanding pending messages to be consumed measured in bytes.

Advanced: create_subscription.

Upstream documentation

generate

Generates messages at a given interval using a Bloblang mapping executed without a context. This allows you to generate messages for testing your pipeline configs.

connector: generate · URI connect+generate://

FieldTypeDefaultDescription
mappingstringrequiredA Bloblang mapping to use for generating messages.
intervalstring"1s"The time interval at which messages should be generated, expressed either as a duration string or as a cron expression.
countint0An optional number of messages to generate, if set above 0 the specified number of messages is generated and then the input will shut down.
batch_sizeint1The number of generated messages that should be accumulated into each batch flushed at the specified interval.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Upstream documentation

git

A Git input that clones (or pulls) a repository and reads the repository contents.

connector: git · URI connect+git://

FieldTypeDefaultDescription
repository_urlstringrequiredThe URL of the Git repository to clone.
branchstring"main"The branch to check out.
poll_intervalstring"10s"Duration between polling attempts
include_patternslist of string[]A list of file patterns to include (e.g., ‘**/.md’, ‘configs/.yaml’).
exclude_patternslist of string[]A list of file patterns to exclude (e.g., ‘.git/’, ‘/*.png’).
max_file_sizeint10485760The maximum size of files to include in bytes.
checkpoint_cachestringA cache resource to store the last processed commit hash, allowing the input to resume from where it left off after a restart.
checkpoint_keystring"git_last_commit"The key to use when storing the last processed commit hash in the cache.
authobjectAuthentication options for the Git repository
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Upstream documentation

hdfs

Reads files from a HDFS directory, where each discrete file will be consumed as a single message payload.

connector: hdfs · URI connect+hdfs://

FieldTypeDefaultDescription
hostslist of stringrequiredA list of target host addresses to connect to.
userstring""A user ID to connect as.
directorystringrequiredThe directory to consume from.

Upstream documentation

http_client

Connects to a server and continuously performs requests for a single message.

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

FieldTypeDefaultDescription
urlstringrequiredThe URL to connect to.
verbstring"GET"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.
payloadstringAn optional payload to deliver for each request.
streamobjectAllows you to set streaming mode, where requests are kept open and messages are processed line-by-line.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

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, drop_empty_bodies.

Upstream documentation

http_server

Receive messages POSTed over HTTP(S). 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"/post"The endpoint path to listen for POST requests.
ws_pathstring"/post/ws"The endpoint path to create websocket connections from.
allowed_verbslist of string["POST"]An array of verbs that are allowed for the path endpoint.
timeoutstring"5s"Timeout for requests.
rate_limitstring""An optional rate limit to throttle requests by.

Advanced: ws_welcome_message, ws_rate_limit_message, cert_file, key_file, cors, sync_response, tcp.

Upstream documentation

inproc

Directly connect to an output within a Redpanda Connect process by referencing it by a chosen ID.

connector: inproc · URI connect+inproc://

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

Upstream documentation

mongodb

Executes a query and creates a message for each document received.

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 collection to select from.
querystringrequiredBloblang expression describing MongoDB query.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
batch_sizeintA explicit number of documents to batch up before flushing them for processing.
sortmap of intAn object specifying fields to sort by, and the respective sort order (1 ascending, -1 descending).
limitintAn explicit maximum number of documents to return.

Advanced: app_name, aws, operation, json_marshal_mode.

Upstream documentation

mqtt

Subscribe to topics on MQTT brokers.

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.
topicslist of stringrequiredA list of topics to consume from.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: dynamic_client_id_suffix, will, user, password, keepalive, tls, qos, clean_session.

Upstream documentation

nanomsg

Consumes messages via Nanomsg sockets (scalability protocols).

connector: nanomsg · URI connect+nanomsg://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to (or as).
bindbooltrueWhether the URLs provided should be connected to, or bound as.
socket_typestring"PULL"The socket type to use.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
sub_filterslist of string[]A list of subscription topic filters to use when consuming from a SUB socket.

Advanced: poll_timeout.

Upstream documentation

nats

Subscribe to a NATS subject.

connector: nats · URI connect+nats://

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
subjectstringrequiredA subject to consume from.
queuestringAn optional queue group to consume as.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
send_ackbooltrueControl whether ACKS are sent as a reply to each message.

Advanced: max_reconnects, nak_delay, prefetch_count, tls, tls_handshake_first, auth, extract_tracing_map.

Upstream documentation

nats_jetstream

Reads messages from NATS JetStream subjects.

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

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
queuestringAn optional queue group to consume as.
subjectstringA subject to consume from.
durablestringPreserve the state of your consumer under a durable name.
streamstringA stream to consume from.
bindboolIndicates that the subscription should use an existing consumer.
deliverstring"all"Determines which messages to deliver when consuming without a durable subscriber.

Advanced: max_reconnects, create_stream, ack_wait, max_ack_pending, tls, tls_handshake_first, auth, extract_tracing_map.

Upstream documentation

nats_kv

Watches for updates 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.
keystring">"Key to watch for updates, can include wildcards.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: max_reconnects, ignore_deletes, include_history, meta_only, tls, tls_handshake_first, auth.

Upstream documentation

nats_stream

Subscribe to a NATS Stream subject. Joining a queue is optional and allows multiple clients of a subject to consume using queue semantics.

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

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
cluster_idstringrequiredThe ID of the cluster to consume from.
client_idstring""A client ID to connect as.
queuestring""The queue to consume from.
subjectstring""A subject to consume from.
durable_namestring""Preserve the state of your consumer under a durable name.
unsubscribe_on_closeboolfalseWhether the subscription should be destroyed when this client disconnects.

Advanced: max_reconnects, start_from_oldest, max_inflight, ack_wait, tls, tls_handshake_first, auth, extract_tracing_map.

Upstream documentation

nsq

Subscribe to an NSQ instance topic and channel.

connector: nsq · URI connect+nsq://

FieldTypeDefaultDescription
nsqd_tcp_addresseslist of stringrequiredA list of nsqd addresses to connect to.
lookupd_http_addresseslist of stringrequiredA list of nsqlookupd addresses to connect to.
topicstringrequiredThe topic to consume from.
channelstringrequiredThe channel to consume from.
user_agentstringA user agent to assume when connecting.
max_in_flightint100The maximum number of pending messages to consume at any given time.
max_attemptsint5The maximum number of attempts to successfully consume a messages.

Advanced: tls.

Upstream documentation

parquet

Reads and decodes Parquet files into a stream of structured messages.

connector: parquet · URI connect+parquet://

FieldTypeDefaultDescription
pathslist of stringrequiredA list of file paths to read from.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: batch_count.

Upstream documentation

pulsar

Reads messages from an Apache Pulsar server.

connector: pulsar · URI connect+pulsar://

FieldTypeDefaultDescription
urlstringrequiredA URL to connect to.
topicslist of stringA list of topics to subscribe to.
topics_patternstringA regular expression matching the topics to subscribe to.
subscription_namestringrequiredSpecify the subscription name for this consumer.
subscription_typestring"shared"Specify the subscription type for this consumer.
subscription_initial_positionstring"latest"Specify the subscription initial position for this consumer.
tlsobjectSpecify the path to a custom CA certificate to trust broker TLS service.

Advanced: auth.

Upstream documentation

read_until

Reads messages from a child input until a consumed message passes a Bloblang query, at which point the input closes. It is also possible to configure a timeout after which the input is closed if no new messages arrive in that period.

connector: read_until · URI connect+read-until://

FieldTypeDefaultDescription
inputinputrequiredThe child input to consume from.
checkstringA Bloblang query that should return a boolean value indicating whether the input should now be closed.
idle_timeoutstringThe maximum amount of time without receiving new messages after which the input is closed.
restart_inputboolfalseWhether the input should be reopened if it closes itself before the condition has resolved to true.

Upstream documentation

redis_list

Pops messages from the beginning of a Redis list using the BLPop command.

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

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
keystringrequiredThe key of a list to read from.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

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

Upstream documentation

redis_pubsub

Consume from a Redis publish/subscribe channel using either the SUBSCRIBE or PSUBSCRIBE commands.

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

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
channelslist of stringrequiredA list of channels to consume from.
use_patternsboolfalseWhether to use the PSUBSCRIBE command, allowing for glob-style patterns within target channel names.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: kind, master, client_name, tls.

Upstream documentation

redis_scan

Scans the set of keys in the current selected database and gets their values, using the Scan and Get commands.

connector: redis_scan · URI connect+redis-scan://

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
matchstring""Iterates only elements matching the optional glob-style pattern.

Advanced: kind, master, client_name, tls.

Upstream documentation

redis_streams

Pulls messages from Redis (v5.0+) streams with the XREADGROUP command. The client_id should be unique for each consumer of a group.

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

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
body_keystring"body"The field key to extract the raw message from.
streamslist of stringrequiredA list of streams to consume from.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
limitint10The maximum number of messages to consume from a single request.
client_idstring""An identifier for the client connection.
consumer_groupstring""An identifier for the consumer group of the stream.

Advanced: kind, master, client_name, tls, create_streams, start_from_oldest, commit_period, timeout.

Upstream documentation

resource

Resource is an input type that channels messages from a resource input, identified by its name.

connector: resource · URI connect+resource://

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

Upstream documentation

sequence

Reads messages from a sequence of child inputs, starting with the first and once that input gracefully terminates starts consuming from the next, and so on.

connector: sequence · URI connect+sequence://

FieldTypeDefaultDescription
inputslist of inputrequiredAn array of inputs to read from sequentially.

Advanced: sharded_join.

Upstream documentation

sftp

Consumes files from 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.
pathslist of stringrequiredA list of paths to consume sequentially.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
scannerscanner{"to_the_end": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.
watcherobjectAn experimental mode whereby the input will periodically scan the target paths for new files and consume them, when all files are consumed the input will continue polling for new files.

Advanced: connection_timeout, max_sftp_sessions, delete_on_finish.

Upstream documentation

socket

Connects to a tcp or unix socket and consumes a continuous stream of messages.

connector: socket · URI connect+socket://

FieldTypeDefaultDescription
networkstringrequiredA network type to assume (unix|tcp).
addressstringrequiredThe address to connect to.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
open_message_mappingstringAn optional Bloblang mapping which should evaluate to a string which will be sent upstream before the downstream data flow starts.
scannerscanner{"lines": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.

Advanced: tls.

Upstream documentation

socket_server

Creates a server that receives a stream of messages over a TCP, UDP or Unix socket.

connector: socket_server · URI connect+socket-server://

FieldTypeDefaultDescription
networkstringrequiredA network type to accept.
addressstringrequiredThe address to listen from.
address_cachestringAn optional cache within which this input should write it’s bound address once known.
tlsobjectTLS specific configuration, valid when the network is set to tls.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.
scannerscanner{"lines": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.

Advanced: tcp.

Upstream documentation

spicedb_watch

Consume messages from the Watch API from SpiceDB.

connector: spicedb_watch · URI connect+spicedb-watch://

FieldTypeDefaultDescription
endpointstringrequiredThe SpiceDB endpoint.
bearer_tokenstring""The SpiceDB Bearer token used to authenticate against the SpiceDB instance.
cachestringrequiredA cache resource to use for performing unread message backfills, the ID of the last message received will be stored in this cache and used for subsequent requests.

Advanced: max_receive_message_bytes, cache_key, tls.

Upstream documentation

sql_raw

Executes a select query and creates a message for each row received.

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

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
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.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

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

Upstream documentation

sql_select

Executes a select query and creates a message for each row received.

connector: sql_select · URI connect+sql-select://

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
tablestringrequiredThe table to select from.
columnslist of stringrequiredA list of columns to select.
wherestringAn optional where clause to add.
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 where.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

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

Upstream documentation

stdin

Consumes data piped to stdin, chopping it into individual messages according to the specified scanner.

connector: stdin · URI connect+stdin://

FieldTypeDefaultDescription
scannerscanner{"lines": {}}The scanner by which the stream of bytes consumed will be broken out into individual messages.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Upstream documentation

subprocess

Beta. Executes a command, runs it as a subprocess, and consumes messages from it over stdout.

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 consumed from the subprocess.
restart_on_exitboolfalseWhether the command should be re-executed each time the subprocess ends.

Advanced: max_buffer.

Upstream documentation

timeplus

Executes a query on Timeplus Enterprise and creates a message from each row received

connector: timeplus · URI connect+timeplus://

FieldTypeDefaultDescription
querystringrequiredThe query to run
urlstring"tcp://localhost:8463"The url should always include schema and host.
workspacestringID of the workspace.
apikeystringThe API key.
usernamestringThe username.
passwordstringThe password.

Upstream documentation

Experimental. Consumes tweets matching a given search using the Twitter recent search V2 API.

connector: twitter_search · URI connect+twitter-search://

FieldTypeDefaultDescription
querystringrequiredA search expression to use.
tweet_fieldslist of string[]An optional list of additional fields to obtain for each tweet, by default only the fields id and text are returned.
poll_periodstring"1m"The length of time (as a duration string) to wait between each search request.
backfill_periodstring"5m"A duration string indicating the maximum age of tweets to acquire when starting a search.
cachestringrequiredA cache resource to use for request pagination.
api_keystringrequiredAn API key for OAuth 2.0 authentication.
api_secretstringrequiredAn API secret for OAuth 2.0 authentication.

Advanced: cache_key, rate_limit.

Upstream documentation

websocket

Connects to a websocket server and continuously receives messages.

connector: websocket · URI connect+websocket://

FieldTypeDefaultDescription
urlstringrequiredThe URL to connect to.
auto_replay_nacksbooltrueWhether messages that are rejected (nacked) at the output level should be automatically replayed indefinitely, eventually resulting in back pressure if the cause of the rejections is persistent.

Advanced: proxy_url, open_message, open_message_type, max_message_size, tls, connection, oauth, basic_auth, jwt.

Upstream documentation