Skip to main content
Version: 3.0 (next)

Amazon Kinesis Amazon Kinesis Integration Guide

Connect to Amazon Kinesis Data Streams to write and read streaming records, and to consume records into event-driven pipelines. This guide covers connection setup, function configuration, and pipeline integration.

Overview​

Amazon Kinesis Data Streams is a serverless, high-throughput streaming service — an AWS-native alternative to Kafka for telemetry and analytics pipelines. The connector provides:

  • Put records — single (PutRecord) or batched up to 500 (PutRecords) with a configurable partition-key strategy
  • Read records on demand with LATEST / TRIM_HORIZON / AT_TIMESTAMP shard iterators
  • Stream consumption — a long-running consumer that fires a pipeline once per record via the Kinesis Trigger
  • Stream discovery via ListStreams
  • Automatic shard discovery and re-discovery across resharding
  • Flexible authentication — IAM access keys or the AWS SDK default credential chain (env vars, IRSA, instance profile)
  • Custom endpoint for LocalStack and AWS-compatible services
Streaming semantics

Kinesis orders records within a shard by partition key. Writes are not idempotent — a replayed PutRecord creates a duplicate — so attach a stable partition key when ordering matters. For continuous consumption use the trigger; for ad-hoc reads use the Get Records function.

Connection Configuration​

Creating an Amazon Kinesis Connection​

Navigate to Connections → New Connection → Amazon Kinesis 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 Kinesis connection

2. Connection​

FieldDefaultDescription
AWS Regionus-east-1AWS region where the Kinesis streams live, e.g. us-east-1, eu-central-1 (required)

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. Required when Access Key ID is set; leave empty to use the default credential chain. Masked on edit
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 automatically, so no secrets are stored in MaestroHub. The identity needs kinesis:PutRecord*, kinesis:GetRecords, kinesis:GetShardIterator, kinesis:DescribeStream*, and kinesis:ListStreams as appropriate.

4. Endpoint​

