Skip to main content
Version: 3.0 (next)

Amazon Timestream Amazon Timestream Integration Guide

Connect to Amazon Timestream for LiveAnalytics to write telemetry from a pipeline and query it back with Timestream SQL. This guide covers connection setup, function configuration, and pipeline integration.

Overview​

Timestream for LiveAnalytics is a serverless time-series database. Recent data sits in a fast memory store and ages automatically into cheaper magnetic storage, so you size retention rather than servers. The connector provides:

  • Writes — single-measure and multi-measure records with dimensions, batched to the 100-record WriteRecords quota
  • Queries — Timestream SQL, which adds interpolation, binning and per-series latest-value functions on top of standard SQL
  • Discovery — list databases, list tables, and read a table's retention windows
  • Flexible authentication — IAM role, instance profile or IRSA through the AWS SDK default credential chain, or static access keys
  • Endpoint override — point both APIs at a VPC endpoint or a local emulator
Two Timestream products

This connector targets Timestream for LiveAnalytics (the original Amazon Timestream), which speaks the WriteRecords and Query APIs. It does not target Timestream for InfluxDB; for that, use the InfluxDB connector.

Connection Configuration​

Creating an Amazon Timestream Connection​

Navigate to Connections → New Connection → Amazon Timestream and configure the following:

1. Profile Information​

FieldDefaultDescription
Profile Name-A descriptive name for this connection profile (required, max 100 characters)
Description-Optional description for this Timestream connection

2. Connection​

FieldDefaultDescription
AWS Regionus-east-1The AWS region hosting the Timestream database, e.g. eu-central-1 (required)
Default Database-Database used by functions that do not name one. Optional if every function sets its own
Default Table-Table used by functions that do not name one. Optional if every function sets its own

3. Authentication​

FieldDefaultDescription
Access Key ID-AWS Access Key ID. Masked on edit. Leave empty to use the AWS SDK default credential chain (env vars, shared config, IAM role, IRSA)
Secret Access Key-AWS Secret Access Key. Masked on edit; leave empty to use the default credential chain
Session Token-Session token for temporary STS credentials (optional). Masked on edit
Prefer IAM roles over static keys

On EC2/ECS/EKS, leave the key fields empty and attach an IAM role (or IRSA) — the connector picks up credentials from the AWS SDK default chain, so no secrets are stored in MaestroHub. Use static access keys only when a role is not available.

The identity needs timestream:DescribeEndpoints and timestream:ListDatabases (Test Connection uses the latter), plus timestream:WriteRecords for writes and timestream:Select for queries. timestream:ListTables and timestream:DescribeTable cover the discovery functions.

Access Key ID and Secret Access Key must be filled in together. Filling only one is rejected at save time, because the connection would silently fall back to the default credential chain and run as an identity you did not choose.

4. Advanced​

FieldDefaultDescription
Request Timeout1mHow long a single Timestream API call may take before it is abandoned (1s–1h)
Max Result Rows1000Maximum rows a query returns before truncating (1–100000)
Custom Endpoint-Endpoint URL for both the ingest and query APIs. Leave empty for AWS
Custom Endpoint turns off endpoint discovery

Timestream normally asks DescribeEndpoints which cell to talk to and routes each call there. Setting a Custom Endpoint disables that discovery so every request goes to the URL you gave — which is exactly what a VPC endpoint or a local emulator needs, and exactly wrong against a plain AWS account. Leave the field empty unless you are pointing at a specific endpoint on purpose.

Notes
  • Required fields: Profile Name and AWS Region. Everything else has a default or falls back to the credential chain.
  • Database and table: a write or describe needs both. Set them on the connection, on the function, or split them — the function value wins.
  • Row cap: result sets larger than Max Result Rows are truncated, and the call reports _metadata.truncated: true.
  • Security: the access key, secret key and session token are encrypted at rest and masked on edit. Leave a secret empty to keep the stored value.

Function Builder​

Creating Timestream Functions​

Once a connection exists, create reusable write, query and discovery functions:

  1. Open the connection and go to its Functions tab → New Function
  2. Choose a Timestream function type
  3. Configure the function parameters
Timestream function type selection

Choose from Write Records, Query, List Databases, List Tables and Describe Table function types

Write Records Function​

Purpose: write one or more records into a Timestream table. Each record carries its dimensions (the metadata that identifies the series), one or more measures, and a timestamp.

Configuration Fields

