Cloud

Components Catalog

Use the following table to search for available inputs, outputs, and processors.

Name Connector Type

memory

Buffer, Cache

none

Buffer, Metric, Tracer

system_window

Buffer

aws_dynamodbAWS DynamoDB Amazon DynamoDB DynamoDB

Cache, Output

aws_s3AWS S3 Amazon S3 S3 Simple Storage Service

Cache, Input, Output

gcp_cloud_storageGCP Cloud Storage Google Cloud Storage GCS

Cache, Input, Output

lru

Cache

memcached

Cache

mongodbMongo

Cache, Input, Output, Processor

multilevel

Cache

nats_kvNATS KV

Cache, Input, Output, Processor

noop

Cache, Processor

redis

Cache, Processor, Rate_limit

redpanda

Cache, Input, Output, Tracer

ristretto

Cache

sql

Cache

ttlru

Cache

amqp_0_9RabbitMQ AMQP

Input, Output

aws_cloudwatch_logsAWS CloudWatch Logs Amazon CloudWatch Logs

Input

aws_dynamodb_cdcAmazon DynamoDB CDC

Input

aws_kinesisAWS Kinesis Amazon Kinesis Kinesis

Input, Output

aws_sqsAWS SQS Amazon SQS SQS Simple Queue Service

Input, Output

azure_blob_storageAzure Blob Storage Microsoft Azure Storage

Input, Output

azure_cosmosdbMicrosoft Azure Azure

Input, Output, Processor

azure_queue_storageAzure Queue Storage Microsoft Azure Queue

Input, Output

azure_table_storageAzure Table Storage Microsoft Azure Table

Input, Output

batched

Input

broker

Input, Output

csvComma-Separated Values

Scanner

gateway

Input

gcp_bigquery_selectGCP BigQuery Google Cloud GCP

Input, Processor

gcp_pubsubGCP PubSub Google Cloud Pub/Sub GCP Pub/Sub Google Pub/Sub

Input, Output

gcp_spanner_cdcGoogle Cloud GCP

Input

generate

Input

git

Input

http_clientHTTP REST API REST

Input, Output

http_serverHTTP REST API REST Gateway

Input

inproc

Input, Output

jiraAtlassian Jira

Input, Processor

kafkaApache Kafka

Input, Output

kafka_franzApache Kafka Kafka

Input, Output

microsoft_sql_server_cdc

Input

mongodb_cdcMongoDB CDC

Input

mqtt

Input, Output

mysql_cdc

Input

natsNATS.io

Input, Output

nats_jetstreamNATS JetStream NATS

Input, Output

oracledb_cdcOracle CDC OracleDB CDC Oracle Database CDC

Input

otlp_grpcOpenTelemetry OTLP OTel gRPC

Input, Output

otlp_httpOpenTelemetry OTLP OTel

Input, Output

pg_stream

postgres_cdc

Input

read_until

Input

redis_listRedis List Redis Lists Redis

Input, Output

redis_pubsubRedis PubSub Redis Pub/Sub Redis

Input, Output

redis_scanRedis

Input

redis_streamsRedis Streams Redis

Input, Output

redpanda_common

Input, Output

redpanda_migrator

Input, Output

resource

Input, Output, Processor

salesforce

Input

salesforce_cdcSalesforce Salesforce CDC

Input

salesforce_graphqlSalesforce Salesforce GraphQL

Input

schema_registry

Input, Output

sequence

Input

sftp

Input, Output

slack

Input

slack_usersSlack Users

Input

spicedb_watch

Input

splunk

Input

sql_rawSQL PostgreSQL MySQL Microsoft SQL Server ClickHouse Trino

Input, Output, Processor

sql_selectSQL PostgreSQL MySQL Microsoft SQL Server ClickHouse Trino

Input, Processor

timeplus

Input, Output

websocket

open_telemetry_collectorOpenTelemetry

Metric, Tracer

prometheus

Metric

arc

Output

aws_kinesis_firehoseAWS Kinesis Firehose Amazon Kinesis Firehose Kinesis Firehose

Output

aws_snsAWS SNS Amazon SNS SNS Simple Notification Service

Output

azure_data_lake_gen2Microsoft Azure Azure

Output

cache

Output, Processor

cyborgdb

Output

drop

Output

drop_on

Output

elasticsearch_v8

Output

fallback

Output

gcp_bigqueryGCP BigQuery Google BigQuery BigQuery

Output

gcp_bigquery_write_apiGCP BigQuery

Output

icebergApache Iceberg Apache Polaris AWS Glue Databricks Unity Catalog

Output

opensearch

Output

pinecone

Output

qdrant

Output, Processor

questdb

Output

redis_hashRedis Hash Redis

Output

reject

Output

reject_errored

Output

retry

Output, Processor