FieldDefaultDescription
Custom Endpoint-Custom Kinesis endpoint URL for LocalStack or other AWS-compatible services (e.g. http://localhost:4566). Leave empty for AWS Kinesis

5. Advanced​

FieldDefaultDescription
Request Timeout (seconds)30Upper bound on every Kinesis API call on this connection, including the connect check (1–3600). A function's Timeout Override can shorten it, not extend it.
Max Retries3Retries after the first attempt for transient failures such as ProvisionedThroughputExceeded (0 = no retries, up to 10)
Notes
  • Required fields: Profile Name and AWS Region. Credentials are optional when the default credential chain is available.
  • Throughput: Each shard sustains 1 MB/s or 1000 records/s for writes and 2 MB/s / 5 GetRecords calls per second for reads — batch with Put Records and tune the poll interval accordingly.
  • Security: Access keys and the session token are encrypted and stored securely, masked on edit. Leave a secret empty to keep the stored value.

Function Builder​

Creating Kinesis Functions​

Once a connection exists, create reusable record and stream functions:

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

Choose from Put Record, Put Records (Batch), Get Records, List Streams, and Consume (Trigger) function types

Put Record Function​

Purpose: Write one record to a Kinesis Data Stream via PutRecord. The partition key decides which shard the record lands on; records with the same key go to the same shard in order.

Configuration Fields

FieldTypeRequiredDefaultDescription
Stream NameStringYes-Name of the Kinesis Data Stream to write to. Supports ((param)) templating
DataStringYes-The record payload (up to 1 MB). Strings are sent as-is; objects and arrays are JSON-encoded automatically, so passing ((rowObject)) works without pre-stringifying
Partition KeyStringNorandomDetermines the target shard. Records with the same key are ordered within a shard. Leave empty to generate a random key
Explicit Hash KeyStringNo-Advanced: override the partition-key hash to target a specific shard by its hash-key range
Timeout Override (seconds)NumberNoconnection defaultBound on this operation (0–3600; 0 or empty uses only the connection's Request Timeout). The connection's Request Timeout still applies and the shorter one wins.

Example Configuration

{
"streamName": "telemetry",
"data": "((reading))",
"partitionKey": "((deviceId))"
}

Response Format

{
"shardId": "shardId-000000000000",
"sequenceNumber": "49590338...66",
"partitionKey": "press-02"
}

Use Cases:

  • Forward processed telemetry to a Kinesis stream for S3 + Athena
  • Emit a per-device event keyed by device ID for ordered processing
  • Bridge MQTT messages into a Kinesis analytics stream

Put Records (Batch) Function​

Purpose: Write a batch of records via PutRecords (up to 500 records / 5 MB per call). Each entry carries its own partition key. PutRecords can partially fail — the result reports the failed count and per-record errors so the caller can retry only the failures.

Configuration Fields

FieldTypeRequiredDefaultDescription
Stream NameStringYes-Name of the Kinesis Data Stream to write to. Supports ((param)) templating
Records (JSON array)StringYes-JSON array of records. Raw payloads ([{"msg":"hello"}, {"msg":"world"}] — each object becomes one record) or envelopes ([{"data":"...","partitionKey":"..."}]) when you need to control the partition key. Max 500 entries / 5 MB total
Timeout Override (seconds)NumberNoconnection defaultBound on this operation. The connection's Request Timeout still applies and the shorter one wins.

Response Format

{
"records": [
{ "sequenceNumber": "49590338...66", "shardId": "shardId-000000000000" },
{ "errorCode": "ProvisionedThroughputExceededException", "errorMessage": "Rate exceeded" }
],
"failedRecordCount": 1,
"successCount": 1
}

Use Cases:

  • Flush a buffered window of sensor readings in one call
  • Bulk-load enriched events into an analytics stream

Get Records Function​

Purpose: Read a batch of records on demand via GetShardIterator + GetRecords across the stream's shards. Choose where to start: LATEST (only new records), TRIM_HORIZON (oldest retained), or AT_TIMESTAMP (from a point in time). For continuous consumption use the Kinesis Trigger instead.

Configuration Fields

FieldTypeRequiredDefaultDescription
Stream NameStringYes-Name of the Kinesis Data Stream to read from. Supports ((param)) templating
Start PositionEnumNoLATESTLATEST reads only records added after the call; TRIM_HORIZON reads from the oldest retained record; AT_TIMESTAMP reads from the given timestamp
Start TimestampStringNo-RFC3339 timestamp (e.g. 2026-01-01T00:00:00Z). Required when Start Position is AT_TIMESTAMP
Shard IDStringNoall shardsRead only this shard (e.g. shardId-000000000000). Leave empty to read across all shards
Max RecordsNumberNo100Maximum total records to return (1–10000). Applied as the per-shard GetRecords limit so the read is bounded at the source

Response Format

{
"records": [
{
"shardId": "shardId-000000000000",
"data": "{\"temperature\":42.7}",
"partitionKey": "press-02",
"sequenceNumber": "49590338...66"
}
],
"count": 1
}

Use Cases:

  • Inspect the most recent records on a stream for debugging
  • Pull a controlled batch from the start of the retention window

List Streams Function​

Purpose: List the Kinesis Data Streams visible to the configured credentials. Use this to discover available streams before writing, or to audit stream inventory for a region.

Configuration Fields

FieldTypeRequiredDefaultDescription
Max ItemsNumberNo100Maximum number of stream names to return (1–10000). The connector paginates AWS responses internally and stops once this budget is reached

Response Format

{
"streams": ["telemetry", "processed-events", "audit-log"],
"count": 3
}

Use Cases:

  • Discover available streams before writing
  • Audit stream inventory for a region

Consume (Trigger) Function​

Purpose: A long-running consumer that polls every shard of a stream and fires the pipeline once per record. Used as the first node in a pipeline via the Kinesis Trigger — it has no upstream input and takes no dynamic parameters.

Configuration Fields

FieldTypeRequiredDefaultDescription
Stream NameStringYes-Name of the Kinesis Data Stream to consume
Start PositionEnumNoLATESTWhere the consumer begins on first start: LATEST (only new records), TRIM_HORIZON (oldest retained), or AT_TIMESTAMP
Start TimestampStringNo-RFC3339 timestamp. Required when Start Position is AT_TIMESTAMP
Records Per PollNumberNo100Maximum records fetched per GetRecords call per shard (1–10000)
Poll Interval (ms)NumberNo1000Delay between polling rounds when a shard is caught up (200–60000). Lower = lower latency but more API calls (mind the 5 GetRecords/s per-shard limit)
Single-instance consumer

The consumer is single-instance: durable cross-restart checkpointing and enhanced fan-out are not yet supported. Without durable checkpointing, restarting the consumer re-applies the configured Start Position.

Use Cases:

  • Trigger a pipeline on each telemetry record landing on a stream
  • Fan Kinesis events into UNS topics in real time

Using Parameters​

The ((parameterName)) syntax turns a function into a dynamic, reusable building block. Parameters are auto-detected from the stream name, payload, and 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 providedtelemetry, press-02
DescriptionHelp text for users"Target stream", "Device partition key"
Kinesis Parameter Configuration

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

Parameter Availability

The Put Record, Put Records, and Get Records functions accept ((parameter)) templating in the stream name, payload, partition key, and other templated fields. The Consume trigger function takes no dynamic parameters — it is configured once and runs as a pipeline trigger.

Pipeline Integration​

Use the Kinesis functions you create here as nodes inside the Pipeline Designer. Drag a Put, Get Records, or List Streams node onto the canvas and bind its parameters to upstream outputs, or start a pipeline from the Kinesis Trigger to process records as they arrive.

Common patterns include:

  • Collect → Put: Gather values from OPC UA, MQTT, or Modbus and stream them into Kinesis with Put Record / Put Records
  • Consume → Transform → Act: Start a pipeline from the Kinesis Trigger, transform each record, and write it onward
  • Discover → Read: List streams, then pull a controlled batch from one with Get Records

For broader orchestration patterns that combine Kinesis with other connector steps, see the Connector Nodes page, the Kinesis node reference, and the Kinesis Trigger.

Kinesis Put Record node in the pipeline designer

Kinesis Put Record node with connection, function, and parameter bindings

Common Use Cases​

Bridging OT Telemetry into an Analytics Stream​

Scenario: Forward temperature, pressure, and vibration readings from plant equipment into a Kinesis stream that feeds S3 + Athena for analytics.

Put Record Configuration:

{
"streamName": "equipment-telemetry",
"data": "((reading))",
"partitionKey": "((machineId))"
}

Pipeline Integration: Connect after OPC UA or Modbus read nodes to continuously stream equipment telemetry keyed by machine.


Event-Driven Stream Processing​

Scenario: Process each record landing on a Kinesis stream in real time.

Start the pipeline with the Kinesis Trigger on the source stream, transform each record, and route the result to UNS topics, a database, or another stream.


Bulk Ingestion with Partial-Failure Handling​

Scenario: Flush a buffered window of readings in one call and retry only the records that failed.

Use Put Records (Batch) and branch on failedRecordCount — when non-zero, re-submit the entries whose result carries an errorCode.