尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

AWS SDK for Java V2 批处理工具(BatchManager)设计决策全解析:从决策日志到 SQS 源码落地

AWS SDK for Java V2 批处理工具(BatchManager)设计决策全解析:从决策日志到 SQS 源码落地 AWS SDK for Java V2 批处理工具BatchManager设计决策全解析从决策日志到 SQS 源码落地【免费下载链接】aws-sdk-java-v2The official AWS SDK for Java - Version 2项目地址: https://gitcode.com/GitHub_Trending/aw/aws-sdk-java-v2本篇技术指南以 AWS SDK for Java V2 仓库中的《Decision Log for SDK V2 Batch Utility》为核心骨架系统梳理自动批处理Automatic Request Batching工具在 2021 年 6 月至 7 月期间做出的全部架构与 API 设计决策并结合docs/design/core/batch-manager/Design.md设计文档与 SQS 模块下batchmanager包的源码实现逐一验证这些决策在SqsAsyncBatchManager、RequestBatchManager、BatchOverrideConfiguration等类中的落地方式。读者读完后将完整掌握 BatchManager 的命名由来、请求映射机制、缓冲与刷盘flush策略、重试与同步/异步接口划分的设计逻辑以及如何通过源码理解其真实调用链。背景为什么要一个批处理工具决策日志中并未单独记录需求来源但与其配套的 Design.md 说明了动机部分客户需要在多个 AWS 服务上执行批量写操作而 v2 SDK 缺少相应能力成为采用 v2 的阻碍。v1 中 SQS 服务以AmazonSQSBufferedAsyncClient形式实现了此功能但等价能力没有移植到 v2。由于许多 AWS 服务本身已提供批量写 API如 SQS 的sendMessageBatch、Kinesis 的putRecords因此可以设计一套通用的自动批处理方案不局限于 SQS任何受益的服务都可以接入。BatchManager 的目标是帮助客户降低成本减少请求数、提升性能、简化实现。调用方先经由管理器缓冲buffer请求再由管理器以批请求形式发送给服务客户端侧缓冲被泛化实现并允许达到服务允许的最大批量请求数例如 SQS 最多 10 条。命名决策为什么叫{Service}BatchManager2021 年 7 月 15 日日志中最新的会议讨论的唯一决策是工具命名。结论命名为{Service}BatchManager如 SQS 即SqsBatchManager/SqsAsyncBatchManager。理由如下批处理工具运行在低层服务客户端low-level service client之上因此它不太像一个 utility但它又不替换低层客户端实现的方法因此也不是 enhanced client增强客户端它的运作方式与TransferManager传输管理器类似因此命名为BatchManager与现有生态保持一致。该决策的最终形态可以在源码中验证。SQS 模块下公开接口为software.amazon.awssdk.services.sqs.batchmanager.SqsAsyncBatchManager见 SqsAsyncBatchManager.java标注了SdkPublicApi从SqsAsyncClient出发即可获得实例——测试代码SqsAsyncBatchManagerTest.java第 68 行展示了client.batchManager()的用法见 SqsAsyncBatchManagerTest.java。独立工具 vs 直接写在客户端上六个关键设计决策2021-06-296 月 29 日的会议围绕批处理初始设计对应当时 PR #2563展开产生了 7 个关闭决策Closed Decisions。这些决策奠定了 BatchManager 的整体架构逐条分析如下。决策 1独立 utility而非直接实现于客户端批处理方法不直接写在客户端上而是放在独立的工具类中。这样与其他 API如 Waiters保持一致如果要改变这一点就需要改动所有 API。Design.md 的 FAQ 中给出了更完整的方案对比Option 1直接在低层客户端上实现批处理如sqsAsync.automaticSendMessageBatch(...)优点是更容易被发现Option 2独立手工编写的高层库更贴近 v1 的AmazonSQSBufferedAsyncClient迁移成本低Option 3独立的批处理管理器类BatchManager某服务所有批处理特性都内聚在该服务的工具类中与 Waiters 一致易于配置和扩展到多个服务。最终选择Option 3理由与决策日志一致紧密跟随 v2 的整体风格尤其类似 Waiters 抽象且能在不增加使用复杂度的前提下跨服务扩展。落地形态就是SqsAsyncBatchManager接口及其默认实现DefaultSqsAsyncBatchManager见 DefaultSqsAsyncBatchManager.java。决策 2使用batch.json存储默认值与服务特定的批处理方法使用配置文件batch.json或batcher.json来存放默认值和特定服务的批处理方法不把提供这些值的负担转嫁给客户这与 Waiters 等其他 API 一致并考虑在整个 SDK 范围内推广。决策 3为批响应创建包装类避免泛型嵌套为CompletableFutureBatchResponseSendMessageBatchResultEntry这类复杂泛型创建包装类如SQSBatchResponse使客户更容易理解。从当前源码看sendMessage()直接返回CompletableFutureSendMessageResponse1:1 对应单个响应包装思路体现在响应映射与按 batchKey 分组的设计中。决策 4支持手动 flush且可选择 flush 指定 bufferBatchManager必须支持手动刷新manual flush以及指定某个 buffer 的刷新。这保证了与 v1 的功能对齐并给客户提供额外能力。该决策在 Design.md 中被列为核心使用场景自动批处理 手动刷新源码中的实现则更为完整DefaultSqsAsyncBatchManager.close()会先执行closeAndDispatch()排空所有缓冲区、派发所有批请求再等待超时最后cancelPending()取消仍在进行的请求见 DefaultSqsAsyncBatchManager.javaRequestBatchManager内部还实现了按字节大小溢出即刷盘的路径extractBatchIfSizeExceeded见 RequestBatchManager.java。决策 5批重试由客户端低层 client处理批请求的重试retry由客户端处理这与 SDK 中其他部分的处理方式一致这些重试也可能被批量处理These retries could possibly be batched as well。从源码看SendMessageBatchManager将组装好的SendMessageBatchRequest直接交给asyncClient.sendMessage(batchRequest)重试逻辑完全复用低层异步客户端的能力见 SendMessageBatchManager.java。决策 6分别维护 sync 与 async 两套接口批处理工具同时提供同步与异步接口即使当前它们看起来非常相似也要分别维护原因有二预期接口会随功能增加而分化两者的 builder 也不同。落地即SqsBatchManager与SqsAsyncBatchManager两个接口并存。SqsAsyncBatchManager的 Builder 提供了client(SqsAsyncClient)、overrideConfiguration(...)、scheduledExecutor(...)等方法见 SqsAsyncBatchManager.java。决策 7异步客户端不串行发送请求针对异步客户端是否通过一条条串行发送来应对限流异常throttling的问题答案是否——这会违背异步客户端的初衷。实际并发上限只受客户端允许的最大连接数限制并且低层客户端本身已经提供了限流throttling支持。这一设计保证 BatchManager 在重负载下仍能充分利用连接池并发发送批次。遗留的开放问题6/29该次会议同时记录了三个新开放问题Open Decisions在 7 月 13 日的会议上得到解决批处理工具是否只使用一个sendMessage()方法以保持请求与响应 1:1 关联若否则客户如何关联请求与响应是否接受 streams 或 iterators 传入sendMessages()与问题 1 互斥批处理工具该如何命名7 月 15 日已解决请求映射与边界情况六个后续决策2021-07-137 月 13 日的后续会议解决了上轮遗留问题并处理了新议题产生 6 个关闭决策。决策 1尊重RequestOverrideConfiguration并按 override 分组批处理是否尊重SendMessageRequest中的RequestOverrideConfiguration字段是。带有不同 override 配置的请求将分别批处理分组依据是AwsRequestOverrideConfiguration自带的equals/hashCode方法。源码验证SendMessageBatchManager.getBatchKey()的实现为——若请求带 override 配置则 batchKey 为queueUrl overrideConfiguration.hashCode()否则仅为queueUrl见 SendMessageBatchManager.java。即队列 URL override 配置哈希共同决定请求归属于哪个批次。组装批请求时代码取批次中第一条请求的 override 配置注释明确所有请求必须具有相同的 overrideConfiguration并注入 BatchManager 专属的 user-agent API 名称ApiName hll、versionabm即 Automatic Batching Manager见 SendMessageBatchManager.java 与 RequestBatchManager.java。决策 2通过名称匹配生成请求映射SendMessageRequest与SendMessageBatchRequestEntry之间的映射采用名称匹配name matching加某种定制配置如batch.json自动生成。源码中的手工映射即按名称逐字段搬运createSendMessageBatchRequestEntry()将messageBody、delaySeconds、messageAttributes、messageSystemAttributesWithStrings、messageDeduplicationId、messageGroupId一一映射为SendMessageBatchRequestEntry的同名字段见 SendMessageBatchManager.java。决策 3字段不一致时先抛异常方案 a如果SendMessageRequest中存在SendMessageBatchRequestEntry没有的字段或反之怎么办当时讨论了四个选项(a) 接收SendMessageRequest若存在无法映射到SendMessageBatchRequestEntry的字段按名称则抛异常(b) 同 (a)但额外提供(QueueUrl, SendMessageBatchRequestEntry)方法(c) 仅接收(QueueUrl, SendMessageBatchRequestEntry)(d) 接收SendMessageRequest若不匹配则构建失败并调查或将字段加入白名单排除。结论现阶段采用 (a)但若字段开始分化则准备好转向 (b)。这体现了先保持简单、预留演进路径的设计原则。决策 4暂不允许客户覆盖映射行为是否允许客户覆盖SendMessageRequest → SendMessageBatchRequestEntry的映射以改变固执的opinionated行为或获得转换过程可见性暂时不允许但未来若字段分化可能需要支持。决策 5只使用sendMessage()保持 1:1 关联批处理工具是否只提供一次只发送一个请求的sendMessage()方法使请求与响应保持 1:1 关联是现阶段如此但对该领域保持开放、欢迎客户反馈。这也解决了 6/29 的开放问题 1。SqsAsyncBatchManager接口中sendMessage(SendMessageRequest)返回CompletableFutureSendMessageResponse正是 1:1 关联的体现见 SqsAsyncBatchManager.java。SendMessageBatchManager.mapBatchResponse()展示了 1:1 回填的完整逻辑批量响应中成功的条目按 id 映射为SendMessageResponse携带messageId、md5OfMessageBody、sequenceNumber等并回填responseMetadata与sdkHttpResponse失败的条目则按 id 构造SqsException含errorCode、errorMessage异常化对应的 future见 SendMessageBatchManager.java。决策 6暂不接受 streams / iterators是否接受 streams 或 iterators 传给sendMessages()暂时不接受但保持开放。这也解决了 6/29 的开放问题 2。该决策与决策 5 互斥联动——既然保持 1:1 关联就不接受流式/迭代器批量输入。客户如需发送多条消息循环调用sendMessage()即可。设计文档中的其他关键决策FAQ 补充Design.md 的 FAQ 还回答了决策日志未展开的三个问题属于对上述决策的必要补充为什么只支持一次一条而非列表或流支持单个sendMessage()方法让客户更容易将请求消息与各自响应相关联例如 SQS 中发送一个SendMessageRequest返回一个SendMessageResponse而不是批响应包装类。多条消息只需循环调用sendMessage()因此未来在保留该方法的基础上再支持流或列表也容易。为什么同时支持同步与异步既保证两种客户端的 API 不分叉又与 v1 的缓冲客户端对齐。实现上只是分别调用同步/异步客户端的底层方法。为什么同步和异步的sendMessage()都返回CompletableFuturesendMessage()会先自动缓冲请求直到缓冲区满或超时因此同步客户端的调用可能阻塞长达整个超时周期。为减少长时间阻塞两个接口的sendMessage()都返回CompletableFuture在底层批请求发送并收到响应后完成。同步与异步管理器的真正区别在于底层调用的是各自客户端的sendMessageBatch。从决策到实现BatchManager 在源码中的完整形态缓冲、批键与定时刷盘的调用链核心抽象类RequestBatchManagerRequestT, ResponseT, BatchResponseT见 RequestBatchManager.java把自动批处理实现为通用骨架子类只需实现三个抽象方法getBatchKey(request)计算请求归属的批次键batchAndSend(identifiedRequests, batchKey)将一组带 id 的请求组装为批请求并发送mapBatchResponse(batchResponse)把批响应按条目 id 拆分为成功响应或异常。batchRequest()的核心流程是按批键放入缓冲 Map → 若超过字节上限maxBatchBytesSize先局部刷盘 → 调度定时刷盘scheduleAtFixedRate周期即sendRequestFrequency→ 若达到条数上限立即刷盘见 RequestBatchManager.java。定时刷盘由ScheduledExecutorService驱动Builder 支持注入自定义 executor见 SqsAsyncBatchManager.java。BatchOverrideConfiguration可调参数与默认值配置类BatchOverrideConfiguration见 BatchOverrideConfiguration.java对应决策默认值放配置文件、不转嫁给客户的落地——所有值可选未指定时使用默认值配置项默认值说明maxBatchSize10单个出站批请求SendMessageBatchRequest、ChangeMessageVisibilityBatchRequest、DeleteMessageBatchRequest最多容纳的条目数必须 ≤ 10构造时校验sendRequestFrequency200ms出站批次保持打开、等待更多请求的最长时间若提前达到maxBatchSize则立即发送。增大该值可减少请求数、提升吞吐但会增加平均消息延迟receiveMessageVisibilityTimeout未设置使用队列配置接收消息时自定义可见性超时receiveMessageMinWaitDuration50ms接收请求的最小等待时间避免线程空转 busy-wait 浪费 CPUreceiveMessageSystemAttributeNames空列表接收时请求的系统属性名与配置不一致的请求将绕过 BatchManager 直连 SQSreceiveMessageAttributeNames空列表接收时请求的消息属性名同样不一致的请求绕过批处理SqsAsyncBatchManager支持的操作面与最初设计只聚焦sendMessage/flush不同当前实现已覆盖出站与入站四个方向见 SqsAsyncBatchManager.javasendMessage(...)→ 组装SendMessageBatchRequestdeleteMessage(...)→ 组装DeleteMessageBatchRequestchangeMessageVisibility(...)→ 组装ChangeMessageVisibilityBatchRequestreceiveMessage(...)→ 缓冲并批量拉取消息每次最多 10 条缓冲区无消息时返回空消息。每个方法都提供直接传请求对象与传ConsumerBuilder配置式两种重载。DefaultSqsAsyncBatchManager为每个方向维护独立的RequestBatchManager实例SendMessageBatchManager、DeleteMessageBatchManager、ChangeMessageVisibilityBatchManager与ReceiveMessageBatchManager其中发送方向额外启用了字节上限MAX_SEND_MESSAGE_PAYLOAD_SIZE_BYTES见 DefaultSqsAsyncBatchManager.java。通用性验证SampleBatchManager决策该方案可推广到任何受益的服务由测试中的SampleBatchManager佐证它把RequestBatchManagerString, String, BatchResponse泛型实例化为字符串请求的演示实现仅需实现batchAndSend、getBatchKey、mapBatchResponse三个方法见 SampleBatchManager.java。这说明任何服务都能通过继承该抽象骨架获得自动批处理能力与 6/29 决策 2跨 SDK 推广的方向一致。延伸阅读决策日志原文DecisionLog.md配套设计文档Design.md公开接口与配置SqsAsyncBatchManager.java、BatchOverrideConfiguration.java核心实现RequestBatchManager.java、SendMessageBatchManager.java、DefaultSqsAsyncBatchManager.java测试与集成验证SqsAsyncBatchManagerTest.java、SampleBatchManager.java、RequestBatchManagerSqsIntegrationTest.java总结从 2021 年 6 月 29 日到 7 月 15 日的三次设计讨论AWS SDK for Java V2 的批处理工具完成了从命名到行为边界的全部关键决策独立工具类对齐 Waiters 风格、{Service}BatchManager命名、按队列 URL 与 override 配置哈希分组、名称匹配的请求映射、字段不一致时抛异常、1:1 的请求/响应关联、暂不支持流式输入、异步客户端并发受限于连接数而非串行、重试交由低层客户端处理。这些决策在今天的services/sqs模块中全部有迹可循——公开接口、配置默认值、缓冲刷盘调用链与测试用例共同构成一个既尊重 v1AmazonSQSBufferedAsyncClient遗产、又面向多服务通用的自动批处理架构。【免费下载链接】aws-sdk-java-v2The official AWS SDK for Java - Version 2项目地址: https://gitcode.com/GitHub_Trending/aw/aws-sdk-java-v2创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表