Skip to main content

aws_dynamodb_partiql processor

Executes a PartiQL expression against a DynamoDB table for each message.

# Common config fields, showing default values
pipeline:
processors:
- label: ""
aws_dynamodb_partiql:
query: "" # No default (required)
args_mapping: ""

Both writes or reads are supported, when the query is a read the contents of the message will be replaced with the result. This processor is more efficient when messages are pre-batched as the whole batch will be executed in a single call.

Examples

Insert

The following example inserts rows into the table footable with the columns foo, bar and baz populated with values extracted from messages:

pipeline:
processors:
- aws_dynamodb_partiql:
query: "INSERT INTO footable VALUE {'foo':'?','bar':'?','baz':'?'}"
args_mapping: |
root = [
{ "S": this.foo },
{ "S": meta("kafka_topic") },
{ "S": this.document.content },
]

Query a GSI for a single record

The following example looks up a single record from the table footable using the global secondary index index_name, matching on the field bar. BatchExecuteStatement can't query a GSI, so use_batch is disabled:

pipeline:
processors:
- aws_dynamodb_partiql:
query: "SELECT * FROM \"footable\".\"index_name\" WHERE bar = ?"
use_batch: false
args_mapping: |
root = [
{ "S": this.bar },
]

Fields

query

A PartiQL query to execute for each message.

Type: string

unsafe_dynamic_query

Whether to enable dynamic queries that support interpolation functions.

Type: bool
Default: false

use_batch

Whether to execute all messages in a batch as a single BatchExecuteStatement call. Set this to false to execute one ExecuteStatement call per message instead, which is required for PartiQL SELECT queries against a global secondary index (GSI) — BatchExecuteStatement does not support querying a GSI. Only the first result row is used when a query returns multiple items.

Type: bool
Default: true

args_mapping

A Bloblang mapping that, for each message, creates a list of arguments to use with the query.

Type: string
Default: ""

region

The AWS region to target.

Type: string

endpoint

Allows you to specify a custom endpoint for the AWS API.

Type: string

tcp

TCP socket configuration.

Type: object

tcp.connect_timeout

Maximum amount of time a dial will wait for a connect to complete. Zero disables.

Type: string
Default: "0s"

tcp.keep_alive

TCP keep-alive probe configuration.

Type: object

tcp.keep_alive.idle

Duration the connection must be idle before sending the first keep-alive probe. Zero defaults to 15s. Negative values disable keep-alive probes.

Type: string
Default: "15s"

tcp.keep_alive.interval

Duration between keep-alive probes. Zero defaults to 15s.

Type: string
Default: "15s"

tcp.keep_alive.count

Maximum unanswered keep-alive probes before dropping the connection. Zero defaults to 9.

Type: int
Default: 9

tcp.tcp_user_timeout

Maximum time to wait for acknowledgment of transmitted data before killing the connection. Linux-only (kernel 2.6.37+), ignored on other platforms. When enabled, keep_alive.idle must be greater than this value per RFC 5482. Zero disables.

Type: string
Default: "0s"

credentials

Optional manual configuration of AWS credentials to use. More information can be found in this document.

Type: object

credentials.profile

A profile from ~/.aws/credentials to use.

Type: string

credentials.id

The ID of credentials to use.

Type: string

credentials.secret

The secret for the credentials being used.

Secret

This field contains sensitive information. Use a secret reference rather than a literal value.

Type: string

credentials.token

The token for the credentials being used, required when using short term credentials.

Secret

This field contains sensitive information. Use a secret reference rather than a literal value.

Type: string

credentials.from_ec2_role

Use the credentials of a host EC2 machine configured to assume an IAM role associated with the instance.

Type: bool

credentials.role

A role ARN to assume.

Type: string

credentials.role_external_id

An external ID to provide when assuming a role.

Type: string