AMQP背景下Kinesis Streams事件路由等效方案问询
Great question! Coming from an AMQP background where routing keys like widget.created or widget.updated let you bind queues to specific events, it makes total sense to look for a similar pattern in Kinesis Streams. Let’s break down your proposed options and add some practical context to help you choose the right approach.
Option 1: Separate Streams for Each Event Type
This is the closest parallel to how you’d use dedicated queues with specific routing keys in AMQP:
- How it works: Create individual Kinesis Streams for each event type (e.g.,
widget-created-stream,widget-updated-stream). Producers send events directly to the stream matching their type, and consumers only listen to the streams they care about. - Pros: No filtering logic needed in consumers—they can process every record they receive. It’s straightforward to scale individual streams based on the volume of their specific event type.
- Cons: Managing dozens (or hundreds) of streams adds operational overhead. Kinesis has default limits on the number of streams per account (though you can request increases), and each stream’s shards add to your cost footprint.
Option 2: Single Stream + Consumer-Side Filtering
This approach consolidates all events into one stream and lets consumers pick out the events they need:
- How it works: Push all widget-related events into a single stream (e.g.,
widget-events-stream), and include aneventTypefield in each record payload (like"eventType": "widget.created"). Consumers (like Lambda functions, or apps using the Kinesis Client Library/KCL) check this field and only process records matching their target event type. - Pros: Far more efficient in terms of resource management—you only maintain one stream, and shards can be shared across multiple event types. Ideal if you have many low-volume event types.
- Cons: Consumers will read all records in their shard, even if they only care about a subset. While this doesn’t add significant cost (Kinesis charges for shard-hours, not individual record reads), it does mean some extra processing logic.
Bonus: Built-In Event Filtering (Lambda + Kinesis)
If you’re using Lambda as your consumer, you can leverage Lambda event filtering to replicate AMQP’s routing key behavior even closer. This lets you configure Lambda to only invoke when a record matches your criteria, so you don’t waste function executions on irrelevant events.
For example, here’s a filter pattern you’d set up in your Lambda trigger to only process widget.created events:
{ "filterPatterns": [ { "eventType": [ "widget.created" ] } ] }
This way, Lambda automatically drops any records that don’t match the pattern before invoking your function—just like how an AMQP queue only receives messages matching its bound routing keys.
Which Should You Choose?
- Go with separate streams if you have a small number of high-volume event types, or if you need strict isolation between event processing pipelines.
- Use a single stream + filtering (especially with Lambda’s built-in filters) if you have many event types with low-to-moderate volume, or if you want to minimize operational overhead and cost.
- For more complex routing needs (like routing to multiple consumers based on event attributes), consider combining a single stream with Kinesis Data Analytics to route records to different destinations, or using the KCL’s built-in filtering capabilities for non-Lambda consumers.
内容的提问来源于stack exchange,提问作者James

