- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome 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.
导读
本文基于本仓库 kotlin/services/firehose 目录下的完整示例代码,系统讲解如何用 AWS SDK for Kotlin 与 Amazon Data Firehose 交互。Amazon Data Firehose 是一项完全托管的服务,用于将实时流式数据可靠地交付到 AWS 目的地(如 Amazon S3)和第三方 HTTP 端点。读完本文,你将掌握 Delivery Stream 的创建、列表查询、单条记录写入(PutRecord)、批量写入(PutRecordBatch)与删除的完整 API 调用链,并理解其配套的测试机制与前置条件。
示例总览
本目录的 README 罗列了 5 个自定义示例,全部位于 src/main/kotlin/com/kotlin/firehose/ 目录下:
| 示例 | 文件名 | 功能 |
|---|---|---|
| CreateDeliveryStream | CreateDeliveryStream.kt | 创建一条交付流(Delivery Stream) |
| DeleteStream | DeleteStream.kt | 删除一条交付流 |
| ListDeliveryStreams | ListDeliveryStreams.kt | 列出账户下所有交付流 |
| PutBatchRecords | PutBatchRecords.kt | 批量写入多条数据记录,并用响应对象逐条核对写入结果 |
| PutRecord | PutRecord.kt | 向交付流写入单条数据记录 |
此外目录中还包含两个数据模型辅助类 StockTrade.kt 与 StockTradeGenerator.kt,用于批量写入示例中模拟实时股票交易数据流。
运行前置条件
根据 kotlin/README.md 的说明,运行本目录示例需要满足:
- 一个有效的 AWS 账户;
- 已按 AWS SDK for Kotlin Developer Guide 配置好默认凭证与默认 AWS Region;
- 推荐使用 Gradle 搭建 AWS SDK for Kotlin 项目的开发环境(参见 Get started with the AWS SDK for Kotlin)。
需要特别注意的几点:
- 费用提示:运行这些示例代码或测试可能对你的 AWS 账户产生费用;
- 最小权限原则:建议按最小权限(least privilege)授予代码运行所需的最低 IAM 权限;
- 区域可用性:示例代码并未在所有 AWS 区域测试过,某些服务可能仅在特定区域可用。从源码看,示例统一将客户端 region 硬编码为
us-west-2(如 CreateDeliveryStream.kt),请确保你的账户支持在us-west-2使用 Data Firehose。
创建 Delivery Stream
创建交付流是使用 Firehose 的第一步。该示例要求三个命令行参数:
Usage: <bucketARN> <roleARN> <streamName> Where: bucketARN - 数据流写入的 Amazon S3 桶的 ARN。 roleARN - IAM 角色的 ARN,该角色持有 Kinesis Data Firehose 所需的权限。 streamName - 交付流的名称。核心代码位于 CreateDeliveryStream.kt:
suspend fun createStream( bucketARNVal: String?, roleARNVal: String?, streamName: String?, ) { val destinationConfiguration = ExtendedS3DestinationConfiguration { bucketArn = bucketARNVal roleArn = roleARNVal } val request = CreateDeliveryStreamRequest { deliveryStreamName = streamName extendedS3DestinationConfiguration = destinationConfiguration deliveryStreamType = DeliveryStreamType.DirectPut } FirehoseClient.fromEnvironment { region = "us-west-2" }.use { firehoseClient -> val streamResponse = firehoseClient.createDeliveryStream(request) println("Delivery Stream ARN is ${streamResponse.deliveryStreamArn}") } }关键点解读:
ExtendedS3DestinationConfiguration指明数据的目的地是 Amazon S3,必须同时提供 S3 桶 ARN(bucketArn)和具有写入权限的 IAM 角色 ARN(roleArn);DeliveryStreamType.DirectPut表示该流通过PutRecord/PutRecordBatch直接写入(而不是从 Kinesis Data Streams 等上游摄取);FirehoseClient.fromEnvironment { region = "us-west-2" }从环境变量/凭证链读取凭证,并显式指定区域;- 创建成功后,响应中的
deliveryStreamArn即为新流的 ARN。
需要注意,创建的流需要一段时间才能进入 ACTIVE 状态。在仓库的测试中,创建后等待了 15 分钟才继续写入操作(见下文测试章节)。
列出全部 Delivery Stream
ListDeliveryStreams.kt 展示了如何分页列出账户下的所有交付流:
suspend fun listStreams() { FirehoseClient.fromEnvironment { region = "us-west-2" }.use { firehoseClient -> val response = firehoseClient.listDeliveryStreams(ListDeliveryStreamsRequest {}) response.deliveryStreamNames?.forEach { item -> println("The delivery stream name is $item") } } }该示例无需命令行参数,直接调用listDeliveryStreams,遍历响应中的deliveryStreamNames并打印每个流的名称。生产环境中,当流数量较多时,可结合ExclusiveStartDeliveryStreamName与Limit参数实现真正的分页遍历。
写入单条记录(PutRecord)
PutRecord.kt 演示了如何向交付流写入一条数据。命令行用法:
Usage: <textValue> <streamName> Where: textValue - 写入数据流的文本内容。 streamName - 数据流名称。核心代码位于 PutRecord.kt:
suspend fun putSingleRecord( textValue: String, streamName: String?, ) { val bytes = textValue.toByteArray() val recordOb = Record { data = bytes } val request = PutRecordRequest { deliveryStreamName = streamName record = recordOb } FirehoseClient.fromEnvironment { region = "us-west-2" }.use { firehoseClient -> val recordResponse = firehoseClient.putRecord(request) println("The record ID is ${recordResponse.recordId}") } }要点:
- Firehose 的
Record只包含一个data字段,类型为字节数组(ByteArray),因此写入前需要将文本内容编码为字节,textValue.toByteArray()默认使用 UTF-8 编码; PutRecordRequest需要同时指定目标deliveryStreamName和record;- 响应中的
recordId由服务端生成,可用于后续追踪该条记录的处理状态(例如配合 S3 对象路径中的 Firehose 记录 ID)。
批量写入记录(PutRecordBatch)
PutBatchRecords.kt 演示了批量写入场景——这是吞吐量敏感的生产场景中推荐的做法。该示例模拟了一个股票交易行情源,持续生成交易数据并分批发送:
suspend fun addStockTradeData(streamName: String?) { try { val recordList = mutableListOf<Record>() // 每 100 毫秒发送一条模拟股票交易 val stockTradeGenerator = StockTradeGenerator() val index = 100 // 用 StockTrade 数据填充列表 for (x in 0 until index) { val trade = stockTradeGenerator.randomTrade val bytes = trade.toJsonAsBytes() val myRecord = Record { data = bytes } println("Adding trade: $trade") recordList.add(myRecord) delay(100) } val request = PutRecordBatchRequest { deliveryStreamName = streamName records = recordList } FirehoseClient.fromEnvironment { region = "us-west-2" }.use { firehoseClient -> val recordResponse = firehoseClient.putRecordBatch(request) println("The number of records added is ${recordResponse.requestResponses?.size}") } } catch (e: InterruptedException) { println(e.localizedMessage) exitProcess(0) } }核心要点:
- 数据由 StockTradeGenerator.kt 生成:它内置了 25 只常见股票的基准价格(如 AAPL 119.72、GOOG 527.83 等),以 ±20% 的随机波动生成价格,以 40% 概率生成 SELL、其余为 BUY,并在 1~10000 股之间随机生成交易数量,同时用
AtomicLong生成递增的交易 ID; - 每条交易经 StockTrade.kt 的
toJsonAsBytes()方法(基于 JacksonObjectMapper)序列化为 JSON 字节数组后封装进Record; - 循环内使用
delay(100)(kotlinx.coroutines)模拟 100 毫秒的发送间隔; PutRecordBatchRequest.records接收List<Record>,一次请求最多携带 500 条记录或 4 MB 数据(服务端限制);- 响应中的
requestResponses数组与请求记录一一对应,README 特别强调可通过它逐条核对每条记录的写入结果(检查是否存在ErrorCode),这是批量写入场景中定位失败记录的关键手段。
删除 Delivery Stream
DeleteStream.kt 演示了清理操作,命令行只需要一个参数:
Usage: <streamName> Where: streamName - 交付流名称。核心代码位于 DeleteStream.kt:
suspend fun delStream(streamName: String) { val request = DeleteDeliveryStreamRequest { deliveryStreamName = streamName } FirehoseClient.fromEnvironment { region = "us-west-2" }.use { firehoseClient -> firehoseClient.deleteDeliveryStream(request) println("Delivery Stream $streamName is deleted") } }删除是幂等的清理操作,配合上面三个示例,即可构成“创建 → 写入 → 列出 → 删除”的完整生命周期闭环。
测试与验证机制
仓库为该目录提供了 JUnit 集成测试 FirehoseTest.kt,它按顺序编排了全部 5 个 API 的端到端验证:
createDeliveryStreamTest(Order 1):创建一条带随机 UUID 后缀的交付流;putRecordsTest(Order 2):先等待 15 分钟让流进入可用状态,再调用putSingleRecord写入单条记录;putBatchRecordsTest(Order 3):调用addStockTradeData批量写入模拟股票数据;listDeliveryStreamsTest(Order 4):调用listStreams验证列表接口;deleteStreamTest(Order 5):调用delStream删除测试创建的流。
测试的输入参数并非硬编码,而是从 AWS Secrets Manager 中的test/firehose密钥读取(bucketARN、roleARN、newStream、textValue),再通过 Gson 反序列化到SecretValues数据类。测试使用了@TestInstance(Lifecycle.PER_CLASS)与@TestMethodOrder(OrderAnnotation::class)确保顺序执行,并用runBlocking包装协程。
这与 kotlin/README.md 描述的运行方式一致:既可以在 IntelliJ 等 IDE 中直接运行 JUnit 测试,也可以从命令行执行;测试运行时会在控制台输出Test N passed之类的进度信息。需要注意的是,运行这些测试同样可能产生 AWS 账户费用,且要求对应的 Secrets Manager 密钥已存在。
运行方式与延伸资源
运行示例:针对每个示例文件,先设置 AWS 凭证,再按各文件注释中的Usage传入对应参数执行。示例源码中的main函数均为suspend fun,说明需要运行在协程上下文中(可直接用runBlocking包装,参考测试代码的写法)。
延伸学习资源(README 中给出,本文不展开):
- Data Firehose User Guide
- Data Firehose API Reference
- SDK for Kotlin Data Firehose reference
若要继续深入,仓库 kotlin 根目录还提供了完整的 SDK for Kotlin 示例总览与 Docker 镜像(Beta)使用说明;本目录的 README.md 则是本次讲解的原始依据文档。
Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. SPDX-License-Identifier: Apache-2.0
- 示例工程
- 教程
- 后端
【免费下载链接】aws-doc-sdk-examples
Welcome 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.
相关推荐
AWS SDK for Java 2.x 实战:使用 Amazon Data Firehose PutRecord 与 PutRecordBatch 批量写入 Delivery Stream
AWS SDK for Java 2.x 实战:使用 Amazon Data Firehose PutRecord 与 PutRecordBatch 批量写入
示例工程教程后端使用 AWS SDK for JavaScript (v3) 操作 Amazon Transcribe:转写任务创建、查询与删除实战
使用 AWS SDK for JavaScript v3 操作 Amazon Transcribe:转写任务创建、查询与删除实战 Amazon Transcri
示例工程教程后端Klipper省耗材指南:3D打印浪费少一半
Klipper省耗材指南:3D打印浪费少一半 Klipper 固件提供三项 3D 打印省耗材工具:压力提前(Pressure Advance)让喷嘴流量与移动精
示例工程教程后端
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考