FieldTypeRequiredDefaultDescription
RecordsJSONYes-A record object or an array of them. Supports ((param)) templating — ((records)) takes the whole upstream batch
DatabaseStringNoconnection defaultOverride the connection's default database
TableStringNoconnection defaultOverride the connection's default table
Multi-measure NameStringNometricsMeasure name given to multi-measure records that do not set their own
Time UnitEnumNoMILLISECONDSUnit for numeric timestamps: MILLISECONDS, SECONDS, MICROSECONDS, NANOSECONDS
Common DimensionsJSONNo-Dimensions merged into every record in the batch
TimeoutDurationNo30mBound on this single operation

Record shape

A record is single-measure or multi-measure. Single-measure writes one value per row:

{
"dimensions": { "machine_id": "press-02", "site": "ankara" },
"measureName": "temperature",
"measureValue": 42.5,
"time": "2026-03-01T10:00:00Z"
}

Multi-measure writes several values in one row, which is cheaper to store and faster to query:

{
"dimensions": { "machine_id": "press-02" },
"measures": { "temperature": 42.5, "rpm": 1200, "state": "running" },
"time": "2026-03-01T10:00:00Z"
}

Rules the connector applies:

  • time is optional. An RFC 3339 string or a real timestamp is exact and sent in milliseconds; a bare number is interpreted in the op's Time Unit. A record with no time is stamped with the moment the write was produced — the same on every store-and-forward replay.
  • Measure types are inferred. A JSON number becomes DOUBLE, a string VARCHAR, a boolean BOOLEAN. Whole numbers stay DOUBLE on purpose: a column's type is fixed by its first write, so a series that starts at 0 and later reads 0.5 would be rejected forever if the first row had chosen BIGINT.
  • Dimensions are strings. Numbers and booleans are rendered; a null dimension is dropped rather than written as text.
  • snake_case also works. measure_name, measure_value, tags, fields and timestamp are read as aliases, so a row from a SQL or MQTT node usually needs no reshaping.
  • Batching is automatic. WriteRecords accepts at most 100 records per call, so a larger payload is split. _metadata.batches reports how many calls went out.

Example Configuration

{
"records": "((records))",
"database": "ot_telemetry",
"table": "readings",
"commonDimensions": { "site": "ankara" }
}

Response Format

{
"recordsWritten": 120,
"memoryStore": 118,
"magneticStore": 2,
"batches": 2
}
Rejected records are permanent

Timestream refuses a record whose timestamp falls outside the table's memory-store or magnetic-store retention, whose type disagrees with an established column, or that duplicates an existing record at the same version. Retrying sends the identical rows and gets the identical refusal, so the connector classifies RejectedRecordsException as permanent and the batch goes to the dead-letter queue rather than burning the retry budget. The error names the offending record index and reason; Describe Table shows the retention windows.

If a later batch fails after earlier ones succeeded, the earlier records are already durable — there is no rollback. The error says how far the write got (batch 2 of 3, 100 already written) and _metadata.recordsWritten carries the same number, so a replay can skip what landed.

Query Function​

Purpose: run a Timestream SQL statement and return the result rows as structured records.

Configuration Fields

FieldTypeRequiredDefaultDescription
SQL QueryStringYes-Timestream SQL statement. Supports ((param)) templating
Max RowsNumberNoconnection defaultMaximum rows to return (1–100000)
TimeoutDurationNo30mBound on this single operation

Table names are quoted as "database"."table".

Example Configuration

{
"sql": "SELECT BIN(time, 5m) AS t, AVG(measure_value::double) AS avg_temp FROM \"ot_telemetry\".\"readings\" WHERE machine_id = '((machine))' AND time > ago(1h) GROUP BY BIN(time, 5m) ORDER BY t",
"maxRows": 500
}

Response Format

{
"rows": [
{ "t": "2026-03-01 10:00:00.000000000", "avg_temp": "42.7" }
],
"columns": [
{ "name": "t", "type": "TIMESTAMP" },
{ "name": "avg_temp", "type": "DOUBLE" }
],
"rowCount": 1
}

Timestream returns scalars as strings with no type tag on the value itself, so the rows carry strings and columns carries the type. Arrays, rows and time-series columns unwrap into nested JSON.

List Databases Function​

Purpose: list the Timestream databases visible to the connection's credentials.

Takes no configuration beyond an optional timeout.

Response Format

{
"databases": [
{ "name": "ot_telemetry", "tableCount": 4 }
],
"count": 1
}

List Tables Function​

Purpose: list the tables in a Timestream database along with their status.

