Building a Real-Time Analytics Pipeline with AWS Kinesis: From Ingestion to Dashboard
How we built a real-time analytics pipeline processing 2.4 million events per minute using Kinesis Data Streams, Firehose, and Lambda with sub-second latency.

Our product analytics team was stuck in a 15-minute feedback loop. Events hit our API, queued in SQS, landed in S3 via batch jobs, got picked up by Spark ETL, and finally appeared in Redshift. By the time a dashboard updated, the moment had passed. When our growth team ran a flash sale, they could not see conversion impact until the promotion was already over.
We rebuilt the pipeline around Kinesis Data Streams and cut end-to-end latency from 15 minutes to 1.2 seconds at P99. Here is the architecture, the performance data, and the sharp edges we hit along the way.
The Problem: Batch Analytics Cannot Drive Real-Time Decisions
Our existing pipeline processed events in 5-minute micro-batches. Acceptable for historical reporting, useless for:
- Real-time A/B test monitoring (detecting regressions before they affect thousands of users)
- Live conversion funnel tracking during promotions
- Fraud detection requiring sub-10-second response times
- Operational dashboards showing current system health
The business requirement was clear: events must be queryable within 2 seconds of occurrence, at a sustained throughput of 2 million events per minute with burst capacity to 5 million.
Architecture: Three-Layer Streaming Design
The final architecture separates concerns into ingestion, processing, and serving layers:
Ingestion Layer
- Kinesis Data Streams with 64 shards (initial provisioning)
- Producer SDK with PutRecords batching and retry logic
- Partition key strategy based on tenant_id for even distribution
Processing Layer
- Lambda consumers with Enhanced Fan-Out for dedicated throughput
- Kinesis Data Analytics (Apache Flink) for windowed aggregations
- Dead letter queue in SQS for failed processing attempts
Serving Layer
- Kinesis Firehose to S3 (Parquet format) for cold storage
- Direct Lambda-to-DynamoDB writes for hot queries
- OpenSearch for full-text search and complex aggregations
Producer Configuration: Maximizing Throughput
The default Kinesis producer settings leave significant throughput on the table. Each shard supports 1MB/second or 1,000 records/second for writes. With 64 shards, our theoretical maximum is 64MB/s or 64,000 records/second. But record aggregation changes the math entirely:
import { KinesisClient, PutRecordsCommand } from '@aws-sdk/client-kinesis';
interface AnalyticsEvent {
eventId: string;
tenantId: string;
eventType: string;
timestamp: number;
payload: Record<string, unknown>;
}
class KinesisProducer {
private buffer: AnalyticsEvent[] = [];
private readonly maxBatchSize = 500;
private readonly maxBatchBytes = 4 * 1024 * 1024; // 4MB limit for PutRecords
private readonly flushIntervalMs = 200;
private readonly client: KinesisClient;
private readonly streamName: string;
constructor(streamName: string, region: string) {
this.client = new KinesisClient({ region });
this.streamName = streamName;
setInterval(() => this.flush(), this.flushIntervalMs);
}
async publish(event: AnalyticsEvent): Promise<void> {
this.buffer.push(event);
if (this.buffer.length >= this.maxBatchSize) {
await this.flush();
}
}
private async flush(): Promise<void> {
if (this.buffer.length === 0) return;
const batch = this.buffer.splice(0, this.maxBatchSize);
const records = batch.map((event) => ({
Data: Buffer.from(JSON.stringify(event)),
PartitionKey: event.tenantId,
}));
const command = new PutRecordsCommand({
StreamName: this.streamName,
Records: records,
});
const response = await this.client.send(command);
if (response.FailedRecordCount && response.FailedRecordCount > 0) {
const failedRecords = response.Records!
.map((r, i) => (r.ErrorCode ? batch[i] : null))
.filter(Boolean) as AnalyticsEvent[];
// Exponential backoff retry for failed records
await this.retryWithBackoff(failedRecords);
}
}
private async retryWithBackoff(
events: AnalyticsEvent[],
attempt = 0
): Promise<void> {
if (attempt >= 3) {
console.error(`Failed to publish ${events.length} events after 3 retries`);
// Send to DLQ
return;
}
await new Promise((r) => setTimeout(r, Math.pow(2, attempt) * 100));
for (const event of events) {
this.buffer.push(event);
}
}
}
The key decisions here: 200ms flush interval balances latency against throughput, partition key on tenant_id ensures related events land on the same shard for ordered processing, and the retry logic handles the inevitable ProvisionedThroughputExceededException during bursts.
Consumer Configuration: Enhanced Fan-Out for Dedicated Throughput
Standard consumers share the 2MB/second read capacity per shard across all consumers. With 4 consuming applications, each only gets 500KB/second per shard. Enhanced Fan-Out gives each consumer its own dedicated 2MB/second pipe:
import {
KinesisClient,
SubscribeToShardCommand,
} from '@aws-sdk/client-kinesis';
async function processWithEnhancedFanOut(
consumerArn: string,
shardId: string
): Promise<void> {
const client = new KinesisClient({ region: 'us-east-1' });
const command = new SubscribeToShardCommand({
ConsumerARN: consumerArn,
ShardId: shardId,
StartingPosition: { Type: 'LATEST' },
});
const response = await client.send(command);
if (response.EventStream) {
for await (const event of response.EventStream) {
if ('SubscribeToShardEvent' in event) {
const records = event.SubscribeToShardEvent!.Records || [];
const batchStartTime = Date.now();
for (const record of records) {
const payload = JSON.parse(
Buffer.from(record.Data!).toString('utf-8')
);
await processEvent(payload);
}
const processingTime = Date.now() - batchStartTime;
console.log(
`Processed ${records.length} records in ${processingTime}ms`
);
}
}
}
}
Throughput Benchmarks
We ran a 7-day sustained load test simulating production traffic patterns including daily peaks and synthetic burst events:
| Metric | Target | Achieved |
|---|---|---|
| Sustained throughput | 2M events/min | 2.4M events/min |
| Burst throughput (5-min) | 5M events/min | 5.8M events/min |
| End-to-end P50 latency | < 2s | 340ms |
| End-to-end P99 latency | < 5s | 1,200ms |
| Data loss rate | 0% | 0% (over 24B events) |
| Consumer lag (steady state) | < 1s | 180ms avg |
| Monthly cost (64 shards + EFO) | < $8,000 | $6,240 |
The cost breakdown: $2,688 for shard hours, $1,920 for Enhanced Fan-Out, $1,632 for PUT payload units. On-Demand mode would have cost approximately $9,100 for the same workload, so provisioned capacity with auto-scaling saved us 31%.
Shard Scaling Strategy
We implemented auto-scaling using a CloudWatch alarm that monitors IncomingRecords and WriteProvisionedThroughputExceeded:
- Scale-out trigger: WriteProvisionedThroughputExceeded > 0 for 3 consecutive minutes
- Scale-in trigger: IncomingRecords < 40% capacity for 30 consecutive minutes
- Minimum shards: 32 (baseline night traffic)
- Maximum shards: 128 (handles 4x burst)
Shard splitting doubles capacity but takes 20-30 seconds. We keep a 20% headroom buffer to absorb bursts while scaling operations complete.
Sharp Edges and Lessons Learned
Hot shards are silent killers. If your partition key distribution is skewed, one shard hits limits while others sit idle. We added a monitoring dashboard showing per-shard utilization. One tenant generating 30% of events required a composite partition key (tenant_id + random suffix) with client-side deaggregation.
Iterator age is your most important metric. GetRecords.IteratorAgeMilliseconds tells you how far behind your consumer is. We alert at 5 seconds and page at 30 seconds. If iterator age grows linearly, you need more shards or faster processing.
Kinesis Data Firehose buffering adds latency. Firehose buffers for either 60 seconds or 1MB (configurable). For our S3 cold storage path, 60-second buffering is fine. But if you route Firehose to OpenSearch, that 60-second delay matters. We reduced the buffer to 15 seconds for the search path.
Cross-region replication is not built-in. Unlike DynamoDB Global Tables, Kinesis does not replicate across regions. We built a Lambda consumer that forwards critical events to a stream in our DR region at approximately 200ms additional latency.
Conclusion
Kinesis Data Streams delivers real-time analytics at scale when you get three things right: partition key distribution for even shard utilization, Enhanced Fan-Out for dedicated consumer throughput, and proactive scaling based on iterator age rather than reactive scaling on throttle events. Our pipeline processes 2.4 million events per minute at 340ms P50 latency for $6,240/month, replacing a batch system that cost $4,800/month but delivered 15-minute latency. The 30% cost increase bought us a 750x latency improvement, and the real-time visibility paid for itself within the first week when we caught a payment processing regression in 4 seconds instead of 15 minutes.
Recommended reading

Per-Team Cost Allocation in Shared Kubernetes Clusters: From Chaos to Clarity
Implementing accurate per-namespace cost allocation in multi-tenant Kubernetes clusters, covering request vs. usage attribution, shared resource amortization, and building showback dashboards that drive accountability.

Measuring and Eliminating Toil: From 40% to 12% of Engineering Time
A systematic approach to identifying, measuring, and automating toil—the repetitive operational work that scales linearly with service growth and prevents engineers from doing creative work.

Serverless Postgres in Production: Branching, Scale-to-Zero, and the End of Database Provisioning
Running Neon serverless Postgres in production for 8 months — covering database branching workflows, scale-to-zero economics, connection pooling, and migration from RDS.

Comments
No comments yet. Be the first to share your thoughts.