AWS Kinesis Streams logs
AWS Kinesis Streams (logs)
AWS Kinesis Streams is used when you want near-real-time log ingestion with ordered, partitioned delivery. Worker batches log events, encodes them, and pushes them as Kinesis records. Key design levers are:
- Partitioning (how logs are distributed across shards)
- Batch size vs latency (Kinesis has strict per-record and per-request limits)
- Backpressure & buffering (what happens when Kinesis throttles)
- Concurrency & retry behavior (how aggressively worker pushes)
This sink is typically used as a fan-out ingestion layer before Lambda, Kinesis Analytics, Firehose, or custom consumers.
Core destination & routing
stream_name
The target Kinesis data stream. This is the logical pipe logs are written into. Shard count, retention, and scaling behavior are defined on the stream side, not here.
region
AWS region where the stream lives. Must match the region of the Kinesis stream, otherwise requests will be signed incorrectly.
endpoint
Overrides the AWS Kinesis endpoint. Used mainly for:
- AWS-compatible environments
- Local testing / emulators
- Private endpoints or special partitions
Partitioning & ordering (very important for Kinesis)
partition_key_field
Defines which log field is used as the Kinesis partition key.
Why this matters:
- Records with the same partition key always go to the same shard
- Ordering is preserved per partition key
- Hot partition keys can overload a shard
If not set:
- Worker generates a random/unique key
- Distribution is even, but ordering is lost
Typical strategies:
- Use tenant ID, service name, or host ID
- Avoid high-cardinality + bursty fields that create shard hot-spots
Batching (latency vs throughput trade-off)
batch.max_bytes
Maximum size of a batch before it is sent. This protects you from exceeding Kinesis request limits.
Larger values:
- Better throughput
- Higher latency
- Higher risk of throttling
Smaller values:
- Lower latency
- More API calls
batch.max_events
Maximum number of log events per batch. Useful when events are small but numerous.
batch.timeout_secs
Maximum time worker waits before flushing a batch.
Key behavior:
- Kinesis favors small, frequent batches
- Default is very low (1s) to keep latency tight
Buffering (what happens when Kinesis slows down)
buffer.type
Where events are queued before delivery:
- memory: fast, volatile
- disk: durable, survives restarts
Kinesis throttling is common → disk buffer is safer for production pipelines.
buffer.max_size
Hard cap on buffer storage. Once reached, worker must either block or drop events.
buffer.max_events
Event-count cap for memory buffers. Prevents unbounded memory growth.
buffer.when_full
Behavior when the buffer is exhausted:
- block: apply backpressure upstream (preferred for reliability)
- drop_newest: keep running but lose data
Compression
compression
Controls whether records are compressed before sending.
Options like gzip or zstd:
- Reduce payload size
- Lower shard bandwidth usage
- Add CPU overhead
Default is none because:
- Kinesis record size limits are strict
- Many consumers expect raw payloads
Compression choice must match consumer expectations.
Encoding (what is actually sent into Kinesis)
encoding.codec
Defines the wire format of each record:
- json, text, avro, cef, protobuf, etc.
This directly impacts:
- Consumer parsing logic
- Schema evolution
- Interoperability with AWS services
encoding.* (codec-specific options)
Each codec has its own tuning knobs:
- CSV field order
- JSON formatting
- Avro schema
- CEF metadata
These do not affect delivery mechanics, only payload structure.
Authentication & identity
auth.*
Controls how worker authenticates to AWS:
- Static keys
- Assume-role
- Instance metadata (IMDS)
- Named credential profiles
Important operational notes:
- Assume-role adds STS calls (latency + retry impact)
- IMDS requires network access to the metadata service
- Region mismatches can break signing
Request behavior & throughput control
request.concurrency
How many requests worker sends in parallel.
- Higher concurrency = higher throughput
- Too high = throttling, retries, shard pressure
request.rate_limit_num / request.rate_limit_duration_secs
Explicit rate limiting to protect:
- Kinesis shard limits
- Downstream AWS account quotas
request.timeout_secs
How long worker waits for AWS to respond.
Too low:
- Causes retries and duplicates
Too high:
- Slows failure detection
Retries & partial failures
request.retry_attempts
Maximum retry count for failed requests.
request_retry_partial
Whether worker retries partially successful Kinesis responses.
Why this exists:
- Kinesis can accept some records and reject others in the same request
- Retrying partials improves reliability but increases duplicates risk
request.retry_* backoff controls
Shape how retries ramp:
- Initial delay
- Max delay
- Jitter to avoid retry storms
Proxy & TLS (network path)
proxy.*
Routes traffic through HTTP(S) proxies. Useful in restricted or enterprise networks.
tls.*
Controls certificate trust, hostname verification, and client auth.
Critical for:
- Private endpoints
- Man-in-the-middle proxies
- AWS-compatible non-AWS backends