You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何从Amazon Kinesis按需流中顺序读取数据?

AWS Kinesis按需流实现全局顺序读取的方案

问题背景

基于AWS Kinesis官方示例实现流的生产和消费逻辑后,已成功向按需流写入50条数据,但遇到全局顺序读取的核心问题:

  • 按需流会自动创建多个分片,get-records API必须指定分片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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 06:52:54