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

The 86 Redpanda Connect processors 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.

16 of them are also exported as their own connect_<name> middleware and run on any endpoint, native ones included. Every other processor that keeps, rewrites or drops each message runs inside a connect middleware; one that splits or merges messages belongs in a connect endpoint’s pipeline. Middleware needs plugin 0.1.1 or newer.

Own middleware: connect_avro, connect_bloblang, connect_branch, connect_cached, connect_dedupe, connect_grok, connect_http, connect_javascript, connect_jmespath, connect_jq, connect_json_schema, connect_log, connect_mapping, connect_msgpack, connect_mutation, connect_parse_log.

input:
  kafka: { url: localhost:9092, topic: orders }
  middlewares:
    - connect_mapping: 'root = this.merge({"received_at": now()})'

Processors in bold are also their own middleware.

archive

Archives all the messages of a batch into a single message according to the selected archive format.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
formatstringrequiredThe archiving format to apply.
pathstring""The path to set for each message in the archive (when applicable).

Upstream documentation

avro

Performs Avro based operations on messages based on a schema.

Middleware: connect_avro, or inside a connect middleware.

FieldTypeDefaultDescription
operatorstringrequiredThe operator to execute
encodingstring"textual"An Avro encoding format to use for conversions to and from a schema.
schemastring""A full Avro schema to use.
schema_pathstring""The path of a schema document to apply.

Upstream documentation

awk

Executes an AWK program on messages. This processor is very powerful as it offers a range of custom functions for querying and mutating message contents and metadata.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
codecstringrequiredA codec defines how messages should be inserted into the AWK program as variables.
programstringrequiredAn AWK program to execute

Upstream documentation

aws_bedrock_chat

Generates responses to messages in a chat conversation, using the AWS Bedrock API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
modelstringrequiredThe model ID to use.
promptstringThe prompt you want to generate a response for.
system_promptstringThe system prompt to submit to the AWS Bedrock LLM.
max_tokensintThe maximum number of tokens to allow in the generated response.
temperaturefloatThe likelihood of the model selecting higher-probability options while generating a response.

Advanced: region, endpoint, tcp, credentials, stop, top_p.

Upstream documentation

aws_bedrock_embeddings

Computes vector embeddings on text, using the AWS Bedrock API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
modelstringrequiredThe model ID to use.
textstringThe text you want to compute vector embeddings for.
input_typestringSpecifies the type of input passed to the model.

Advanced: region, endpoint, tcp, credentials.

Upstream documentation

aws_lambda

Invokes an AWS lambda for each message. The contents of the message is the payload of the request, and the result of the invocation will become the new contents of the message.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
parallelboolfalseWhether messages of a batch should be dispatched in parallel.
functionstringrequiredThe function to invoke.

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

Upstream documentation

azure_cosmosdb

Creates or updates messages as JSON documents in Azure CosmosDB.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
endpointstringCosmosDB endpoint.
account_keystringAccount key.
connection_stringstringConnection string.
databasestringrequiredDatabase.
containerstringrequiredContainer.
partition_keys_mapstringrequiredA Bloblang mapping which should evaluate to a single partition key value or an array of partition key values of type string, integer or boolean.
operationstring"Create"Operation.
item_idstringID of item to replace or delete.

Advanced: patch_operations, patch_condition, auto_id, enable_content_response_on_write.

Upstream documentation

benchmark

Logs basic throughput statistics of messages that pass through this processor.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
intervalstring"5s"How often to emit rolling statistics.
count_bytesbooltrueWhether or not to measure the number of bytes per second of throughput.

Upstream documentation

bloblang

Executes a Bloblang mapping on messages.

Middleware: connect_bloblang, or inside a connect middleware.

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

Upstream documentation

bounds_check

Removes messages (and batches) that do not fit within certain size boundaries.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
max_part_sizeint1073741824The maximum size of a message to allow (in bytes)
min_part_sizeint1The minimum size of a message to allow (in bytes)

