Connect

aws_dynamodb_partiql

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

Introduced in version 3.48.0.

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.

  • Common

  • Advanced

processor:
  label: ""
  aws_dynamodb_partiql:
    query: "" # No default (required)
    args_mapping: ""
processor:
  label: ""
  aws_dynamodb_partiql:
    query: "" # No default (required)
    unsafe_dynamic_query: false
    use_batch: true
    args_mapping: ""
    region: "" # No default (optional)
    endpoint: "" # No default (optional)
    tcp:
      connect_timeout: 0s
      keep_alive:
        idle: 15s
        interval: 15s
        count: 9
      tcp_user_timeout: 0s
    credentials:
      profile: "" # No default (optional)
      id: "" # No default (optional)
      secret: "" # No default (optional)
      token: "" # No default (optional)
      from_ec2_role: false # No default (optional)
      role: "" # No default (optional)
      role_external_id: "" # No default (optional)

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

args_mapping

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

Type: string

Default: ""

credentials

Manually configure the AWS credentials to use (optional). For more information, see the Amazon Web Services guide.

Type: object

credentials.from_ec2_role

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

Requires version 4.2.0 or later.

Type: bool

credentials.id

The ID of the AWS credentials to use.

Type: string

credentials.profile

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

Type: string

credentials.role

The ARN of the role to assume.

Type: string

credentials.role_external_id

An external ID to use when assuming a role.

Type: string

credentials.secret

The secret for the AWS credentials in use.

This field contains sensitive information that usually shouldn’t be added to a configuration directly. For more information, see Secrets.

Type: string

credentials.token

The token for the AWS credentials in use. Required only when using short-term credentials.

Type: string

endpoint

A custom endpoint URL for AWS API requests. Use this to connect to AWS-compatible services or local testing environments instead of the standard AWS endpoints.

Type: string

query

A PartiQL query to execute for each message.

Type: string

region

The AWS region in which your resources are hosted.

Type: string

tcp

Configure TCP socket-level settings to optimize network performance and reliability. These low-level controls are useful for:

  • Unresponsive hosts: Set connect_timeout to limit how long a connection attempt can take (the default 0s sets no limit)

  • Long-lived connections: Configure keep_alive settings to detect and recover from stale connections

  • Unstable networks: Tune keep-alive probes to balance between quick failure detection and avoiding false positives

  • Linux systems with specific requirements: Use tcp_user_timeout (Linux 2.6.37+) to control data acknowledgment timeouts

Most users should keep the default values. Only modify these settings if you’re experiencing connection stability issues or have specific network requirements.

Requires version 4.69.0 or later.

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.count

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

Type: int

Default: 9

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.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

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), because BatchExecuteStatement does not support querying a GSI. Only the first result row is used when a query returns multiple items.

Requires version 4.102.0 or later.

Type: bool

Default: true