salesforce_sinkSalesforce Salesforce Sink

Output

slack_postSlack Post

Output

slack_reactionSlack Reaction

Output

snowflake_putSnowflake

Output

snowflake_streamingSnowflake Streaming

Output

splunk_hecSplunk

Output

sql_insertSQL PostgreSQL MySQL Microsoft SQL Server ClickHouse Trino

Output, Processor

switch

Output, Processor, Scanner

sync_response

Output, Processor

a2a_message

Processor

archiveZIP TAR GZIP

Processor

avro

Processor, Scanner

aws_bedrock_chatAmazon AWS Bedrock Chat

Processor

aws_bedrock_embeddingsAmazon AWS Bedrock Embeddings

Processor

aws_dynamodb_partiqlAmazon AWS DynamoDB PartiQL

Processor

aws_lambdaAWS Lambda Amazon Lambda Lambda

Processor

benchmark

Processor

bloblang

Processor

bounds_check

Processor

branch

Processor

cached

Processor

catch

Processor

cohere_chat

Processor

cohere_embeddings

Processor

cohere_rerank

Processor

compress

Processor

decompress

Processor, Scanner

dedupe

Processor

for_each

Processor

gcp_vertex_ai_chatGCP Vertex AI Google Cloud GCP

Processor

gcp_vertex_ai_embeddingsGoogle Cloud GCP

Processor

google_drive_download

Processor

google_drive_list_labels

Processor

google_drive_search

Processor

group_by

Processor

group_by_value

Processor

http

Processor

insert_part

Processor

jmespath

Processor

jq

Processor

json_schemaJSON Schema

Processor

log

Processor

mapping

Processor

metric

Processor

mutation

Processor

nats_request_replyNATS Request Reply

Processor

ollama_chat

Processor

ollama_embeddings

Processor

ollama_moderation

Processor

openai_chat_completion

Processor

openai_embeddings

Processor

openai_image_generation

Processor

openai_speech

Processor

openai_transcription

Processor

openai_translation

Processor

parallel

Processor

parquet_decode

Processor

parquet_encode

Processor

parse_log

Processor

processors

Processor

rate_limit

Processor

redis_scriptRedis Script

Processor

schema_registry_decode

Processor

schema_registry_encode

Processor

select_parts

Processor

slack_threadSlack Thread

Processor

sleep

Processor

split

Processor

string_split

Processor

text_chunker

Processor

try

Processor

try_catch

Processor

unarchiveZIP TAR GZIP Archive

Processor

while

Processor

workflow

Processor

xml

Processor

local

Rate_limit

chunker

Scanner

json_array

Scanner

json_documents

Scanner

lines

Scanner

re_match

Scanner

skip_bom

Scanner

tar

Scanner

to_the_end

Scanner

gcp_cloudtraceGCP Cloud Trace

Tracer

sql_driver_clickhouseClickHouse

sql_driver_mysqlMYSQL

sql_driver_oracleOracle

sql_driver_postgresPostgreSQL

sql_driver_sqliteSQLite

About Components

Every Redpanda Connect pipeline has at least one input, an optional buffer, an output and any number of processors:

input:
  kafka:
    addresses: [ TODO ]
    topics: [ foo, bar ]
    consumer_group: foogroup

buffer:
  type: none

pipeline:
  processors:
  - mapping: |
      message = this
      meta.link_count = links.length()

output:
  aws_s3:
    bucket: TODO
    path: '${! meta("kafka_topic") }/${! json("message.id") }.json'

These are the main components within Redpanda Connect and they provide the majority of useful behavior.

Observability components

There are also the observability components: logger, metrics, and tracing, which allow you to specify how Redpanda Connect exposes observability data.

http:
  address: 0.0.0.0:4195
  enabled: true
  debug_endpoints: false

logger:
  format: json
  level: WARN

metrics:
  statsd:
    address: localhost:8125
    flush_period: 100ms

tracer:
  jaeger:
    agent_address: localhost:6831

Resource components

Finally, there are caches and rate limits. These are components that are referenced by core components and can be shared.

input:
  http_client: # This is an input
    url: TODO
    rate_limit: foo_ratelimit # This is a reference to a rate limit

pipeline:
  processors:
    - cache: # This is a processor
        resource: baz_cache # This is a reference to a cache
        operator: add
        key: '${! json("id") }'
        value: "x"
    - mapping: root = if errored() { deleted() }

rate_limit_resources:
  - label: foo_ratelimit
    local:
      count: 500
      interval: 1s

cache_resources:
  - label: baz_cache
    memcached:
      addresses: [ localhost:11211 ]

It’s also possible to configure inputs, outputs and processors as resources which allows them to be reused throughout a configuration with the resource input, resource output and resource processor respectively.

For more information about any of these component types check out their sections: