AWS Kinesis Data Firehose logs
AWS Kinesis Data Firehose (logs) — quick summary
Firehose is the “managed delivery” version of Kinesis: Kron TP sends logs to a Firehose delivery stream, and Firehose handles buffering + delivery into a destination (S3 / OpenSearch / Splunk / HTTP endpoint / etc.). So the main knobs on the worker side are about how the worker batches/encodes/sends, not how data lands in S3 (that’s configured in Firehose itself).
Big takeaways:
- You still tune batch / buffer / request for performance and reliability.
- partition_key_field exists here too (distribution behavior is handled by Firehose internally, but partition key still affects record grouping behavior on the ingest side).
- request_retry_partial is the big “duplication risk” toggle.
Core destination parameters
stream_name (required)
The Firehose delivery stream name (not a Kinesis Data Stream). This must match exactly what’s created in AWS Firehose.
region (optional)
AWS region where the Firehose stream exists. If omitted, AWS SDK resolution + environment config may decide it, but explicit is safer.
endpoint (optional)
Override the Firehose API endpoint. Typical use cases:
- AWS-compatible endpoints
- Testing / private environments
Partitioning / keying
partition_key_field (optional)
Chooses the log field used as the “partition key” value for records.
Practical meaning:
- If set, records that share the same key are logically grouped on the producer side.
- If not set, worker generates a unique key per record, which tends to spread records more evenly.
When to set it
- Multi-tenant pipelines (tenant_id / customer_id)
- Service grouping (service.name / namespace)
- Host grouping (host / node)
When not to
- If you have a risk of “hot keys” (one tenant/service dominates volume), because that can amplify throttling/partial failures.
Batching (throughput vs latency)
Firehose ingest calls have their own constraints, so batching is your “API efficiency” lever.
batch.max_bytes (optional, default ~4 MiB)
Max uncompressed batch payload size before flush.
- Larger = fewer requests, better throughput
- Too large = higher chance of request rejection / partial failures
batch.max_events (optional, default 500)
Flush after N events even if bytes limit is not reached.
- Useful when logs are small but frequent
batch.timeout_secs (optional, default 1s)
Time-based flush to keep latency bounded.
- Lower = faster delivery, more API calls
- Higher = fewer calls, more delay and larger bursts
Buffering (backpressure & durability)
This controls what happens inside worker when Firehose is slow, throttled, or unreachable.
buffer.type (optional, default memory)
- memory: fast, loses data on crash/restart
- disk: durable, survives restarts (safer for production pipelines)
buffer.max_size (required)
Hard cap of how much worker can buffer.
This is the safety valve:
- Too small → drops or backpressure happens frequently
- Too large → disk/memory pressure (depends on buffer.type)
buffer.max_events (optional, default 500 for memory)
Extra cap relevant mainly for memory buffering.
buffer.when_full (optional, default block)
- block: apply backpressure upstream (reliable)
- drop_newest: keep running, lose newest logs (performance-first)
Compression
compression (optional, default none)
Compresses outgoing payloads.
- Compression reduces network + payload size
- Adds CPU overhead
- Can complicate consumers if you later route Firehose to places expecting plaintext (depends on Firehose destination handling)
Encoding (payload format)
encoding.codec (required)
Defines the serialized record format sent to Firehose:
- json, text, avro, cef, protobuf, etc.
This is a contract with whatever reads downstream (or what Firehose delivers to).
encoding.* (codec-specific)
All the schema/format knobs (CSV field order, Avro schema, JSON options, etc.). These affect structure, not delivery reliability.
Authentication
auth.*
How worker authenticates to AWS:
- static keys
- assume role
- IMDS (instance metadata)
- credentials profile
Operationally:
- assume_role adds STS dependency
- IMDS requires network reachability + correct instance role permissions
Proxy & TLS (network path)
proxy.*
Send AWS API traffic via HTTP(S) proxy. Useful in enterprise egress-controlled networks.
tls.*
TLS verification, custom CAs, SNI, etc. Relevant when:
- corporate MITM proxies
- custom endpoints / private endpoints
Request / retry behavior (most important for “at-least-once”)
request.concurrency (default adaptive)
Controls how many concurrent requests worker tries to keep in flight.
- More concurrency = higher throughput
- Too much = throttling → retries → duplicates risk
request.rate_limit_num / request.rate_limit_duration_secs
Hard caps request rate to avoid:
- hitting account limits
- causing repeated throttles
request.timeout_secs
Timeout for each request. Too low → unnecessary retries/duplicates. Too high → slow failure detection.
request.retry_attempts, request.retry_* backoff/jitter
Controls retry persistence + retry storm behavior.
request_retry_partial (default false) ⚠️
If true, worker will retry requests that partially succeeded (some records accepted, some rejected).
- Pros: higher chance everything eventually gets delivered
- Cons: duplicates become more likely (because accepted records may be resent)
If your downstream cannot deduplicate, keep this false unless you have a strong reason.