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.

#aws#kinesis#streaming#analytics
Cover image for the article: Building a Real-Time Analytics Pipeline with AWS Kinesis: From Ingestion to Dashboard

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.

Kinesis Pipeline Architecture

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:

MetricTargetAchieved
Sustained throughput2M events/min2.4M events/min
Burst throughput (5-min)5M events/min5.8M events/min
End-to-end P50 latency< 2s340ms
End-to-end P99 latency< 5s1,200ms
Data loss rate0%0% (over 24B events)
Consumer lag (steady state)< 1s180ms 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%.

Kinesis Throughput Benchmarks

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.

Comments

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