Advanced: max_parts, min_parts.

Upstream documentation

branch

The branch processor allows you to create a new request message via a Bloblang mapping, execute a list of processors on the request messages, and, finally, map the result back into the source message using another mapping.

Middleware: connect_branch, or inside a connect middleware.

FieldTypeDefaultDescription
request_mapstring""A Bloblang mapping that describes how to create a request payload suitable for the child processors of this branch.
processorslist of processorrequiredA list of processors to apply to mapped requests.
result_mapstring""A Bloblang mapping that describes how the resulting messages from branched processing should be mapped back into the original payload.

Upstream documentation

cache

Performs operations against a cache resource for each message, allowing you to store or retrieve data within message payloads.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
resourcestringrequiredThe cache resource to target with this processor.
operatorstringrequiredThe operation to perform with the cache.
keystringrequiredA key to use with the cache.
valuestringA value to use with the cache (when applicable).

Advanced: ttl.

Upstream documentation

cached

Cache the result of applying one or more processors to messages identified by a key. If the key already exists within the cache the contents of the message will be replaced with the cached result instead of applying the processors. This component is therefore useful in situations where an expensive set of processors need only be executed periodically.

Middleware: connect_cached, or inside a connect middleware.

FieldTypeDefaultDescription
cachestringrequiredThe cache resource to read and write processor results from.
skip_onstringA condition that can be used to skip caching the results from the processors.
keystringrequiredA key to be resolved for each message, if the key already exists in the cache then the cached result is used, otherwise the processors are applied and the result is cached under this key.
ttlstringAn optional expiry period to set for each cache entry.
processorslist of processorrequiredThe list of processors whose result will be cached.

Upstream documentation

catch

Applies a list of child processors only when a previous processing step has failed.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

cohere_chat

Generates responses to messages in a chat conversation, using the Cohere API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
base_urlstring"https://api.cohere.com"The base URL to use for API requests.
api_keystringrequiredThe API key for the Cohere API.
modelstringrequiredThe name of the Cohere model to use.
promptstringThe user prompt you want to generate a response for.
system_promptstringThe system prompt to submit along with the user prompt.
max_tokensintThe maximum number of tokens that can be generated in the chat completion.
temperaturefloatWhat sampling temperature to use, between 0 and 2.
response_formatstring"text"Specify the model’s output format.
json_schemastringThe JSON schema to use when responding in json_schema format.
max_tool_callsint10Maximum number of tool calls the model can do.
toolslist of object[]The tools to allow the LLM to invoke.

Advanced: schema_registry, top_p, frequency_penalty, presence_penalty, seed, stop.

Upstream documentation

cohere_embeddings

Generates vector embeddings to represent input text, using the Cohere API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
base_urlstring"https://api.cohere.com"The base URL to use for API requests.
api_keystringrequiredThe API key for the Cohere API.
modelstringrequiredThe name of the Cohere model to use.
text_mappingstringThe text you want to generate a vector embedding for.
input_typestring"search_document"Specifies the type of input passed to the model.
dimensionsintThe number of dimensions of the output embedding.

Upstream documentation

cohere_rerank

Ranks a list of documents by their relevance to a query, using the Cohere API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
base_urlstring"https://api.cohere.com"The base URL to use for API requests.
api_keystringrequiredThe API key for the Cohere API.
modelstringrequiredThe name of the Cohere model to use.
querystringrequiredThe search query
documentsstringrequiredA list of texts that will be compared to the query.
top_nstring"0"The number of documents to return, if 0 all documents are returned.
max_tokens_per_docint4096Long documents will be automatically truncated to the specified number of tokens.

Upstream documentation

command

Executes a command for each message.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
namestringrequiredThe name of the command to execute.
args_mappingstringAn optional Bloblang mapping that, when specified, should resolve into an array of arguments to pass to the command.

Upstream documentation

compress

