如何从Amazon Kinesis按需流中顺序读取数据?
AWS Kinesis按需流实现全局顺序读取的方案
问题背景
基于AWS Kinesis官方示例实现流的生产和消费逻辑后,已成功向按需流写入50条数据,但遇到全局顺序读取的核心问题:
- 按需流会自动创建多个分片,
get-recordsAPI必须指定分片ID才能读取数据 - 遍历所有分片读取时,仅能保证单个分片内部数据有序,无法实现整个流的全局顺序
- 不指定分片ID调用API会直接抛出异常
- 切换为单分片预置模式可实现顺序读取,但希望保留按需模式的同时解决顺序问题
当前消费代码如下:
public static void getStockTrades(KinesisClient kinesisClient, String streamName) { String lastShardId = null; // Retrieve the Shards from a Stream DescribeStreamRequest describeStreamRequest = DescribeStreamRequest.builder() .streamName(streamName) .build(); List<Shard> shards = new ArrayList<>(); DescribeStreamResponse streamRes; do { streamRes = kinesisClient.describeStream(describeStreamRequest); shards.addAll(streamRes.streamDescription().shards()); if (shards.size() > 0) { lastShardId = shards.get(shards.size() - 1).shardId(); } } while (streamRes.streamDescription().hasMoreShards()); shards.forEach(shard -> { String shardIterator; String shardId = shard.shardId(); GetShardIteratorRequest itReq = GetShardIteratorRequest.builder() .streamName(streamName) .shardIteratorType("TRIM_HORIZON") .shardId(shardId) .build(); GetShardIteratorResponse shardIteratorResult = kinesisClient.getShardIterator(itReq); shardIterator = shardIteratorResult.shardIterator(); // Create new GetRecordsRequest with existing shardIterator. // Set maximum records to return to 1000. GetRecordsRequest recordsRequest = GetRecordsRequest.builder() .shardIterator(shardIterator) .limit(1000) .build(); GetRecordsResponse result = kinesisClient.getRecords(recordsRequest); // Put result into record list. Result may be empty. List<Record> recordsList = result.records(); System.out.printf("Shared id: %s, Number of records: %d%n", shardId, recordsList.size()); for (Record record : recordsList) { SdkBytes byteBuffer = record.data(); System.out.printf("Seq No: %s - %s%n", record.sequenceNumber(), new String(byteBuffer.asByteArray())); } }); }
解决方案
要在Kinesis按需流中实现全局顺序读取,核心是通过分区键控制数据分片路由,并处理分片分裂后的继承链读取逻辑,具体步骤如下:
1. 生产者端统一分区键
在调用PutRecord或PutRecords写入数据时,为所有记录指定相同的Partition Key。Kinesis会根据Partition Key的哈希值路由分片,相同Partition Key的记录会被分配到同一个分片(除非该分片因吞吐量阈值触发分裂),从源头保证全局顺序。
示例修改(生产者端):
// 所有记录使用同一个全局分区键 String globalPartitionKey = "single-order-key"; PutRecordRequest putRecordRequest = PutRecordRequest.builder() .streamName(streamName) .data(SdkBytes.fromUtf8String(recordContent)) .partitionKey(globalPartitionKey) .build(); kinesisClient.putRecord(putRecordRequest);
2. 消费端处理分片分裂场景
如果按需流的单个分片因吞吐量触发分裂,会生成两个子分片。此时需要按照分片继承链顺序读取:先读取原分片的所有数据,再依次读取子分片的数据,才能保证全局顺序。
修改后的消费代码:
public static void getStockTradesInOrder(KinesisClient kinesisClient, String streamName) { // 获取所有分片 List<Shard> allShards = getAllShards(kinesisClient, streamName); // 遍历所有根分片(无父分片的分片),递归读取其所有子分片 for (Shard shard : allShards) { if (shard.parentShardId() == null) { readShardAndChildren(kinesisClient, streamName, shard); } } } // 递归读取分片及其所有子分片 private static void readShardAndChildren(KinesisClient kinesisClient, String streamName, Shard shard) { // 先读取当前分片数据 readSingleShard(kinesisClient, streamName, shard); // 查找当前分片的子分片并按顺序递归读取 List<Shard> childShards = getAllShards(kinesisClient, streamName).stream() .filter(s -> shard.shardId().equals(s.parentShardId())) .sorted(Comparator.comparing(Shard::shardId)) .toList(); for (Shard child : childShards) { readShardAndChildren(kinesisClient, streamName, child); } } // 读取单个分片数据(保留原逻辑) private static void readSingleShard(KinesisClient kinesisClient, String streamName, Shard shard) { String shardId = shard.shardId(); GetShardIteratorRequest itReq = GetShardIteratorRequest.builder() .streamName(streamName) .shardIteratorType("TRIM_HORIZON") .shardId(shardId) .build(); GetShardIteratorResponse shardIteratorResult = kinesisClient.getShardIterator(itReq); String shardIterator = shardIteratorResult.shardIterator(); GetRecordsRequest recordsRequest = GetRecordsRequest.builder() .shardIterator(shardIterator) .limit(1000) .build(); GetRecordsResponse result = kinesisClient.getRecords(recordsRequest); List<Record> recordsList = result.records(); System.out.printf("分片ID: %s, 记录数: %d%n", shardId, recordsList.size()); for (Record record : recordsList) { SdkBytes byteBuffer = record.data(); System.out.printf("序列号: %s - %s%n", record.sequenceNumber(), new String(byteBuffer.asByteArray())); } } // 获取所有分片的工具方法 private static List<Shard> getAllShards(KinesisClient kinesisClient, String streamName) { List<Shard> shards = new ArrayList<>(); String exclusiveStartShardId = null; do { DescribeStreamRequest describeStreamRequest = DescribeStreamRequest.builder() .streamName(streamName) .exclusiveStartShardId(exclusiveStartShardId) .build(); DescribeStreamResponse streamRes = kinesisClient.describeStream(describeStreamRequest); shards.addAll(streamRes.streamDescription().shards()); exclusiveStartShardId = streamRes.streamDescription().hasMoreShards() ? shards.get(shards.size() - 1).shardId() : null; } while (exclusiveStartShardId != null); return shards; }
3. 关键注意事项
- 该方案会牺牲按需流的水平扩展性,所有数据集中在一个分片(或其分裂后的子分片链),吞吐量受限于单个分片的上限(按需模式下单分片最大1MB/s写入、2MB/s读取)
- 如果业务需要全局顺序且高吞吐量,Kinesis按需流并非最优选择,建议考虑单分片预置模式或其他支持全局顺序的消息队列服务
- 分片分裂后,必须严格遵循「原分片→子分片」的读取顺序,否则会出现数据乱序
内容的提问来源于stack exchange,提问作者Amudhan
相关产品推荐
相关产品推荐

