ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

使用 AWS SDK for Rust 操作 Amazon Kinesis:数据流创建、描述、写入与清理实战指南

使用 AWS SDK for Rust 操作 Amazon Kinesis:数据流创建、描述、写入与清理实战指南 示例工程教程后端【免费下载链接】aws-doc-sdk-examplesWelcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.项目地址https://gitcode.com/gh_mirrors/aw/aws-doc-sdk-examples点击查看免费下载导读本文围绕 AWS SDK for Rust 的 Kinesis 代码示例该文档原位于rust_dev_preview/examples/kinesis/现已迁移至rustv1/examples/kinesis/展开系统讲解如何用 Rust 对 Amazon Kinesis 数据流执行创建CreateStream、删除DeleteStream、描述DescribeStream、列举ListStreams与写入PutRecord五个核心单动作操作。读完本文你将掌握每个示例的可运行命令、命令行参数语义、底层 SDK 调用链以及把它们串联成一个建流 → 写数 → 描述 → 清理完整流程的实战能力。Amazon Kinesis 与 SDK for Rust 示例概览Amazon Kinesis 用于实时收集、处理和分析视频与数据流。本仓库提供的示例聚焦于Kinesis Data Streams的五个基础服务调用每个示例是一个独立的 Rust 二进制程序源码位于 rustv1/examples/kinesis/src/bin/示例文件对应 Kinesis API功能create-stream.rsCreateStream创建数据流并指定分片数delete-stream.rsDeleteStream删除数据流describe-stream.rsDescribeStream查询流的名称、状态、分片、保留期、加密方式list-streams.rsListStreams列出当前 Region 下所有流名称put-record.rsPutRecord向指定流写入一条带分区键的数据记录这些示例通过 Cargo.toml 统一组织包名为kinesis-code-examples完整依赖如下[dependencies] aws-config { version 1.0.1, features [behavior-version-latest] } aws-sdk-kinesis { version 1.3.0 } tokio { version 1.20.1, features [full] } clap { version 4.4, features [derive] } tracing-subscriber { version 0.3.15, features [env-filter] }依赖分工值得注意aws-config负责加载共享配置与 Region 解析aws-sdk-kinesis提供 Kinesis 客户端tokiofull 特性提供异步运行时所有示例都以#[tokio::main]作为入口clap的 derive 特性用于声明式命令行参数解析tracing-subscriber初始化日志输出。前置条件运行这些示例需要满足AWS 账户与凭证必须拥有 AWS 账户并按 Getting started with the AWS SDK for Rust。成本提醒运行示例代码、特别是运行测试都可能产生 AWS 账户费用。权限最小化建议按最小权限原则仅授予完成任务所需的最低权限例如kinesis:CreateStream、kinesis:DeleteStream、kinesis:DescribeStream、kinesis:ListStreams、kinesis:PutRecord详见 AWS IAM 的 Grant least privilege 最佳实践。Region 可用性示例代码并非在每个 AWS Region 都经过测试使用前请确认目标服务在所选 Region 可用。五个单动作示例详解所有示例都遵循同一套代码骨架用 clap 定义命令行参数 → 构建 Region 提供链 → 加载共享配置 → 创建Client→ 调用异步函数。下面逐一切入。创建数据流CreateStreamcreate-stream.rs 的核心调用位于make_stream函数对应文档标注的行号 L26 起// snippet-start:[kinesis.rust.create-stream] async fn make_stream(client: Client, stream: str) - Result(), Error { client .create_stream() .stream_name(stream) .shard_count(4) .send() .await?; println!(Created stream); Ok(()) } // snippet-end:[kinesis.rust.create-stream]要点.shard_count(4)将新流的初始分片数设为 4。分片是 Kinesis 数据流的基本吞吐单位分片数量直接决定流的容量每分片提供 1 MB/s 写入、2 MB/s 读取和约 1000 条/s 的 PUT 能力生产环境应基于预估写入量规划。send().await返回Result出错时通过?向上传播最终由main返回的Result(), Error体现。删除数据流DeleteStreamdelete-stream.rs 的remove_stream函数L26 起是建流的逆操作// snippet-start:[kinesis.rust.delete-stream] async fn remove_stream(client: Client, stream: str) - Result(), Error { client.delete_stream().stream_name(stream).send().await?; println!(Deleted stream.); Ok(()) } // snippet-end:[kinesis.rust.delete-stream]DeleteStream是异步操作SDK 返回后流并不会立即消失实际删除进度需通过DescribeStream观察状态字段。描述数据流DescribeStreamdescribe-stream.rs 的show_stream函数L26 起把返回的StreamDescription逐字段打印// snippet-start:[kinesis.rust.describe-stream] async fn show_stream(client: Client, stream: str) - Result(), Error { let resp client.describe_stream().stream_name(stream).send().await?; let desc resp.stream_description.unwrap(); println!(Stream description:); println!( Name: {}:, desc.stream_name()); println!( Status: {:?}, desc.stream_status()); println!( Open shards: {:?}, desc.shards.len()); println!( Retention (hours): {}, desc.retention_period_hours()); println!( Encryption: {:?}, desc.encryption_type.unwrap()); Ok(()) } // snippet-end:[kinesis.rust.describe-stream]这里的stream_status可呈现CREATING、ACTIVE、DELETING等生命周期状态正好用于验证建流与删流的异步进展retention_period_hours反映数据在流中的保留时长默认 24 小时encryption_type则体现服务端加密配置NONE或KMS。字段通过desc.shards.len()统计开放分片数印证了 CreateStream 时设置分片数的效果。列举数据流ListStreamslist-streams.rs 的show_streams函数L22 起无需任何流名参数// snippet-start:[kinesis.rust.list-streams] async fn show_streams(client: Client) - Result(), Error { let resp client.list_streams().send().await?; println!(Stream names:); let streams resp.stream_names; for stream in streams { println!( {}, stream); } println!(Found {} stream(s), streams.len()); Ok(()) } // snippet-end:[kinesis.rust.list-streams]响应中的stream_names是当前 Region 下的流名称集合循环打印后统计总数是快速盘点环境内流资源的首选操作。写入数据记录PutRecordput-record.rs 是五个示例中参数最多的一个add_record函数L35 起演示了如何把数据与分区键组合写入// snippet-start:[kinesis.rust.put-record] async fn add_record(client: Client, stream: str, key: str, data: str) - Result(), Error { let blob Blob::new(data); client .put_record() .data(blob) .partition_key(key) .stream_name(stream) .send() .await?; println!(Put data into stream.); Ok(()) } // snippet-end:[kinesis.rust.put-record]三个概念值得展开Blob来自aws_sdk_kinesis::primitives::Blob将字符串数据包装成 SDK 的二进制载荷类型是PutRecord请求data字段的标准承载方式。分区键partition keyKinesis 依据分区键的哈希值决定数据落在哪个分片相同分区键的数据保证进入同一分片并保持相对顺序因此要按业务维度如用户 ID、设备 ID设计分区键避免单分片热点。无读回确认示例只调用put_record后打印成功未读取返回的SequenceNumber。Kinesis 的PutRecord响应会包含ShardId与SequenceNumber可用于后续定位记录示例中以最小化形式演示了写入调用本身。命令行参数与运行方式五个示例共用一个参数风格均由 clap 解析参数含义适用范围-s, --stream-name STREAM-NAME流名称必填create、delete、describe、put-record-d, --data DATA要写入流的数据必填put-record-k, --key KEY-NAME分区键名称必填put-record-r, --region REGION客户端所在 AWS Region可选全部-v, --verbose是否显示附加信息可选全部Region 解析采用统一的RegionProviderChain三级回退逻辑见各main函数let region_provider RegionProviderChain::first_try(region.map(Region::new)) .or_default_provider() .or_else(Region::new(us-west-2));即优先使用-r显式传入的 Region → 否则读取AWS_REGION环境变量默认提供者→ 仍未设置则回退到us-west-2。随后通过aws_config::from_env().region(region_provider).load().await加载共享配置并创建Client::new(shared_config)。开启-v时示例会打印 SDK 包版本aws_sdk_kinesis::meta::PKG_VERSION、最终生效的 Region、流名、数据与分区键等调试信息。运行示例在 rustv1/examples/kinesis/ 目录下使用cargo run --bin 名称 -- 参数即可运行# 1. 创建名为 my-stream 的数据流4 个分片 cargo run --bin create-stream -- -s my-stream # 2. 列举当前 Region 下的所有流 cargo run --bin list-streams # 3. 查看 my-stream 的详细信息 cargo run --bin describe-stream -- -s my-stream # 4. 以 user-001 为分区键写入一条记录 cargo run --bin put-record -- -s my-stream -k user-001 -d hello kinesis # 5. 清理资源删除数据流 cargo run --bin delete-stream -- -s my-stream需要指定 Region 时追加-r REGION如-r us-west-2想查看调试输出则加-v。运行测试可参照 rustv1/README.md 的说明测试cargo test -- --ignored可能产生 AWS 费用。串联为完整的数据流生命周期流程虽然 README 将五个示例归类为单动作但它们天然构成一个 Kinesis 数据流的完整生命周期可组合成一条可复用的操作链建流create-stream -s my-stream创建流分片数为 4随后describe-stream -s my-stream可确认状态从CREATING变为ACTIVE写数put-record -s my-stream -k partition-key -d data按分区键写入记录巡检list-streams快速盘点describe-stream深入查看分片、保留期与加密状态清理delete-stream -s my-stream删除流避免长期占用产生费用。这一链路对应的正是从 rustv1/examples/kinesis/README.md 的Single actions列表CreateStream、DeleteStream、DescribeStream、ListStreams、PutRecord 五个条目延伸出的实操路径每个环节都能在对应src/bin/*.rs文件中找到独立的实现证据。注意事项与成本提醒费用数据流按分片小时计费PutRecord 按写入量计费。创建后不删除会持续产生费用测试完毕务必执行delete-stream清理。测试集成测试以cargo test -- --ignored运行可能在账户中产生资源与费用运行前请确认。异步语义CreateStream与DeleteStream都是异步操作返回不代表立即完成需用DescribeStream的stream_status字段确认终态。权限遵循最小权限原则只为执行这些操作授予对应的kinesis:*单动作权限。延伸阅读本示例的源码清单rustv1/examples/kinesis/src/bin/ 下五个.rs文件依赖与包配置rustv1/examples/kinesis/Cargo.toml整个 SDK for Rust 示例集的前置条件与测试说明rustv1/README.md相关官方资料Kinesis Developer Guide、Kinesis API Reference、SDK for Rust Kinesis referencedocs.rs/aws-sdk-kinesisCopyright Amazon.com, Inc. or its affiliates. All Rights Reserved.SPDX-License-Identifier: Apache-2.0赞分享示例工程教程后端【免费下载链接】aws-doc-sdk-examplesWelcome to the AWS Code Examples Repository. This repo contains code examples used in the AWS documentation, AWS SDK Developer Guides, and more. For more information, see the Readme.md file below.项目地址https://gitcode.com/gh_mirrors/aw/aws-doc-sdk-examples点击查看免费下载相关推荐使用 AWS SDK for Java 2.x 操作 Amazon Kinesis数据流创建、写入与读取实战指南使用 AWS SDK for Java 2.x 操作 Amazon Kinesis数据流创建、写入与读取实战指南 Amazon Kinesis 让开发者能够实示例工程教程后端使用 AWS SDK for Kotlin 操作 Amazon Kinesis数据流完整实战指南使用 AWS SDK for Kotlin 操作 Amazon Kinesis数据流完整实战指南 Amazon Kinesis 让开发者能够实时收集、处理和分示例工程教程后端AWS SDK for PHP v3 操作 Amazon Kinesis 数据流实战指南php/example_code/kinesisAWS SDK for PHP v3 操作 Amazon Kinesis 数据流实战指南php/example_code/kinesis Amazon Ki示例工程教程后端上一篇三步把 NCM 转成 MP3ncmdumpGUI 完整使用指南下一篇ThinkPad风扇控制终极指南TPFanCtrl2双风扇快速上手教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表