Compresses messages according to the selected algorithm. Supported compression algorithms are: [flate gzip lz4 pgzip snappy zlib]

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
algorithmstringrequiredThe compression algorithm to use.
levelint-1The level of compression to use.

Upstream documentation

couchbase

Performs operations against Couchbase for each message, allowing you to store or retrieve data within message payloads.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
urlstringrequiredCouchbase connection string.
usernamestringUsername to connect to the cluster.
passwordstringPassword to connect to the cluster.
bucketstringrequiredCouchbase bucket.
idstringrequiredDocument id.
contentstringDocument content.
operationstring"get"Couchbase operation to perform.

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

Upstream documentation

crash

Beta. Crashes the process using a fatal log message. The log message can be set using function interpolations described in Bloblang queries which allows you to log the contents and metadata of messages.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

decompress

Decompresses messages according to the selected algorithm. Supported decompression algorithms are: [bzip2 flate gzip lz4 pgzip snappy zlib]

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
algorithmstringrequiredThe decompression algorithm to use.

Upstream documentation

dedupe

Deduplicates messages by storing a key value in a cache using the add operator. If the key already exists within the cache it is dropped.

Middleware: connect_dedupe, or inside a connect middleware.

FieldTypeDefaultDescription
cachestringrequiredThe cache resource to target with this processor.
keystringrequiredAn interpolated string yielding the key to deduplicate by for each message.
drop_on_errbooltrueWhether messages should be dropped when the cache returns a general error such as a network issue.

Upstream documentation

for_each

A processor that applies a list of child processors to messages of a batch as though they were each a batch of one message.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

gcp_bigquery_select

Executes a SELECT query against BigQuery and replaces messages with the rows returned.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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.
job_labelsmap of string{}A list of labels to add to the query job.
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_vertex_ai_chat

Generates responses to messages in a chat conversation, using the Vertex AI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
projectstringrequiredGCP project ID to use
credentials_jsonstringAn optional field to set google Service Account Credentials json.
locationstringrequiredThe location of the model if using a fined tune model.
modelstringrequiredThe name of the LLM to use.
promptstringThe prompt you want to generate a response for.
historystringHistorical messages to include in the chat request.
attachmentstringAdditional data like an image to send with the prompt to the model.
temperaturefloatControls the randomness of predications.
max_tokensintThe maximum number of output tokens to generate per message.
response_formatstring"text"The response format of generated type, the model must also be prompted to output the appropriate response type.
toolslist of object[]The tools to allow the LLM to invoke.

Advanced: system_prompt, top_p, top_k, stop, presence_penalty, frequency_penalty, max_tool_calls.

Upstream documentation

gcp_vertex_ai_embeddings

Generates vector embeddings to represent input text, using the Vertex AI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
projectstringrequiredGCP project ID to use
credentials_jsonstringAn optional field to set google Service Account Credentials json.
locationstring"us-central1"The location of the model.
modelstringrequiredThe name of the LLM to use.
task_typestring"RETRIEVAL_DOCUMENT"The way to optimize embeddings that the model generates for specific use cases.
textstringThe text you want to compute vector embeddings for.
output_dimensionsintThe maximum length for the output embedding size.

Upstream documentation

grok

Parses messages into a structured format by attempting to apply a list of Grok expressions, the first expression to result in at least one value replaces the original message with a JSON object containing the values.

Middleware: connect_grok, or inside a connect middleware.

FieldTypeDefaultDescription
expressionslist of stringrequiredOne or more Grok expressions to attempt against incoming messages.
pattern_definitionsmap of string{}A map of pattern definitions that can be referenced within patterns.
pattern_pathslist of string[]A list of paths to load Grok patterns from.

Advanced: named_captures_only, use_default_patterns, remove_empty_values.

Upstream documentation

group_by

Splits a batch of messages into N batches, where each resulting batch contains a group of messages determined by a Bloblang query.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
checkstringrequiredA Bloblang query that should return a boolean value indicating whether a message belongs to a given group.
processorslist of processor[]A list of processors to execute on the newly formed group.

