☰
使用 AWS SDK for Kotlin 操作 Amazon Data Firehose:创建、写入与删除 Delivery Stream 实战指南
2026/9/26 15:47:58 网站建设 项目流程
  • 示例工程
  • 教程
  • 后端

【免费下载链接】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.

项目地址:https://gitcode.com/gh_mirrors/aw/aws-doc-sdk-examples
点击查看免费下载

导读

本文基于本仓库 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/ 目录下:

示例文件名功能
CreateDeliveryStreamCreateDeliveryStream.kt创建一条交付流(Delivery Stream)
DeleteStreamDeleteStream.kt删除一条交付流
ListDeliveryStreamsListDeliveryStreams.kt列出账户下所有交付流
PutBatchRecordsPutBatchRecords.kt批量写入多条数据记录,并用响应对象逐条核对写入结果
PutRecordPutRecord.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 的端到端验证:

  1. createDeliveryStreamTest(Order 1):创建一条带随机 UUID 后缀的交付流;
  2. putRecordsTest(Order 2):先等待 15 分钟让流进入可用状态,再调用putSingleRecord写入单条记录;
  3. putBatchRecordsTest(Order 3):调用addStockTradeData批量写入模拟股票数据;
  4. listDeliveryStreamsTest(Order 4):调用listStreams验证列表接口;
  5. 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.

项目地址:https://gitcode.com/gh_mirrors/aw/aws-doc-sdk-examples
点击查看免费下载
上一篇:3分钟上手HRNet:AI面部关键点检测终极指南
下一篇:如何突破平台限制?WorkshopDL跨平台工具让你畅享所有游戏模组

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询