FieldTypeRequiredDefaultDescription
DatabaseStringNoconnection defaultDatabase whose tables to list. Supports ((param)) templating
TimeoutDurationNo30mBound on this single operation

Response Format

{
"tables": [
{ "name": "readings", "status": "ACTIVE" }
],
"count": 1
}

Describe Table Function​

Purpose: read one table's status and retention settings. This is the first thing to check when writes are being rejected.

FieldTypeRequiredDefaultDescription
DatabaseStringNoconnection defaultDatabase holding the table. Supports ((param)) templating
TableStringNoconnection defaultTable to describe. Supports ((param)) templating
TimeoutDurationNo30mBound on this single operation

Response Format

{
"name": "readings",
"status": "ACTIVE",
"arn": "arn:aws:timestream:eu-central-1:123456789012:database/ot_telemetry/table/readings",
"memoryStoreRetentionHours": 12,
"magneticStoreRetentionDays": 365
}

A write whose timestamp is older than memoryStoreRetentionHours (and not covered by magnetic-store writes) is rejected. A write further back than magneticStoreRetentionDays is rejected outright.

Using Parameters​

The ((parameterName)) syntax turns a function into a dynamic, reusable building block. Parameters are auto-detected from the SQL, the records payload and the other templated fields, and can be configured with:

ConfigurationDescriptionExample
TypeData type validationstring, number, boolean, datetime, json, buffer
RequiredMake the parameter mandatory or optionalRequired / Optional
Default ValueFallback value if not providedreadings, 2026-03-01, 1000
DescriptionHelp text for users"Machine identifier", "Target table"
Timestream parameter configuration

Parameters detected from the records payload and other templated fields are configured with type, requiredness, and defaults

Parameter availability

Write Records, Query, List Tables and Describe Table accept ((parameter)) templating. List Databases has no templatable field and therefore takes no parameters.

Pipeline Integration​

Use the Timestream functions you create here as nodes inside the Pipeline Designer. Drop a write or query node onto the canvas, bind its parameters to upstream outputs or constants, and configure error handling as needed.

Common patterns:

  • Trigger → Transform → Write: take readings off MQTT or OPC UA, reshape them into records, and persist them with a Write Records node
  • Schedule → Query → Act: recompute a KPI on a schedule and push the series into a dashboard or a notification
  • Describe → Branch: read the retention window before a backfill, and route rows that are too old somewhere else

For broader orchestration patterns, see the Connector Nodes page and the Amazon Timestream node reference.

Amazon Timestream write node in the pipeline designer

Timestream write node with connection, function, and parameter bindings

Common Use Cases​

Persisting line telemetry​

Scenario: readings arrive on MQTT and need to land in Timestream for dashboards.

Write Records Configuration:

{
"records": "((records))",
"commonDimensions": { "site": "ankara", "line": "A" }
}

Pipeline Integration: MQTT trigger → a transform that maps each message to {dimensions, measures, time} → Write Records. Records batch automatically at 100 per call.


Downsampling for a dashboard​

Scenario: a dashboard needs five-minute averages, not raw samples.

Query Configuration:

{
"sql": "SELECT machine_id, BIN(time, 5m) AS t, AVG(measure_value::double) AS avg_temp FROM \"ot_telemetry\".\"readings\" WHERE time > ago(6h) GROUP BY machine_id, BIN(time, 5m) ORDER BY t"
}

Diagnosing rejected writes​

Scenario: a backfill is bouncing and nobody knows why.

Run Describe Table, read memoryStoreRetentionHours and magneticStoreRetentionDays, and compare them against the oldest timestamp in the batch. Records older than both windows can never be accepted; widen the retention or drop them upstream.

Troubleshooting​

SymptomLikely causeWhat to do
Test Connection fails with AccessDeniedExceptionThe identity lacks timestream:ListDatabasesGrant it, or attach a role that has it
Test Connection hangs, then times outCustom Endpoint points somewhere unreachableClear the field for AWS, or fix the URL
The table does not exist on writeDatabase or table resolved to the wrong valueCheck both the connection defaults and the function overrides — the function wins
RejectedRecordsException naming a timestampThe record is outside the retention windowRun Describe Table and compare; widen retention or filter upstream
RejectedRecordsException naming a typeThe column's type was fixed by an earlier writeQuery the column, then either match the type or write to a new measure name
Query returns fewer rows than expectedThe result hit Max Result RowsRaise the cap, or check _metadata.truncated