Upstream documentation

group_by_value

Splits a batch of messages into N batches, where each resulting batch contains a group of messages determined by a function interpolated string evaluated per message.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
valuestringrequiredThe interpolated string to group based on.

Upstream documentation

http

Performs an HTTP request using a message batch as the request body, and replaces the original message parts with the body of the response.

Middleware: connect_http, or inside a connect middleware.

FieldTypeDefaultDescription
urlstringrequiredThe URL to connect to.
verbstring"POST"A verb to connect with
headersmap of string{}A map of headers to add to the request.
rate_limitstringAn optional rate limit to throttle requests by.
timeoutstring"5s"A static timeout to apply to requests.
parallelboolfalseWhen processing batched messages, whether to send messages of the batch in parallel, otherwise they are sent serially.

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.

Upstream documentation

insert_part

Insert a new message into a batch at an index. If the specified index is greater than the length of the existing batch it will be appended to the end.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
indexint-1The index within the batch to insert the message at.
contentstring""The content of the message being inserted.

Upstream documentation

javascript

Executes a provided JavaScript code block or file for each message.

Middleware: connect_javascript, or inside a connect middleware.

FieldTypeDefaultDescription
codestringAn inline JavaScript program to run.
filestringA file containing a JavaScript program to run.
global_folderslist of string[]List of folders that will be used to load modules from if the requested JS module is not found elsewhere.

Upstream documentation

jmespath

Executes a JMESPath query on JSON documents and replaces the message with the resulting document.

Middleware: connect_jmespath, or inside a connect middleware.

FieldTypeDefaultDescription
querystringrequiredThe JMESPath query to apply to messages.

Upstream documentation

jq

Transforms and filters messages using jq queries.

Middleware: connect_jq, or inside a connect middleware.

FieldTypeDefaultDescription
querystringrequiredThe jq query to filter and transform messages with.

Advanced: raw, output_raw.

Upstream documentation

json_schema

Checks messages against a provided JSONSchema definition but does not change the payload under any circumstances. If a message does not match the schema it can be caught using error handling methods.

Middleware: connect_json_schema, or inside a connect middleware.

FieldTypeDefaultDescription
schemastringA schema to apply.
schema_pathstringThe path of a schema document to apply.

Upstream documentation

log

Prints a log event for each message. Messages always remain unchanged. The log message can be set using function interpolations described in Bloblang queries which allows you to log the contents and metadata of messages.

Middleware: connect_log, or inside a connect middleware.

FieldTypeDefaultDescription
levelstring"INFO"The log level to use.
fields_mappingstringAn optional Bloblang mapping that can be used to specify extra fields to add to the log.
messagestring""The message to print.

Upstream documentation

mapping

Executes a Bloblang mapping on messages, creating a new document that replaces (or filters) the original message.

Middleware: connect_mapping, or inside a connect middleware.

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

Upstream documentation

metric

Emit custom metrics by extracting values from messages.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
typestringrequiredThe metric type to create.
namestringrequiredThe name of the metric to create, this must be unique across all Redpanda Connect components otherwise it will overwrite those other metrics.
labelsmap of stringA map of label names and values that can be used to enrich metrics.
valuestring""For some metric types specifies a value to set, increment.

Upstream documentation

mongodb

Performs operations against MongoDB for each message, allowing you to store or retrieve data within message payloads.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Advanced: app_name, aws, json_marshal_mode.

Upstream documentation

msgpack

Converts messages to or from the MessagePack format.

Middleware: connect_msgpack, or inside a connect middleware.

FieldTypeDefaultDescription
operatorstringrequiredThe operation to perform on messages.

Upstream documentation

mutation

Executes a Bloblang mapping and directly transforms the contents of messages, mutating (or deleting) them.

Middleware: connect_mutation, or inside a connect middleware.

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

Upstream documentation

nats_kv

