多Shard连接Kinesis后调用client.destroy()无法终止脚本问题
问题:连接多个Kinesis Shard后调用
client.destroy()无法终止脚本 使用@aws-sdk/client-kinesis (3.369.0)时,连接单个Shard后调用client.destroy()能正常终止脚本,但连接多个Shard后,调用该方法脚本会持续运行。事件消费功能正常,仅终止环节存在问题。
初始化代码
const region = ""; const accountId = ""; const streamName = "" const consumerName = ""; const streamARN = "arn:aws:kinesis:" + region + ":" + accountId + ":stream/" + streamName; const client = new KinesisClient({ region }); // 恢复已有消费者或创建新消费者 const { Consumers } = await client.send( new ListStreamConsumersCommand({ streamARN }) ); let consumer = Consumers.find((i) => i.ConsumerName == consumerName); if (!consumer) { const { Consumer } = await client.send( new RegisterStreamConsumerCommand({ streamARN, consumerName }) ); consumer = Consumer; }
正常终止场景(单个Shard)
// 获取Shards const { Shards } = await client.send( new ListShardsCommand({ streamARN, streamName, }) ); // 订阅Shards for (const shard of Shards) { const { EventStream } = await client.send( new SubscribeToShardCommand({ ConsumerARN: consumer.ConsumerARN, ShardId: shard.ShardId, StartingPosition: { Type: "LATEST" }, }) ); console.log("subscribed", shard.ShardId); break; // 只订阅第一个Shard } setTimeout(() => { console.log('stopping..') client.destroy(); }, 30000)
输出:
subscribed shardId-000000000000 stopping..
脚本可正常终止。
无法终止场景(多个Shard)
移除上述代码中的break后:
// 获取Shards const { Shards } = await client.send( new ListShardsCommand({ StreamARN, StreamName, }) ); // 订阅Shards for (const shard of Shards) { const { EventStream } = await client.send( new SubscribeToShardCommand({ ConsumerARN: consumer.ConsumerARN, ShardId: shard.ShardId, StartingPosition: { Type: "LATEST" }, }) ); console.log("subscribed", shard.ShardId); } setTimeout(() => { console.log('stopping..') client.destroy(); }, 30000)
输出:
subscribed shardId-000000000000 subscribed shardId-000000000001 subscribed shardId-000000000002 subscribed shardId-000000000003 stopping..
脚本持续运行无法终止。
原因分析
每个SubscribeToShardCommand返回的EventStream是独立的异步流,会占用客户端连接资源并维持事件循环活跃。仅调用client.destroy()不会自动关闭所有已创建的EventStream,导致事件循环仍有未处理资源,脚本无法退出。
解决方案
保存所有订阅得到的EventStream实例,终止时先手动关闭每个流,再调用client.destroy():
// 获取Shards const { Shards } = await client.send( new ListShardsCommand({ streamARN, streamName, }) ); // 保存所有EventStream实例 const eventStreams = []; // 订阅Shards for (const shard of Shards) { const { EventStream } = await client.send( new SubscribeToShardCommand({ ConsumerARN: consumer.ConsumerARN, ShardId: shard.ShardId, StartingPosition: { Type: "LATEST" }, }) ); eventStreams.push(EventStream); console.log("subscribed", shard.ShardId); } setTimeout(async () => { console.log('stopping..') // 先关闭所有EventStream for (const stream of eventStreams) { await stream.destroy(); } // 再销毁客户端 client.destroy(); }, 30000)
此处理会正确关闭所有异步流,释放事件循环资源,脚本即可正常终止。
内容的提问来源于stack exchange,提问作者keja
相关产品推荐
相关产品推荐

