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"]
| Category | Components |
|---|---|
| AWS | aws_cloudwatch_logs, aws_kinesis, aws_s3, aws_sqs |
| Azure | azure_blob_storage, azure_cosmosdb, azure_queue_storage, azure_table_storage |
| GCP | gcp_bigquery_select, gcp_cloud_storage, gcp_pubsub |
| Local | csv, file, parquet, stdin |
| Network | http_client, http_server, nanomsg, sftp, socket, socket_server, websocket |
| Services | amqp_0_9, amqp_1, aws_cloudwatch_logs, aws_kinesis, aws_s3, aws_sqs, azure_blob_storage, azure_queue_storage, azure_table_storage, beanstalkd, cassandra, cockroachdb_changefeed, discord, gcp_bigquery_select, gcp_cloud_storage, gcp_pubsub, git, hdfs, mongodb, mqtt, nats, nats_jetstream, nats_kv, nats_stream, nsq, pulsar, redis_list, redis_pubsub, redis_scan, redis_streams, spicedb_watch, sql_raw, sql_select, timeplus, twitter_search |
| Social | discord, twitter_search |
| SpiceDB | spicedb_watch |
| Utility | batched, broker, dynamic, generate, inproc, read_until, resource, sequence, subprocess |
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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
queue | string | required | An AMQP queue to consume from. |
consumer_tag | string | "" | A consumer tag. |
prefetch_count | int | 10 | The maximum number of pending messages to have consumed at a time. |
Advanced: queue_declare, bindings_declare, auto_ack, nack_reject_patterns, prefetch_size, tls.
amqp_1
Reads messages from 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. | |
source_address | string | required | The source address to consume from. |
Advanced: azure_renew_lock, read_header, credit, tls, sasl.
aws_cloudwatch_logs
Consumes log events from AWS CloudWatch Logs.
connector: aws_cloudwatch_logs · URI connect+aws-cloudwatch-logs://
| Field | Type | Default | Description |
|---|---|---|---|
log_group_name | string | required | The name of the CloudWatch Log Group to consume from. |
log_stream_names | list of string | An optional list of log stream names to consume from. | |
log_stream_prefix | string | An optional log stream name prefix to filter streams. | |
filter_pattern | string | An optional CloudWatch Logs filter pattern to apply when querying log events. | |
start_time | string | The time to start consuming log events from. | |
poll_interval | string | "5s" | The interval at which to poll for new log events. |
auto_replay_nacks | bool | true | Whether 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.
aws_kinesis
Receive messages from one or more Kinesis streams.
connector: aws_kinesis · URI connect+aws-kinesis://
| Field | Type | Default | Description |
|---|---|---|---|
streams | list of string | required | One or more Kinesis data streams to consume from. |
dynamodb | object | Determines the table used for storing and accessing the latest consumed sequence for shards, and for coordinating balanced consumers of streams. | |
checkpoint_limit | int | 1024 | The maximum gap between the in flight sequence versus the latest acknowledged sequence at a given time. |
auto_replay_nacks | bool | true | Whether 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_period | string | "5s" | The period of time between each update to the checkpoint table. |
steal_grace_period | string | "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_oldest | bool | true | Whether to consume from the oldest message when a sequence does not yet exist for the stream. |
batching | object | Allows you to configure a batching policy. |
Advanced: poll_period, enhanced_fan_out, rebalance_period, lease_period, region, endpoint, tcp, credentials.
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://
| Field | Type | Default | Description |
|---|---|---|---|
bucket | string | "" | The bucket to consume from. |
prefix | string | "" | An optional path prefix, if set only objects with the prefix are consumed when walking a bucket. |
scanner | scanner | {"to_the_end": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
sqs | object | Consume SQS messages in order to trigger key downloads. |
Advanced: region, endpoint, tcp, credentials, force_path_style_urls, delete_objects.
aws_sqs
Consume messages from an AWS SQS URL.
connector: aws_sqs · URI connect+aws-sqs://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The SQS URL to consume from. |
max_outstanding_messages | int | 1000 | The 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.
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://
| 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 name of the container from which to download blobs. |
prefix | string | "" | An optional path prefix, if set only objects with the prefix are consumed. |
scanner | scanner | {"to_the_end": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
targets_input | input | EXPERIMENTAL: An optional source of download targets, configured as a regular Redpanda Connect input. |
Advanced: delete_objects.
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://
| 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. |
query | string | required | The query to execute |
args_mapping | string | A Bloblang mapping that, for each message, creates a list of arguments to use with the query. | |
auto_replay_nacks | bool | true | Whether 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.
azure_queue_storage
Dequeue objects from 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. |
queue_name | string | required | The name of the source storage queue. |
Advanced: dequeue_visibility_timeout, max_in_flight, track_properties.
azure_table_storage
Queries an Azure Storage Account Table, optionally with multiple filters.
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 read messages from. |
Advanced: filter, select, page_size.
batched
Consumes data from a child input and applies a batching policy to the stream.
connector: batched · URI connect+batched://
| Field | Type | Default | Description |
|---|---|---|---|
child | input | required | The child input. |
policy | object | Allows you to configure a batching policy. |
beanstalkd
Reads messages from a Beanstalkd queue.
connector: beanstalkd · URI connect+beanstalkd://
| Field | Type | Default | Description |
|---|---|---|---|
address | string | required | An address to connect to. |
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://
| Field | Type | Default | Description |
|---|---|---|---|
inputs | list of input | required | A list of inputs to create. |
batching | object | Allows you to configure a batching policy. |
Advanced: copies.
cassandra
Executes a find query and creates a message for each row received.
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. |
auto_replay_nacks | bool | true | Whether 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.
cockroachdb_changefeed
Listens to a CockroachDB Core Changefeed and creates a message for each row received.
connector: cockroachdb_changefeed · URI connect+cockroachdb-changefeed://
| Field | Type | Default | Description |
|---|---|---|---|
dsn | string | required | A Data Source Name to identify the target database. |
tables | list of string | required | CSV of tables to be included in the changefeed |
cursor_cache | string | A 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_nacks | bool | true | Whether 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.
csv
Reads one or more CSV files as structured records following the format described in RFC 4180.
connector: csv · URI connect+csv://
| Field | Type | Default | Description |
|---|---|---|---|
paths | list of string | required | A list of file paths to read from. |
parse_header_row | bool | true | Whether to reference the first row as a header row. |
delimiter | string | "," | The delimiter to use for splitting values in each record. |
lazy_quotes | bool | false | If set to true, a quote may appear in an unquoted field and a non-doubled quote may appear in a quoted field. |
auto_replay_nacks | bool | true | Whether 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.
discord
Consumes messages posted in a Discord channel.
connector: discord · URI connect+discord://
| Field | Type | Default | Description |
|---|---|---|---|
channel_id | string | required | A discord channel ID to consume messages from. |
bot_token | string | required | A bot token used for authentication. |
cache | string | required | A 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_nacks | bool | true | Whether 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
inputs | map of input | {} | A map of inputs to statically create. |
prefix | string | "" | A path prefix for HTTP endpoints that are registered. |
file
Consumes data from files on disk, emitting messages according to a chosen codec.
connector: file · URI connect+file://
| Field | Type | Default | Description |
|---|---|---|---|
paths | list of string | required | A list of paths to consume sequentially. |
scanner | scanner | {"lines": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
auto_replay_nacks | bool | true | Whether 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
project | string | required | GCP project where the query job will execute. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
table | string | required | Fully-qualified BigQuery table name to query. |
columns | list of string | required | A list of columns to query. |
where | string | An optional where clause to add. | |
auto_replay_nacks | bool | true | Whether 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_labels | map of string | {} | A list of labels to add to the query job. |
priority | string | "" | The priority with which to schedule the query. |
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 where. | |
prefix | string | An optional prefix to prepend to the select query (before SELECT). | |
suffix | string | An optional suffix to append to the select query. |
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://
| Field | Type | Default | Description |
|---|---|---|---|
bucket | string | required | The name of the bucket from which to download objects. |
prefix | string | "" | An optional path prefix, if set only objects with the prefix are consumed. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
scanner | scanner | {"to_the_end": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
Advanced: delete_objects.
gcp_pubsub
Consumes messages from a GCP Cloud Pub/Sub subscription.
connector: gcp_pubsub · URI connect+gcp-pubsub://
| Field | Type | Default | Description |
|---|---|---|---|
project | string | required | The project ID of the target subscription. |
credentials_json | string | "" | An optional field to set Google Service Account Credentials json. |
subscription | string | required | The target subscription ID. |
endpoint | string | "" | An optional endpoint to override the default of pubsub.googleapis.com:443. |
sync | bool | false | Enable synchronous pull mode. |
max_outstanding_messages | int | 1000 | The maximum number of outstanding pending messages to be consumed at a given time. |
max_outstanding_bytes | int | 1000000000 | The maximum number of outstanding pending messages to be consumed measured in bytes. |
Advanced: create_subscription.
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://
| Field | Type | Default | Description |
|---|---|---|---|
mapping | string | required | A Bloblang mapping to use for generating messages. |
interval | string | "1s" | The time interval at which messages should be generated, expressed either as a duration string or as a cron expression. |
count | int | 0 | An 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_size | int | 1 | The number of generated messages that should be accumulated into each batch flushed at the specified interval. |
auto_replay_nacks | bool | true | Whether 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. |
git
A Git input that clones (or pulls) a repository and reads the repository contents.
connector: git · URI connect+git://
| Field | Type | Default | Description |
|---|---|---|---|
repository_url | string | required | The URL of the Git repository to clone. |
branch | string | "main" | The branch to check out. |
poll_interval | string | "10s" | Duration between polling attempts |
include_patterns | list of string | [] | A list of file patterns to include (e.g., ‘**/.md’, ‘configs/.yaml’). |
exclude_patterns | list of string | [] | A list of file patterns to exclude (e.g., ‘.git/’, ‘/*.png’). |
max_file_size | int | 10485760 | The maximum size of files to include in bytes. |
checkpoint_cache | string | A cache resource to store the last processed commit hash, allowing the input to resume from where it left off after a restart. | |
checkpoint_key | string | "git_last_commit" | The key to use when storing the last processed commit hash in the cache. |
auth | object | Authentication options for the Git repository | |
auto_replay_nacks | bool | true | Whether 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. |
hdfs
Reads files from a HDFS directory, where each discrete file will be consumed as a single message payload.
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 | The directory to consume from. |
http_client
Connects to a server and continuously performs requests for a single message.
connector: http_client · URI connect+http-client://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL to connect to. |
verb | string | "GET" | 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. |
payload | string | An optional payload to deliver for each request. | |
stream | object | Allows you to set streaming mode, where requests are kept open and messages are processed line-by-line. | |
auto_replay_nacks | bool | true | Whether 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
address | string | "" | An alternative address to host from. |
path | string | "/post" | The endpoint path to listen for POST requests. |
ws_path | string | "/post/ws" | The endpoint path to create websocket connections from. |
allowed_verbs | list of string | ["POST"] | An array of verbs that are allowed for the path endpoint. |
timeout | string | "5s" | Timeout for requests. |
rate_limit | string | "" | An optional rate limit to throttle requests by. |
Advanced: ws_welcome_message, ws_rate_limit_message, cert_file, key_file, cors, sync_response, tcp.
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.
mongodb
Executes a query and creates a message for each document received.
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 collection to select from. |
query | string | required | Bloblang expression describing MongoDB query. |
auto_replay_nacks | bool | true | Whether 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_size | int | A explicit number of documents to batch up before flushing them for processing. | |
sort | map of int | An object specifying fields to sort by, and the respective sort order (1 ascending, -1 descending). | |
limit | int | An explicit maximum number of documents to return. |
Advanced: app_name, aws, operation, json_marshal_mode.
mqtt
Subscribe to topics on MQTT brokers.
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. |
topics | list of string | required | A list of topics to consume from. |
auto_replay_nacks | bool | true | Whether 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.
nanomsg
Consumes messages via Nanomsg sockets (scalability protocols).
connector: nanomsg · URI connect+nanomsg://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to (or as). |
bind | bool | true | Whether the URLs provided should be connected to, or bound as. |
socket_type | string | "PULL" | The socket type to use. |
auto_replay_nacks | bool | true | Whether 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_filters | list of string | [] | A list of subscription topic filters to use when consuming from a SUB socket. |
Advanced: poll_timeout.
nats
Subscribe to a 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 | A subject to consume from. |
queue | string | An optional queue group to consume as. | |
auto_replay_nacks | bool | true | Whether 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_ack | bool | true | Control 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.
nats_jetstream
Reads messages from NATS JetStream subjects.
connector: nats_jetstream · URI connect+nats-jetstream://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
queue | string | An optional queue group to consume as. | |
subject | string | A subject to consume from. | |
durable | string | Preserve the state of your consumer under a durable name. | |
stream | string | A stream to consume from. | |
bind | bool | Indicates that the subscription should use an existing consumer. | |
deliver | string | "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.
nats_kv
Watches for updates 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 | ">" | Key to watch for updates, can include wildcards. |
auto_replay_nacks | bool | true | Whether 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
cluster_id | string | required | The ID of the cluster to consume from. |
client_id | string | "" | A client ID to connect as. |
queue | string | "" | The queue to consume from. |
subject | string | "" | A subject to consume from. |
durable_name | string | "" | Preserve the state of your consumer under a durable name. |
unsubscribe_on_close | bool | false | Whether 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.
nsq
Subscribe to an NSQ instance topic and channel.
connector: nsq · URI connect+nsq://
| Field | Type | Default | Description |
|---|---|---|---|
nsqd_tcp_addresses | list of string | required | A list of nsqd addresses to connect to. |
lookupd_http_addresses | list of string | required | A list of nsqlookupd addresses to connect to. |
topic | string | required | The topic to consume from. |
channel | string | required | The channel to consume from. |
user_agent | string | A user agent to assume when connecting. | |
max_in_flight | int | 100 | The maximum number of pending messages to consume at any given time. |
max_attempts | int | 5 | The maximum number of attempts to successfully consume a messages. |
Advanced: tls.
parquet
Reads and decodes Parquet files into a stream of structured messages.
connector: parquet · URI connect+parquet://
| Field | Type | Default | Description |
|---|---|---|---|
paths | list of string | required | A list of file paths to read from. |
auto_replay_nacks | bool | true | Whether 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.
pulsar
Reads messages from an Apache Pulsar server.
connector: pulsar · URI connect+pulsar://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | A URL to connect to. |
topics | list of string | A list of topics to subscribe to. | |
topics_pattern | string | A regular expression matching the topics to subscribe to. | |
subscription_name | string | required | Specify the subscription name for this consumer. |
subscription_type | string | "shared" | Specify the subscription type for this consumer. |
subscription_initial_position | string | "latest" | Specify the subscription initial position for this consumer. |
tls | object | Specify the path to a custom CA certificate to trust broker TLS service. |
Advanced: auth.
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://
| Field | Type | Default | Description |
|---|---|---|---|
input | input | required | The child input to consume from. |
check | string | A Bloblang query that should return a boolean value indicating whether the input should now be closed. | |
idle_timeout | string | The maximum amount of time without receiving new messages after which the input is closed. | |
restart_input | bool | false | Whether the input should be reopened if it closes itself before the condition has resolved to true. |
redis_list
Pops messages from the beginning of a Redis list using the BLPop 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 of a list to read from. |
auto_replay_nacks | bool | true | Whether 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.
redis_pubsub
Consume from a Redis publish/subscribe channel using either the SUBSCRIBE or PSUBSCRIBE commands.
connector: redis_pubsub · URI connect+redis-pubsub://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
channels | list of string | required | A list of channels to consume from. |
use_patterns | bool | false | Whether to use the PSUBSCRIBE command, allowing for glob-style patterns within target channel names. |
auto_replay_nacks | bool | true | Whether 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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
auto_replay_nacks | bool | true | Whether 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. |
match | string | "" | Iterates only elements matching the optional glob-style pattern. |
Advanced: kind, master, client_name, tls.
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://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
body_key | string | "body" | The field key to extract the raw message from. |
streams | list of string | required | A list of streams to consume from. |
auto_replay_nacks | bool | true | Whether 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. |
limit | int | 10 | The maximum number of messages to consume from a single request. |
client_id | string | "" | An identifier for the client connection. |
consumer_group | string | "" | An identifier for the consumer group of the stream. |
Advanced: kind, master, client_name, tls, create_streams, start_from_oldest, commit_period, timeout.
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.
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://
| Field | Type | Default | Description |
|---|---|---|---|
inputs | list of input | required | An array of inputs to read from sequentially. |
Advanced: sharded_join.
sftp
Consumes files from 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. | |
paths | list of string | required | A list of paths to consume sequentially. |
auto_replay_nacks | bool | true | Whether 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. |
scanner | scanner | {"to_the_end": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
watcher | object | An 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.
socket
Connects to a tcp or unix socket and consumes a continuous stream of messages.
connector: socket · URI connect+socket://
| Field | Type | Default | Description |
|---|---|---|---|
network | string | required | A network type to assume (unix|tcp). |
address | string | required | The address to connect to. |
auto_replay_nacks | bool | true | Whether 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_mapping | string | An optional Bloblang mapping which should evaluate to a string which will be sent upstream before the downstream data flow starts. | |
scanner | scanner | {"lines": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
Advanced: tls.
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://
| Field | Type | Default | Description |
|---|---|---|---|
network | string | required | A network type to accept. |
address | string | required | The address to listen from. |
address_cache | string | An optional cache within which this input should write it’s bound address once known. | |
tls | object | TLS specific configuration, valid when the network is set to tls. | |
auto_replay_nacks | bool | true | Whether 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. |
scanner | scanner | {"lines": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
Advanced: tcp.
spicedb_watch
Consume messages from the Watch API from SpiceDB.
connector: spicedb_watch · URI connect+spicedb-watch://
| Field | Type | Default | Description |
|---|---|---|---|
endpoint | string | required | The SpiceDB endpoint. |
bearer_token | string | "" | The SpiceDB Bearer token used to authenticate against the SpiceDB instance. |
cache | string | required | A 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.
sql_raw
Executes a select query and creates a message for each row received.
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 | 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. | |
auto_replay_nacks | bool | true | Whether 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.
sql_select
Executes a select query and creates a message for each row received.
connector: sql_select · URI connect+sql-select://
| 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 select from. |
columns | list of string | required | A list of columns to select. |
where | string | An optional where clause to add. | |
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 where. | |
auto_replay_nacks | bool | true | Whether 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.
stdin
Consumes data piped to stdin, chopping it into individual messages according to the specified scanner.
connector: stdin · URI connect+stdin://
| Field | Type | Default | Description |
|---|---|---|---|
scanner | scanner | {"lines": {}} | The scanner by which the stream of bytes consumed will be broken out into individual messages. |
auto_replay_nacks | bool | true | Whether 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. |
subprocess
Beta. Executes a command, runs it as a subprocess, and consumes messages from it over stdout.
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 consumed from the subprocess. |
restart_on_exit | bool | false | Whether the command should be re-executed each time the subprocess ends. |
Advanced: max_buffer.
timeplus
Executes a query on Timeplus Enterprise and creates a message from each row received
connector: timeplus · URI connect+timeplus://
| Field | Type | Default | Description |
|---|---|---|---|
query | string | required | The query to run |
url | string | "tcp://localhost:8463" | The url should always include schema and host. |
workspace | string | ID of the workspace. | |
apikey | string | The API key. | |
username | string | The username. | |
password | string | The password. |
twitter_search
Experimental. Consumes tweets matching a given search using the Twitter recent search V2 API.
connector: twitter_search · URI connect+twitter-search://
| Field | Type | Default | Description |
|---|---|---|---|
query | string | required | A search expression to use. |
tweet_fields | list of string | [] | An optional list of additional fields to obtain for each tweet, by default only the fields id and text are returned. |
poll_period | string | "1m" | The length of time (as a duration string) to wait between each search request. |
backfill_period | string | "5m" | A duration string indicating the maximum age of tweets to acquire when starting a search. |
cache | string | required | A cache resource to use for request pagination. |
api_key | string | required | An API key for OAuth 2.0 authentication. |
api_secret | string | required | An API secret for OAuth 2.0 authentication. |
Advanced: cache_key, rate_limit.
websocket
Connects to a websocket server and continuously receives messages.
connector: websocket · URI connect+websocket://
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL to connect to. |
auto_replay_nacks | bool | true | Whether 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.