Perform operations on a NATS key-value bucket.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
bucketstringrequiredThe name of the KV bucket.
operationstringrequiredThe operation to perform on the KV bucket.
keystringrequiredThe key for each message.

Advanced: max_reconnects, revision, timeout, tls, tls_handshake_first, auth.

Upstream documentation

nats_request_reply

Sends a message to a NATS subject and expects a reply, from a NATS subscriber acting as a responder, back.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
urlslist of stringrequiredA list of URLs to connect to.
subjectstringrequiredA subject to write to.
headersmap of string{}Explicit message headers to add to messages.
metadataobjectDetermine which (if any) metadata values should be added to messages as headers.
timeoutstring"3s"A duration string is a possibly signed sequence of decimal numbers, each with optional fraction and a unit suffix, such as 300ms, -1.5h or 2h45m.

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

Upstream documentation

noop

Noop is a processor that does nothing, the message passes through unchanged. Why? Sometimes doing nothing is the braver option.

Use it inside a connect middleware or a connect endpoint’s pipeline.

Upstream documentation

ollama_chat

Generates responses to messages in a chat conversation, using the Ollama API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
modelstringrequiredThe name of the Ollama LLM to use.
promptstringThe prompt you want to generate a response for.
imagestringThe image to submit along with the prompt to the model.
response_formatstring"text"The format of the response that the Ollama model generates.
max_tokensintThe maximum number of tokens to predict and output.
temperatureintThe temperature of the model.
save_prompt_metadataboolfalseIf enabled the prompt is saved as @prompt metadata on the output message.
historystringHistorical messages to include in the chat request.
toolslist of object[]The tools to allow the LLM to invoke.
runnerobjectOptions for the model runner that are used when the model is first loaded into memory.
server_addressstringThe address of the Ollama server to use.

Advanced: system_prompt, num_keep, seed, top_k, top_p, repeat_penalty, presence_penalty, frequency_penalty, stop, max_tool_calls, cache_directory, download_url.

Upstream documentation

ollama_embeddings

Generates vector embeddings from text, using the Ollama API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
modelstringrequiredThe name of the Ollama LLM to use.
textstringThe text you want to create vector embeddings for.
runnerobjectOptions for the model runner that are used when the model is first loaded into memory.
server_addressstringThe address of the Ollama server to use.

Advanced: cache_directory, download_url.

Upstream documentation

ollama_moderation

Classifies an LLM response as safe or unsafe, using the Ollama API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
modelstringrequiredThe name of the Ollama LLM to use.
promptstringrequiredThe input prompt that was used with the LLM.
responsestringrequiredThe LLM’s response to classify if it contains safe or unsafe content.
runnerobjectOptions for the model runner that are used when the model is first loaded into memory.
server_addressstringThe address of the Ollama server to use.

Advanced: cache_directory, download_url.

Upstream documentation

openai_chat_completion

Generates responses to messages in a chat conversation, using the OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
promptstringThe user prompt you want to generate a response for.
system_promptstringThe system prompt to submit along with the user prompt.
historystringThe history of the prior conversation.
imagestringAn image to send along with the prompt.
max_tokensintThe maximum number of tokens that can be generated in the chat completion.
temperaturefloatWhat sampling temperature to use, between 0 and 2.
userstringA unique identifier representing your end-user, which can help OpenAI to monitor and detect abuse.
response_formatstring"text"Specify the model’s output format.
json_schemaobjectThe JSON schema to use when responding in json_schema format.
toolslist of objectThe tools to allow the LLM to invoke.

Advanced: schema_registry, top_p, frequency_penalty, presence_penalty, seed, stop.

Upstream documentation

openai_embeddings

Generates vector embeddings to represent input text, using the OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
text_mappingstringThe text you want to generate a vector embedding for.
dimensionsintThe number of dimensions the resulting output embeddings should have.

Upstream documentation

openai_image_generation

