ARTICLE DETAIL

资讯详情

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

Amazon Data Firehose 实战指南:基于 AWS SDK for Java 2.x 的单条与批量数据写入及指标监控

Amazon Data Firehose 实战指南:基于 AWS SDK for Java 2.x 的单条与批量数据写入及指标监控 示例工程教程后端【免费下载链接】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点击查看免费下载导读本文以javav2/example_code/firehose目录下的官方示例为蓝本系统讲解如何使用 AWS SDK for Java 2.x 对接 Amazon Data Firehose全托管实时流数据投递服务从创建 Delivery Stream、向流中写入单条记录PutRecord与批量记录PutRecordBatch到通过 CloudWatch 监控IncomingBytes、IncomingRecords、FailedPutCount等指标最后完成流的删除与资源清理。读完本文你将掌握一套可直接复制运行的 Java 数据投递完整链路以及面对 Firehose 硬性限制单请求最多 500 条记录或 4MB时的分批切分策略。Amazon Data Firehose 与示例概览Amazon Data Firehose 是一项完全托管的服务用于将实时流式数据投递到 AWS 目标服务与第三方 HTTP 端点是把流式数据可靠地加载到数据湖、数据仓库和分析服务的最便捷方式。它可以捕获、转换并把数据加载到 Amazon S3、Amazon Redshift、Amazon OpenSearchElasticsearch与 Splunk 等目标从而借助现有 BI 工具与仪表盘实现近实时分析。本示例的目标非常聚焦使用 AWS SDK for Java 2.x 操作 Data Firehose演示两种核心写操作——单条写入调用PutRecord逐条发送记录批量写入调用PutRecordBatch一次发送一批记录并处理记录数超过 API 上限时的分批问题。围绕这两类写操作示例还覆盖了流的生命周期管理创建、查询、删除与性能监控CloudWatch 指标完整代码位于 javav2/example_code/firehose/src/main/java/com/example/firehose/scenario/FirehoseScenario.java。⚠ 重要注意事项在运行任何代码之前请先了解以下几点会产生 AWS 费用运行本示例以及运行测试都可能对你的 AWS 账户产生费用请留意 Data Firehose、S3、CloudWatch 等服务的计费规则与免费额度。最小权限原则least privilege建议只为代码授予完成该任务所需的最小权限不要使用管理员级凭证。更多说明见 IAM 最佳实践。区域可用性差异本示例代码并未在每一个 AWS 区域都经过测试使用前请确认目标区域是否支持 Data Firehose。此外javav2顶层 READMEjavav2/README.md中描述了运行所有 Java 示例所需的通用前置条件JDK 版本、SDK 凭证配置等运行前请先对照检查。前置条件运行本示例需要AWS 账户与凭证参考 javav2/README.md 的 Prerequisites 章节配置好 AWS 开发环境访问密钥、区域等。JDK 21 与 Mavenpom.xml中将java.version设置为21见 pom.xml。一个可用的 Delivery Stream示例需要一个已存在的 Firehose 投递流。你可以使用 scenarios/features/firehose/README.md 提供的 CloudFormation 模板创建详见下文基础设施与模拟数据准备一节。样本数据文件sample_records.json真实网络记录模拟数据位于 javav2/example_code/firehose/src/main/resources/sample_records.json会被打包进resources并加载。项目构建配置metadata.yaml 列出该模块涉及的全部代码文件创建流、删除流、列出流、批量写入、单条写入、测试等均标记为firehose服务。构建依赖方面pom.xml 通过 BOM 统一管理 AWS SDK 版本software.amazon.awssdk:bom:2.35.10关键依赖包括依赖用途software.amazon.awssdk:firehoseData Firehose 服务客户端software.amazon.awssdk:cloudwatch读取 Firehose 投递指标software.amazon.awssdk:secretsmanager测试时从 Secrets Manager 读取配置com.fasterxml.jackson.core:jackson-databind(2.14.2)JSON 序列化/反序列化记录com.google.code.gson(2.10.1)测试中解析 Secrets Manager 返回的 JSONorg.junit.jupiter:junit-jupiter(5.11.4)集成测试框架单条写入PutRecordPutRecord是向 Firehose 投递流发送单条记录的基础 API。对应实现位于 FirehoseScenario.java#L98 的putRecord方法public static void putRecord(MapString, Object record, String deliveryStreamName) { if (record null || deliveryStreamName null || deliveryStreamName.isEmpty()) { throw new IllegalArgumentException(Invalid input: record or delivery stream name cannot be null/empty); } try { String jsonRecord new ObjectMapper().writeValueAsString(record); Record firehoseRecord Record.builder() .data(SdkBytes.fromByteArray(jsonRecord.getBytes(StandardCharsets.UTF_8))) .build(); PutRecordRequest putRecordRequest PutRecordRequest.builder() .deliveryStreamName(deliveryStreamName) .record(firehoseRecord) .build(); getFirehoseClient().putRecord(putRecordRequest); System.out.println(Record sent: jsonRecord); } catch (Exception e) { throw new RuntimeException(Failed to put record: e.getMessage(), e); } }要点解析入参校验记录为空或流名为空时直接抛出IllegalArgumentException避免无效请求打到服务端。序列化使用 Jackson 的ObjectMapper.writeValueAsString把MapString, Object转成 JSON 字符串再经SdkBytes.fromByteArray包装为 SDK 要求的字节数据。Firehose 的Record.data字段正是这种原始字节形式。UTF-8 编码StandardCharsets.UTF_8保证中文等多字节内容不会乱码。返回结果putRecord成功后服务端返回记录的唯一RecordId示例中仅打印发送的 JSON 内容作为确认。批量写入PutRecordBatch 与分批策略PutRecordBatch允许一次请求携带多条记录吞吐显著高于逐条调用。对应实现位于 FirehoseScenario.java#L131 的putRecordBatch方法public static void putRecordBatch(ListMapString, Object records, int batchSize, String deliveryStreamName) { if (records null || records.isEmpty() || deliveryStreamName null || deliveryStreamName.isEmpty()) { throw new IllegalArgumentException(Invalid input: records or delivery stream name cannot be null/empty); } ObjectMapper objectMapper new ObjectMapper(); try { for (int i 0; i records.size(); i batchSize) { ListMapString, Object batch records.subList(i, Math.min(i batchSize, records.size())); ListRecord batchRecords batch.stream().map(record - { try { String jsonRecord objectMapper.writeValueAsString(record); return Record.builder() .data(SdkBytes.fromByteArray(jsonRecord.getBytes(StandardCharsets.UTF_8))) .build(); } catch (Exception e) { throw new RuntimeException(Error creating Firehose record, e); } }).collect(Collectors.toList()); PutRecordBatchRequest request PutRecordBatchRequest.builder() .deliveryStreamName(deliveryStreamName) .records(batchRecords) .build(); PutRecordBatchResponse response getFirehoseClient().putRecordBatch(request); if (response.failedPutCount() 0) { response.requestResponses().stream() .filter(r - r.errorCode() ! null) .forEach(r - System.err.println(Failed record: r.errorMessage())); } System.out.println(Batch sent with size: batchRecords.size()); } } catch (Exception e) { throw new RuntimeException(Failed to put record batch: e.getMessage(), e); } }关键点Firehose 的硬性限制与分批切分Data Firehose 的PutRecordBatchAPI 有明确上限单次请求最多 500 条记录或总数据量最多 4MB两者任一先达到即受限。因此当记录总数超过上限时必须把请求拆分为多个小批次——这正是本方法的核心设计以batchSize示例中默认 500为步长通过subList(i, Math.min(i batchSize, records.size()))滑动切分对 5550 条样本数据前 100 条走单条PutRecord剩余 5450 条按 500 一批会被切分为 11 个批次10 批满 500 1 批 450failedPutCount用于检测部分失败响应中errorCode不为空的requestResponses即为失败记录逐条打印错误消息便于后续重试或记录。从规格文档 scenarios/features/firehose/SPECIFICATION.md 可以看到官方对这一场景的批量写入还要求优化批大小以平衡吞吐与成本并在失败时实现带指数退避与抖动的重试retries with exponential back-off and jitter。场景主流程从样本数据到指标监控FirehoseScenario.mainFirehoseScenario.java#L34把单条写入、批量写入与指标监控串成一个完整流程。运行方式java FirehoseScenario deliveryStreamName其中deliveryStreamName是目标投递流名称必填。执行步骤依次为读取并解析样本数据readJsonFile(sample_records.json)从 classpath 的/resources目录读取 JSON 文件用 Jackson 解析为ListMapString, Object。样本数据形如[ { ip_address: 42.93.101.96, timestamp: 1733762810, alert_level: Medium }, { ip_address: 186.136.101.148, timestamp: 1733762810, alert_level: High } ]每个记录包含ip_address、timestamp、alert_level三个字段模拟网络告警日志。处理单条记录sampleData.subList(0, 100)取前 100 条逐条调用putRecord任一记录失败仅打印错误、不中断整体流程。监控指标调用monitorMetrics观察写入是否产生流量见下节。处理批量记录putRecordBatch(sampleData.subList(100, sampleData.size()), 500, deliveryStreamName)把剩余记录按 500 一批投递。再次监控指标批量写入后再查一次指标对比前后数据。关闭客户端finally中依次close()Firehose 与 CloudWatch 客户端释放连接资源。客户端通过惰性单例模式构建固定使用Region.US_EAST_1firehoseClient FirehoseClient.builder() .region(Region.US_EAST_1) .build();注意如果你在CreateDeliveryStream.java中使用的是Region.US_WEST_2请确保流创建与写入使用一致的区域。用 CloudWatch 监控投递指标Firehose 会以AWS/Firehose命名空间把指标发布到 CloudWatch。示例中的monitorMetricsFirehoseScenario.java#L182监控三个指标IncomingBytes—— 流入投递流的字节数确认有数据进入IncomingRecords—— 流入投递流的记录数与IncomingBytes一起确认流量FailedPutCount—— 失败的写入次数用于发现批量写入中的问题。实现通过CloudWatchClient.getMetricStatistics查询最近 600 秒10 分钟的数据GetMetricStatisticsRequest request GetMetricStatisticsRequest.builder() .namespace(AWS/Firehose) .metricName(metricName) .dimensions(Dimension.builder().name(DeliveryStreamName).value(deliveryStreamName).build()) .startTime(startTime) .endTime(endTime) .period(60) .statistics(Statistic.SUM) .build();要点以DeliveryStreamName作为维度过滤只统计目标投递流period(60)表示按 60 秒聚合Statistic.SUM累加各数据点的和输出格式如IncomingBytes: 1234567。配套的流管理操作创建、等待激活、列出与删除完整场景还需要围绕投递流进行生命周期管理仓库中提供了三个独立示例。创建投递流ExtendedS3DestinationConfigurationCreateDeliveryStream.java 演示如何创建一条写入 S3 的投递流运行参数为java CreateDeliveryStream bucketARN roleARN streamName参数说明bucketARN数据写入目标的 S3 桶 ARNroleARN授予 Firehose 所需权限的 IAM 角色 ARNstreamName投递流名称核心代码如下ExtendedS3DestinationConfiguration destinationConfiguration ExtendedS3DestinationConfiguration.builder() .bucketARN(bucketARN) .roleARN(roleARN) .build(); CreateDeliveryStreamRequest deliveryStreamRequest CreateDeliveryStreamRequest.builder() .deliveryStreamName(streamName) .extendedS3DestinationConfiguration(destinationConfiguration) .deliveryStreamType(DirectPut) .build(); CreateDeliveryStreamResponse streamResponse firehoseClient.createDeliveryStream(deliveryStreamRequest); System.out.println(Delivery Stream ARN is streamResponse.deliveryStreamARN());这里deliveryStreamType(DirectPut)表示应用直接调用PutRecord/PutRecordBatch写入对应本场景而不是从 Kinesis Data Streams 等上游拉取。创建后流进入CREATING状态需要等待其变为ACTIVE才能接收数据。示例提供的waitForStreamToBecomeActive方法通过describeDeliveryStream轮询状态每 10 秒检查一次最多 60 次约 10 分钟状态变为ACTIVE即返回超时则退出进程。列出投递流ListDeliveryStreams.java 调用listDeliveryStreams()一次性返回当前区域所有投递流名称并逐条打印可用于确认流已创建成功。注意 SDK 的ListDeliveryStreams默认分页返回每条响应最多 10 个流名如需遍历全部流可参考官方分页机制。删除投递流DeleteStream.java 演示删除操作运行参数为流名DeleteDeliveryStreamRequest deleteDeliveryStreamRequest DeleteDeliveryStreamRequest.builder() .deliveryStreamName(streamName) .build(); firehoseClient.deleteDeliveryStream(deleteDeliveryStreamRequest);删除是异步的请求成功仅表示删除流程已开始。为了在示例结束后不遗留资源、避免持续产生费用务必执行删除操作。基础设施与模拟数据准备本场景依赖真实的 AWS 资源官方在 scenarios/features/firehose/README.md 中给出了完整的上游准备流程1. 用 CloudFormation 创建投递流使用模板创建名为FirehoseStack的栈包含 S3 桶、IAM 角色与 Data Firehose 投递流aws cloudformation create-stack \ --stack-name FirehoseStack \ --template-body file://resources/firehose-stack.yml \ --capabilities CAPABILITY_IAMCAPABILITY_IAM是必需的因为模板会创建 IAM 角色。2. 用脚本生成模拟数据运行python resources/mock_data.py依赖 Python ≥ 3.6 与faker库生成包含5550 条伪造网络记录的sample_records.json。该文件的副本被打包进 Java 模块的 resources 目录场景运行时直接读取。3. 执行各语言脚本按各语言实现运行写入脚本观察记录投递与末尾的指标输出。完整的实现规格见 SPECIFICATION.md其中对PutRecord与PutRecordBatch的输入、输出、错误处理指数退避 抖动重试、分页切分与批量大小优化都提出了明确要求。集成测试与 Secrets Manager 配置仓库为示例编写了完整的集成测试 FirehoseTest.java测试按Order顺序执行顺序测试验证内容1testCreateDeliveryStream创建投递流并轮询等待ACTIVE2testPutRecord读取样本数据对前 100 条逐条调用putRecord3testListDeliveryStreams列出账户下投递流4testDeleteStream删除新建的投递流测试所需的bucketARN、roleARN、流名前缀、文本值通过 AWS Secrets Manager 读取secret 名test/firehose也可写入本地config.propertiesSecretValues values gson.fromJson(json, SecretValues.class); bucketARN values.getBucketARN(); roleARN values.getRoleARN(); newStream values.getNewStream() java.util.UUID.randomUUID();其中流名附加了 UUID 随机后缀避免并发测试时命名冲突。运行测试mvn test⚠ 会产生 AWS 费用。完整调用链与关键代码索引能力位置核心 API单条写入FirehoseScenario.java#L98putRecord批量写入FirehoseScenario.java#L131putRecordBatch分批上限 500/4MB指标监控FirehoseScenario.java#L182CloudWatchgetMetricStatistics创建流CreateDeliveryStream.javacreateDeliveryStream 状态轮询列出流ListDeliveryStreams.javalistDeliveryStreams删除流DeleteStream.javadeleteDeliveryStream集成测试FirehoseTest.java全生命周期验证场景规格scenarios/features/firehose/SPECIFICATION.md官方功能需求总结与最佳实践先建流、再写入、最后清理投递流必须处于ACTIVE状态才能接收PutRecord/PutRecordBatch用完删除避免 S3 与 Firehose 持续计费。尊重 500 条 / 4MB 上限批量写入前务必按记录数与字节数双重评估参考示例的subList切分方式拆分请求。善用failedPutCountPutRecordBatch支持部分成功返回的errorCode/errorMessage是定位单条失败的关键。以指标验证结果IncomingBytes/IncomingRecords增长说明数据确实流入FailedPutCount为 0 说明投递健康。最小权限 区域一致为 IAM 角色只授予 Firehose、S3、CloudWatch 所需的最小权限并确保建流与写流的区域保持一致。Copyright 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 Data Firehose PutRecord 与 PutRecordBatch 批量写入 Delivery StreamAWS SDK for Java 2.x 实战使用 Amazon Data Firehose PutRecord 与 PutRecordBatch 批量写入示例工程教程后端aws-doc-sdk-examples基于 AWS SDK for Java 2.x 实现 Amazon Translate 单文本与批量文档翻译aws doc sdk examples基于 AWS SDK for Java 2.x 实现 Amazon Translate 单文本与批量文档翻译 本文以示例工程教程后端aws-doc-sdk-examples基于 AWS SDK for Java 2.x 的 Amazon Neptune 代码示例实战指南aws doc sdk examples基于 AWS SDK for Java 2.x 的 Amazon Neptune 代码示例实战指南 本文以 javav示例工程教程后端上一篇如何使用param-miner快速发现隐藏HTTP参数完整入门教程下一篇Gravity升级完全指南gh_mirrors/wo/workshop中的7.x版本升级教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表