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.
| Field | Type | Default | Description |
|---|---|---|---|
format | string | required | The archiving format to apply. |
path | string | "" | The path to set for each message in the archive (when applicable). |
avro
Performs Avro based operations on messages based on a schema.
Middleware: connect_avro, or inside a connect middleware.
| Field | Type | Default | Description |
|---|---|---|---|
operator | string | required | The operator to execute |
encoding | string | "textual" | An Avro encoding format to use for conversions to and from a schema. |
schema | string | "" | A full Avro schema to use. |
schema_path | string | "" | The path of a schema document to apply. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
codec | string | required | A codec defines how messages should be inserted into the AWK program as variables. |
program | string | required | An AWK program to execute |
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.
| Field | Type | Default | Description |
|---|---|---|---|
model | string | required | The model ID to use. |
prompt | string | The prompt you want to generate a response for. | |
system_prompt | string | The system prompt to submit to the AWS Bedrock LLM. | |
max_tokens | int | The maximum number of tokens to allow in the generated response. | |
temperature | float | The likelihood of the model selecting higher-probability options while generating a response. |
Advanced: region, endpoint, tcp, credentials, stop, top_p.
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.
| Field | Type | Default | Description |
|---|---|---|---|
model | string | required | The model ID to use. |
text | string | The text you want to compute vector embeddings for. | |
input_type | string | Specifies the type of input passed to the model. |
Advanced: region, endpoint, tcp, credentials.
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.
| Field | Type | Default | Description |
|---|---|---|---|
parallel | bool | false | Whether messages of a batch should be dispatched in parallel. |
function | string | required | The function to invoke. |
Advanced: rate_limit, region, endpoint, tcp, credentials, timeout, retries.
azure_cosmosdb
Creates or updates messages as JSON documents in Azure CosmosDB.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
endpoint | string | CosmosDB endpoint. | |
account_key | string | Account key. | |
connection_string | string | Connection string. | |
database | string | required | Database. |
container | string | required | Container. |
partition_keys_map | string | required | A Bloblang mapping which should evaluate to a single partition key value or an array of partition key values of type string, integer or boolean. |
operation | string | "Create" | Operation. |
item_id | string | ID of item to replace or delete. |
Advanced: patch_operations, patch_condition, auto_id, enable_content_response_on_write.
benchmark
Logs basic throughput statistics of messages that pass through this processor.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
interval | string | "5s" | How often to emit rolling statistics. |
count_bytes | bool | true | Whether or not to measure the number of bytes per second of throughput. |
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
max_part_size | int | 1073741824 | The maximum size of a message to allow (in bytes) |
min_part_size | int | 1 | The minimum size of a message to allow (in bytes) |
Advanced: max_parts, min_parts.
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.
| Field | Type | Default | Description |
|---|---|---|---|
request_map | string | "" | A Bloblang mapping that describes how to create a request payload suitable for the child processors of this branch. |
processors | list of processor | required | A list of processors to apply to mapped requests. |
result_map | string | "" | A Bloblang mapping that describes how the resulting messages from branched processing should be mapped back into the original payload. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
resource | string | required | The cache resource to target with this processor. |
operator | string | required | The operation to perform with the cache. |
key | string | required | A key to use with the cache. |
value | string | A value to use with the cache (when applicable). |
Advanced: ttl.
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.
| Field | Type | Default | Description |
|---|---|---|---|
cache | string | required | The cache resource to read and write processor results from. |
skip_on | string | A condition that can be used to skip caching the results from the processors. | |
key | string | required | A 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. |
ttl | string | An optional expiry period to set for each cache entry. | |
processors | list of processor | required | The list of processors whose result will be cached. |
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
base_url | string | "https://api.cohere.com" | The base URL to use for API requests. |
api_key | string | required | The API key for the Cohere API. |
model | string | required | The name of the Cohere model to use. |
prompt | string | The user prompt you want to generate a response for. | |
system_prompt | string | The system prompt to submit along with the user prompt. | |
max_tokens | int | The maximum number of tokens that can be generated in the chat completion. | |
temperature | float | What sampling temperature to use, between 0 and 2. | |
response_format | string | "text" | Specify the model’s output format. |
json_schema | string | The JSON schema to use when responding in json_schema format. | |
max_tool_calls | int | 10 | Maximum number of tool calls the model can do. |
tools | list of object | [] | The tools to allow the LLM to invoke. |
Advanced: schema_registry, top_p, frequency_penalty, presence_penalty, seed, stop.
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.
| Field | Type | Default | Description |
|---|---|---|---|
base_url | string | "https://api.cohere.com" | The base URL to use for API requests. |
api_key | string | required | The API key for the Cohere API. |
model | string | required | The name of the Cohere model to use. |
text_mapping | string | The text you want to generate a vector embedding for. | |
input_type | string | "search_document" | Specifies the type of input passed to the model. |
dimensions | int | The number of dimensions of the output embedding. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
base_url | string | "https://api.cohere.com" | The base URL to use for API requests. |
api_key | string | required | The API key for the Cohere API. |
model | string | required | The name of the Cohere model to use. |
query | string | required | The search query |
documents | string | required | A list of texts that will be compared to the query. |
top_n | string | "0" | The number of documents to return, if 0 all documents are returned. |
max_tokens_per_doc | int | 4096 | Long documents will be automatically truncated to the specified number of tokens. |
command
Executes a command for each message.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
name | string | required | The name of the command to execute. |
args_mapping | string | An optional Bloblang mapping that, when specified, should resolve into an array of arguments to pass to the command. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
algorithm | string | required | The compression algorithm to use. |
level | int | -1 | The level of compression to use. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | Couchbase connection string. |
username | string | Username to connect to the cluster. | |
password | string | Password to connect to the cluster. | |
bucket | string | required | Couchbase bucket. |
id | string | required | Document id. |
content | string | Document content. | |
operation | string | "get" | Couchbase operation to perform. |
Advanced: collection, scope, transcoder, timeout, ttl.
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
algorithm | string | required | The decompression algorithm to use. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
cache | string | required | The cache resource to target with this processor. |
key | string | required | An interpolated string yielding the key to deduplicate by for each message. |
drop_on_err | bool | true | Whether messages should be dropped when the cache returns a general error such as a network issue. |
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.
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.
| 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. | |
job_labels | map of string | {} | A list of labels to add to the query job. |
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_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.
| Field | Type | Default | Description |
|---|---|---|---|
project | string | required | GCP project ID to use |
credentials_json | string | An optional field to set google Service Account Credentials json. | |
location | string | required | The location of the model if using a fined tune model. |
model | string | required | The name of the LLM to use. |
prompt | string | The prompt you want to generate a response for. | |
history | string | Historical messages to include in the chat request. | |
attachment | string | Additional data like an image to send with the prompt to the model. | |
temperature | float | Controls the randomness of predications. | |
max_tokens | int | The maximum number of output tokens to generate per message. | |
response_format | string | "text" | The response format of generated type, the model must also be prompted to output the appropriate response type. |
tools | list 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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
project | string | required | GCP project ID to use |
credentials_json | string | An optional field to set google Service Account Credentials json. | |
location | string | "us-central1" | The location of the model. |
model | string | required | The name of the LLM to use. |
task_type | string | "RETRIEVAL_DOCUMENT" | The way to optimize embeddings that the model generates for specific use cases. |
text | string | The text you want to compute vector embeddings for. | |
output_dimensions | int | The maximum length for the output embedding size. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
expressions | list of string | required | One or more Grok expressions to attempt against incoming messages. |
pattern_definitions | map of string | {} | A map of pattern definitions that can be referenced within patterns. |
pattern_paths | list of string | [] | A list of paths to load Grok patterns from. |
Advanced: named_captures_only, use_default_patterns, remove_empty_values.
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.
| Field | Type | Default | Description |
|---|---|---|---|
check | string | required | A Bloblang query that should return a boolean value indicating whether a message belongs to a given group. |
processors | list of processor | [] | A list of processors to execute on the newly formed group. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
value | string | required | The interpolated string to group based on. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL to connect to. |
verb | string | "POST" | A verb to connect with |
headers | map of string | {} | A map of headers to add to the request. |
rate_limit | string | An optional rate limit to throttle requests by. | |
timeout | string | "5s" | A static timeout to apply to requests. |
parallel | bool | false | When 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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
index | int | -1 | The index within the batch to insert the message at. |
content | string | "" | The content of the message being inserted. |
javascript
Executes a provided JavaScript code block or file for each message.
Middleware: connect_javascript, or inside a connect middleware.
| Field | Type | Default | Description |
|---|---|---|---|
code | string | An inline JavaScript program to run. | |
file | string | A file containing a JavaScript program to run. | |
global_folders | list of string | [] | List of folders that will be used to load modules from if the requested JS module is not found elsewhere. |
jmespath
Executes a JMESPath query on JSON documents and replaces the message with the resulting document.
Middleware: connect_jmespath, or inside a connect middleware.
| Field | Type | Default | Description |
|---|---|---|---|
query | string | required | The JMESPath query to apply to messages. |
jq
Transforms and filters messages using jq queries.
Middleware: connect_jq, or inside a connect middleware.
| Field | Type | Default | Description |
|---|---|---|---|
query | string | required | The jq query to filter and transform messages with. |
Advanced: raw, output_raw.
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.
| Field | Type | Default | Description |
|---|---|---|---|
schema | string | A schema to apply. | |
schema_path | string | The path of a schema document to apply. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
level | string | "INFO" | The log level to use. |
fields_mapping | string | An optional Bloblang mapping that can be used to specify extra fields to add to the log. | |
message | string | "" | The message to print. |
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.
metric
Emit custom metrics by extracting values from messages.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
type | string | required | The metric type to create. |
name | string | required | The name of the metric to create, this must be unique across all Redpanda Connect components otherwise it will overwrite those other metrics. |
labels | map of string | A map of label names and values that can be used to enrich metrics. | |
value | string | "" | For some metric types specifies a value to set, increment. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target MongoDB server. |
database | string | required | The name of the target MongoDB database. |
username | string | "" | The username to connect to the database. |
password | string | "" | The password to connect to the database. |
collection | string | required | The name of the target collection. |
operation | string | "insert-one" | The mongodb operation to perform. |
write_concern | object | The write concern settings for the mongo connection. | |
document_map | string | "" | A bloblang map representing a document to store within MongoDB, expressed as extended JSON in canonical form. |
filter_map | string | "" | A bloblang map representing a filter for a MongoDB command, expressed as extended JSON in canonical form. |
hint_map | string | "" | A bloblang map representing the hint for the MongoDB command, expressed as extended JSON in canonical form. |
upsert | bool | false | The upsert setting is optional and only applies for update-one and replace-one operations. |
Advanced: app_name, aws, json_marshal_mode.
msgpack
Converts messages to or from the MessagePack format.
Middleware: connect_msgpack, or inside a connect middleware.
| Field | Type | Default | Description |
|---|---|---|---|
operator | string | required | The operation to perform on messages. |
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.
nats_kv
Perform operations on a NATS key-value bucket.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| 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. |
operation | string | required | The operation to perform on the KV bucket. |
key | string | required | The key for each message. |
Advanced: max_reconnects, revision, timeout, tls, tls_handshake_first, auth.
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.
| Field | Type | Default | Description |
|---|---|---|---|
urls | list of string | required | A list of URLs to connect to. |
subject | string | required | A subject to write to. |
headers | map of string | {} | Explicit message headers to add to messages. |
metadata | object | Determine which (if any) metadata values should be added to messages as headers. | |
timeout | string | "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.
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
model | string | required | The name of the Ollama LLM to use. |
prompt | string | The prompt you want to generate a response for. | |
image | string | The image to submit along with the prompt to the model. | |
response_format | string | "text" | The format of the response that the Ollama model generates. |
max_tokens | int | The maximum number of tokens to predict and output. | |
temperature | int | The temperature of the model. | |
save_prompt_metadata | bool | false | If enabled the prompt is saved as @prompt metadata on the output message. |
history | string | Historical messages to include in the chat request. | |
tools | list of object | [] | The tools to allow the LLM to invoke. |
runner | object | Options for the model runner that are used when the model is first loaded into memory. | |
server_address | string | The 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.
ollama_embeddings
Generates vector embeddings from text, using the Ollama API.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
model | string | required | The name of the Ollama LLM to use. |
text | string | The text you want to create vector embeddings for. | |
runner | object | Options for the model runner that are used when the model is first loaded into memory. | |
server_address | string | The address of the Ollama server to use. |
Advanced: cache_directory, download_url.
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.
| Field | Type | Default | Description |
|---|---|---|---|
model | string | required | The name of the Ollama LLM to use. |
prompt | string | required | The input prompt that was used with the LLM. |
response | string | required | The LLM’s response to classify if it contains safe or unsafe content. |
runner | object | Options for the model runner that are used when the model is first loaded into memory. | |
server_address | string | The address of the Ollama server to use. |
Advanced: cache_directory, download_url.
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.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
prompt | string | The user prompt you want to generate a response for. | |
system_prompt | string | The system prompt to submit along with the user prompt. | |
history | string | The history of the prior conversation. | |
image | string | An image to send along with the prompt. | |
max_tokens | int | The maximum number of tokens that can be generated in the chat completion. | |
temperature | float | What sampling temperature to use, between 0 and 2. | |
user | string | A unique identifier representing your end-user, which can help OpenAI to monitor and detect abuse. | |
response_format | string | "text" | Specify the model’s output format. |
json_schema | object | The JSON schema to use when responding in json_schema format. | |
tools | list of object | The tools to allow the LLM to invoke. |
Advanced: schema_registry, top_p, frequency_penalty, presence_penalty, seed, stop.
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.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
text_mapping | string | The text you want to generate a vector embedding for. | |
dimensions | int | The number of dimensions the resulting output embeddings should have. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
prompt | string | A text description of the image you want to generate. |
Advanced: quality, size, style.
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.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
input | string | A text description of the audio you want to generate. | |
voice | string | required | The type of voice to use when generating the audio. |
Advanced: response_format.
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.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
file | string | required | The 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.
openai_translation
Translates spoken audio into English, using the OpenAI API.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
server_address | string | "https://api.openai.com/v1" | The Open API endpoint that the processor sends requests to. |
api_key | string | required | The API key for OpenAI API. |
model | string | required | The name of the OpenAI model to use. |
file | string | The 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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
cap | int | 0 | The maximum number of messages to have processing at a given time. |
processors | list of processor | required | A list of child processors to apply. |
parquet
Deprecated. Converts batches of documents to or from Parquet files.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
operator | string | required | Determines whether the processor converts messages into a parquet file or expands parquet files into messages. |
compression | string | "snappy" | The type of compression to use when writing parquet files, this field is ignored when consuming parquet files. |
schema_file | string | A 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. | |
schema | string | 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. |
parquet_decode
Decodes Parquet files into a batch of structured messages.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
handle_logical_types | string | "v1" | Whether to be smart about decoding logical types. |
parquet_encode
Encodes Parquet files from a batch of structured messages.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
schema | list of object | Parquet schema. | |
schema_metadata | string | "" | Optionally specify a metadata field containing a schema definition to use for encoding instead of a statically defined schema. |
default_compression | string | "uncompressed" | The default compression type to use for fields. |
Advanced: default_encoding, default_timestamp_unit.
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.
| Field | Type | Default | Description |
|---|---|---|---|
format | string | required | A common log format to parse. |
Advanced: best_effort, allow_rfc3339, default_year, default_timezone.
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.
qdrant
Query items within a Qdrant collection.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
grpc_host | string | required | The gRPC host of the Qdrant server. |
api_token | string | "" | The Qdrant API token for authentication. |
collection_name | string | required | The name of the collection in Qdrant. |
vector_mapping | string | required | The mapping to extract the search vector from the document. |
filter | string | Additional filtering to perform on the results. | |
payload_fields | list of string | [] | The fields to include or exclude in returned result based on the payload_filter. |
payload_filter | string | "include" | The way the fields in payload_fields are filtered in the result. |
limit | int | 10 | The maximum number of points to return. |
Advanced: tls.
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.
| Field | Type | Default | Description |
|---|---|---|---|
resource | string | required | The target rate_limit resource. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
command | string | The command to execute. | |
args_mapping | string | A 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.
redis_script
Performs actions against Redis using LUA scripts.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
url | string | required | The URL of the target Redis server. |
script | string | required | A script to use for the target operator. |
args_mapping | string | required | A 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_mapping | string | required | A 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.
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.
retry
Beta. Attempts to execute a series of child processors until success.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
backoff | object | Determine time intervals and cut offs for retry attempts. | |
processors | list of processor | required | A list of processors to execute on each message. |
parallel | bool | false | When processing batches of messages these batches are ignored and the processors apply to each message sequentially. |
max_retries | int | 0 | The maximum number of retry attempts before the request is aborted. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
parts | list of int | [] | An array of message indexes of a batch. |
sentry_capture
Captures log events from messages and submits them to Sentry.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
dsn | string | "" | The DSN address to send sentry events to. |
message | string | required | A message to set on the sentry event |
context | string | A mapping that must evaluate to an object-of-objects or deleted(). | |
extras | string | A mapping that must evaluate to an object. | |
tags | map of string | Sets key/value string tags on an event. | |
environment | string | "" | The environment to be sent with events. |
release | string | "" | The version of the code deployed to an environment. |
level | string | "INFO" | Sets the level on sentry events similar to logging levels. |
transport_mode | string | "async" | Determines how events are sent. |
flush_timeout | string | "5s" | The duration to wait when closing the processor to flush any remaining enqueued events. |
sampling_rate | float | 1 | The rate at which events are sent to the server. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
duration | string | required | The duration of time to sleep for each execution. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
size | int | 1 | The target number of messages. |
byte_size | int | 0 | An optional target of total message bytes. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
data_source_name | string | required | Data source name. |
query | string | required | The query to execute. |
args_mapping | string | An optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query. | |
result_codec | string | "none" | Result codec. |
Advanced: unsafe_dynamic_query.
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.
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
dsn | string | required | A Data Source Name to identify the target database. |
table | string | required | The table to insert to. |
columns | list of string | required | A list of columns to insert. |
args_mapping | string | required | A Bloblang mapping which should evaluate to an array of values matching in size to the number of columns specified. |
Advanced: prefix, suffix, options, init_files, init_statement, conn_max_idle_time, conn_max_life_time, conn_max_idle, conn_max_open.
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.
| Field | Type | Default | Description |
|---|---|---|---|
driver | string | required | A database driver to use. |
dsn | string | required | A Data Source Name to identify the target database. |
query | string | The query to execute. | |
args_mapping | string | An optional Bloblang mapping which should evaluate to an array of values matching in size to the number of placeholder arguments in the field query. | |
exec_only | bool | Whether the query result should be discarded. | |
queries | list of object | A 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.
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.
| 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 query. |
columns | list of string | required | A list of columns to query. |
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. |
Advanced: prefix, suffix, init_files, init_statement, conn_max_idle_time, conn_max_life_time, conn_max_idle, conn_max_open.
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.
| Field | Type | Default | Description |
|---|---|---|---|
delimiter | string | "\n" | The delimiter to split the string by. |
empty_as_null | bool | false | When true, empty strings resulting from the split are converted to null. |
Advanced: emit_bytes.
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.
| 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. |
Advanced: max_buffer, codec_send, codec_recv.
switch
Conditionally processes messages based on their contents.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
check | string | "" | A Bloblang query that should return a boolean value indicating whether a message should have the processors of this case executed on it. |
processors | list of processor | [] | A list of processors to execute on a message. |
Advanced: fallthrough, continue.
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
strategy | string | required | |
chunk_size | int | 512 | The maximum size of each chunk. |
chunk_overlap | int | 100 | The number of characters to overlap between chunks. |
separators | list of string | ["\n\n", "\n", " ", ""] | A list of strings that should be considered as separators between chunks. |
length_measure | string | "runes" | The method for measuring the length of a string. |
include_code_blocks | bool | false | Whether to include code blocks in the output. |
keep_reference_links | bool | false | Whether to keep reference links in the output. |
Advanced: token_encoding, allowed_special, disallowed_special.
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.
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.
| Field | Type | Default | Description |
|---|---|---|---|
processors | list of processor | [] | A list of processors to execute on each message. |
catch | list of processor | [] | A list of processors to execute on each message that failed one of the processors above. |
error_metadata | string | "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. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
format | string | required | The unarchiving format to apply. |
wasm
Executes a function exported by a WASM module for each message.
Use it inside a connect middleware or a connect endpoint’s pipeline.
| Field | Type | Default | Description |
|---|---|---|---|
module_path | string | required | The path of the target WASM module to execute. |
function | string | "process" | The name of the function exported by the target WASM module to run for each message. |
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.
| Field | Type | Default | Description |
|---|---|---|---|
at_least_once | bool | false | Whether to always run the child processors at least one time. |
check | string | "" | A Bloblang query that should return a boolean value indicating whether the while loop should execute again. |
processors | list of processor | required | A list of child processors to execute on each loop. |
Advanced: max_loops.
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.
| Field | Type | Default | Description |
|---|---|---|---|
meta_path | string | "meta.workflow" | A dot path indicating where to store and reference structured metadata about the workflow execution. |
order | string | [] | An explicit declaration of branch ordered tiers, which describes the order in which parallel tiers of branches should be executed. |
branches | map of object | {} | An object of named branch processors that make up the workflow. |
Advanced: branch_resources.
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.
| Field | Type | Default | Description |
|---|---|---|---|
operator | string | "" | An XML operation to apply to messages. |
cast | bool | false | Whether to try to cast values that are numbers and booleans to the right type. |