Generates an image from a text description and other attributes, using OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
promptstringA text description of the image you want to generate.

Advanced: quality, size, style.

Upstream documentation

openai_speech

Generates audio from a text description and other attributes, using OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
inputstringA text description of the audio you want to generate.
voicestringrequiredThe type of voice to use when generating the audio.

Advanced: response_format.

Upstream documentation

openai_transcription

Generates a transcription of spoken audio in the input language, using the OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
filestringrequiredThe audio file object (not file name) to transcribe, in one of the following formats: flac, mp3, mp4, mpeg, mpga, m4a, ogg, wav, or webm.

Advanced: language, prompt.

Upstream documentation

openai_translation

Translates spoken audio into English, using the OpenAI API.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
server_addressstring"https://api.openai.com/v1"The Open API endpoint that the processor sends requests to.
api_keystringrequiredThe API key for OpenAI API.
modelstringrequiredThe name of the OpenAI model to use.
filestringThe audio file object (not file name) to translate, in one of the following formats: flac, mp3, mp4, mpeg, mpga, m4a, ogg, wav, or webm.

Advanced: prompt.

Upstream documentation

parallel

A processor that applies a list of child processors to messages of a batch as though they were each a batch of one message (similar to the for_each processor), but where each message is processed in parallel.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
capint0The maximum number of messages to have processing at a given time.
processorslist of processorrequiredA list of child processors to apply.

Upstream documentation

parquet

Deprecated. Converts batches of documents to or from Parquet files.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
operatorstringrequiredDetermines whether the processor converts messages into a parquet file or expands parquet files into messages.
compressionstring"snappy"The type of compression to use when writing parquet files, this field is ignored when consuming parquet files.
schema_filestringA file path containing a schema used to describe the parquet files being generated or consumed, the format of the schema is a JSON document detailing the tag and fields of documents.
schemastringA schema used to describe the parquet files being generated or consumed, the format of the schema is a JSON document detailing the tag and fields of documents.

Upstream documentation

parquet_decode

Decodes Parquet files into a batch of structured messages.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
handle_logical_typesstring"v1"Whether to be smart about decoding logical types.

Upstream documentation

parquet_encode

Encodes Parquet files from a batch of structured messages.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
schemalist of objectParquet schema.
schema_metadatastring""Optionally specify a metadata field containing a schema definition to use for encoding instead of a statically defined schema.
default_compressionstring"uncompressed"The default compression type to use for fields.

Advanced: default_encoding, default_timestamp_unit.

Upstream documentation

parse_log

Parses common log formats into structured data. This is easier and often much faster than grok.

Middleware: connect_parse_log, or inside a connect middleware.

FieldTypeDefaultDescription
formatstringrequiredA common log format to parse.

Advanced: best_effort, allow_rfc3339, default_year, default_timezone.

Upstream documentation

processors

A processor grouping several sub-processors.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

qdrant

Query items within a Qdrant collection.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
grpc_hoststringrequiredThe gRPC host of the Qdrant server.
api_tokenstring""The Qdrant API token for authentication.
collection_namestringrequiredThe name of the collection in Qdrant.
vector_mappingstringrequiredThe mapping to extract the search vector from the document.
filterstringAdditional filtering to perform on the results.
payload_fieldslist of string[]The fields to include or exclude in returned result based on the payload_filter.
payload_filterstring"include"The way the fields in payload_fields are filtered in the result.
limitint10The maximum number of points to return.

Advanced: tls.

Upstream documentation

rate_limit

Throttles the throughput of a pipeline according to a specified rate_limit resource. Rate limits are shared across components and therefore apply globally to all processing pipelines.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
resourcestringrequiredThe target rate_limit resource.

Upstream documentation

redis

Performs actions against Redis that aren’t possible using a cache processor. Actions are performed for each message and the message contents are replaced with the result. In order to merge the result into the original message compose this processor within a branch processor.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
commandstringThe command to execute.
args_mappingstringA Bloblang mapping which should evaluate to an array of values matching in size to the number of arguments required for the specified Redis command.

