Cloud

aws_dynamodb_partiql

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

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.

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 Manage Secrets before adding it to your configuration.

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.

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.

Type: bool

Default: true