Advanced: kind, master, client_name, tls, retries, retry_period.

Upstream documentation

redis_script

Performs actions against Redis using LUA scripts.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
urlstringrequiredThe URL of the target Redis server.
scriptstringrequiredA script to use for the target operator.
args_mappingstringrequiredA Bloblang mapping which should evaluate to an array of values matching in size to the number of arguments required for the specified Redis script.
keys_mappingstringrequiredA Bloblang mapping which should evaluate to an array of keys matching in size to the number of arguments required for the specified Redis script.

Advanced: kind, master, client_name, tls, retries, retry_period.

Upstream documentation

resource

Resource is a processor type that runs a processor resource identified by its label.

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

retry

Beta. Attempts to execute a series of child processors until success.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
backoffobjectDetermine time intervals and cut offs for retry attempts.
processorslist of processorrequiredA list of processors to execute on each message.
parallelboolfalseWhen processing batches of messages these batches are ignored and the processors apply to each message sequentially.
max_retriesint0The maximum number of retry attempts before the request is aborted.

Upstream documentation

select_parts

Cherry pick a set of messages from a batch by their index. Indexes larger than the number of messages are simply ignored.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
partslist of int[]An array of message indexes of a batch.

Upstream documentation

sentry_capture

Captures log events from messages and submits them to Sentry.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
dsnstring""The DSN address to send sentry events to.
messagestringrequiredA message to set on the sentry event
contextstringA mapping that must evaluate to an object-of-objects or deleted().
extrasstringA mapping that must evaluate to an object.
tagsmap of stringSets key/value string tags on an event.
environmentstring""The environment to be sent with events.
releasestring""The version of the code deployed to an environment.
levelstring"INFO"Sets the level on sentry events similar to logging levels.
transport_modestring"async"Determines how events are sent.
flush_timeoutstring"5s"The duration to wait when closing the processor to flush any remaining enqueued events.
sampling_ratefloat1The rate at which events are sent to the server.

Upstream documentation

sleep

Sleep for a period of time specified as a duration string for each message. This processor will interpolate functions within the duration field, you can find a list of functions here.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
durationstringrequiredThe duration of time to sleep for each execution.

Upstream documentation

split

Breaks message batches (synonymous with multiple part messages) into smaller batches. The size of the resulting batches are determined either by a discrete size or, if the field byte_size is non-zero, then by total size in bytes (which ever limit is reached first).

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
sizeint1The target number of messages.
byte_sizeint0An optional target of total message bytes.

Upstream documentation

sql

Deprecated. Runs an arbitrary SQL query against a database and (optionally) returns the result as an array of objects, one for each row returned.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
data_source_namestringrequiredData source name.
querystringrequiredThe query to execute.
args_mappingstringAn optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query.
result_codecstring"none"Result codec.

Advanced: unsafe_dynamic_query.

Upstream documentation

sql_insert

Inserts rows into an SQL database for each message, and leaves the message unchanged.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
tablestringrequiredThe table to insert to.
columnslist of stringrequiredA list of columns to insert.
args_mappingstringrequiredA Bloblang mapping which should evaluate to an array of values matching in size to the number of columns specified.

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

Upstream documentation

sql_raw

Runs an arbitrary SQL query against a database and (optionally) returns the result as an array of objects, one for each row returned.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
querystringThe query to execute.
args_mappingstringAn optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query.
exec_onlyboolWhether the query result should be discarded.
querieslist of objectA list of statements to run in addition to query.

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

Upstream documentation

sql_select

Runs an SQL select query against a database and returns the result as an array of objects, one for each row returned, containing a key for each column queried and its value.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
driverstringrequiredA database driver to use.
dsnstringrequiredA Data Source Name to identify the target database.
tablestringrequiredThe table to query.
columnslist of stringrequiredA list of columns to query.
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.

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

Upstream documentation

string_split

Splits a string by a delimiter into an array. Generally, using bloblang’s split method is preferred. In some high performance use cases this processor can be faster than the equivalent bloblang if there is no additional logic.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
delimiterstring"\n"The delimiter to split the string by.
empty_as_nullboolfalseWhen true, empty strings resulting from the split are converted to null.

Advanced: emit_bytes.

Upstream documentation

subprocess

Executes a command as a subprocess and, for each message, will pipe its contents to the stdin stream of the process followed by a newline.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
namestringrequiredThe command to execute as a subprocess.
argslist of string[]A list of arguments to provide the command.

Advanced: max_buffer, codec_send, codec_recv.

Upstream documentation

switch

Conditionally processes messages based on their contents.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
checkstring""A Bloblang query that should return a boolean value indicating whether a message should have the processors of this case executed on it.
processorslist of processor[]A list of processors to execute on a message.

Advanced: fallthrough, continue.

Upstream documentation

sync_response

Adds the payload in its current state as a synchronous response to the input source, where it is dealt with according to that specific input type.

Use it inside a connect middleware or a connect endpoint’s pipeline.

Upstream documentation

text_chunker

A processor that allows chunking and splitting text based on some strategy. Usually used for creating vector embeddings of large documents.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
strategystringrequired
chunk_sizeint512The maximum size of each chunk.
chunk_overlapint100The number of characters to overlap between chunks.
separatorslist of string["\n\n", "\n", " ", ""]A list of strings that should be considered as separators between chunks.
length_measurestring"runes"The method for measuring the length of a string.
include_code_blocksboolfalseWhether to include code blocks in the output.
keep_reference_linksboolfalseWhether to keep reference links in the output.

Advanced: token_encoding, allowed_special, disallowed_special.

Upstream documentation

try

Executes a list of child processors on messages only if no prior processors have failed (or the errors have been cleared).

Use it inside a connect middleware or a connect endpoint’s pipeline.

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

Upstream documentation

try_catch

Beta. Executes a list of child processors on each message and, if any of them fail, executes a separate list of catch processors to recover from or react to the error.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
processorslist of processor[]A list of processors to execute on each message.
catchlist of processor[]A list of processors to execute on each message that failed one of the processors above.
error_metadatastring"error"The metadata key under which the caught error is stored, as an object with a what field (the error message) plus name, label and path fields describing the component that failed, before the catch processors are executed.

Upstream documentation

unarchive

Unarchives messages according to the selected archive format into multiple messages within a batch.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
formatstringrequiredThe unarchiving format to apply.

Upstream documentation

wasm

Executes a function exported by a WASM module for each message.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
module_pathstringrequiredThe path of the target WASM module to execute.
functionstring"process"The name of the function exported by the target WASM module to run for each message.

Upstream documentation

while

A processor that checks a Bloblang query against each batch of messages and executes child processors on them for as long as the query resolves to true.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
at_least_onceboolfalseWhether to always run the child processors at least one time.
checkstring""A Bloblang query that should return a boolean value indicating whether the while loop should execute again.
processorslist of processorrequiredA list of child processors to execute on each loop.

Advanced: max_loops.

Upstream documentation

workflow

Executes a topology of branch processors, performing them in parallel where possible.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
meta_pathstring"meta.workflow"A dot path indicating where to store and reference structured metadata about the workflow execution.
orderstring[]An explicit declaration of branch ordered tiers, which describes the order in which parallel tiers of branches should be executed.
branchesmap of object{}An object of named branch processors that make up the workflow.

Advanced: branch_resources.

Upstream documentation

xml

Parses messages as an XML document, performs a mutation on the data, and then overwrites the previous contents with the new value.

Use it inside a connect middleware or a connect endpoint’s pipeline.

FieldTypeDefaultDescription
operatorstring""An XML operation to apply to messages.
castboolfalseWhether to try to cast values that are numbers and booleans to the right type.

Upstream documentation