
缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载导读本篇文章以 Hazelcast 仓库中 Jet 作业上传设计文档 为核心完整讲解 Hazelcast Jet 如何通过**客户端二进制协议Client Binary Protocol**让非 Java 客户端例如 Python、Node.js、Go、C# 等任何实现了该协议的语言向集群成员上传 Jar 并启动 Jet 作业以及如何让成员直接执行已存在于本地的 Jar。读完本文你将掌握SubmitJobParameters参数对象与JetClientInstanceImpl.submitJobFromJar私有 API 的用法、uploadJobMetaData与uploadJobMultipart两个协议消息的完整字段语义、服务端分片校验与临时文件管理机制以及分片大小等关键配置项的调优方法。背景与动机现状上传作业必须依赖 hz-cli 与 JVM在设计该功能之前Jet 作业的 Jar 上传只能通过hz-cli命令完成而hz-cli是运行在主机上的 Java 工具这意味着任何想要提交 Jar 的客户端都必须在本地安装 JVM。这对于非 Java 生态Python、Go、Node.js 等的开发者而言是明显阻碍。该设计文档状态为 DRAFT关联 JiraHZOLD-1455明确了新功能的两大目的让非 Java 客户端能够向集群成员上传 Jar并附带参数执行它——上传流程经由客户端二进制协议完成客户端不再需要本地 JVM让集群成员直接执行一份已存在的 Jar——成员执行本地 Jar同样支持传入参数。术语与参与者术语定义Client Binary ProtocolHazelcast 客户端与 Hazelcast Server 之间使用的二进制消息协议该功能面向的参与者Actors是任何实现了 Client Binary Protocol 的语言客户端。功能设计概览新功能为用户提供两类核心能力上传并执行上传一份位于客户端本地的 Jar 文件并传入与hz-cli相同语义的参数作业名、快照名、主类、作业参数等仅执行直接执行一份已位于集群成员上的 Jar 文件同样传入与hz-cli相同语义的参数。设计约束与注意事项设计文档同时给出了两类场景必须遵守的约束Jar On ClientJar 在客户端场景作业运行期间访问的所有资源都必须对服务端可用。例如用于填充IMap的文本文件必须位于服务端或被打包进上传的 Jar 中——因为服务端作业无法访问客户端侧资源客户端必须按顺序发送多个分片multi parts。乱序消息不做处理因为服务端为此需要占用更多资源。Jar On MemberJar 在成员上场景被执行的 Jar 文件必须能够通过file://URL 方案访问。使用方式基于 JetClientInstanceImpl 的私有 API该功能使用的是私有 API因此客户端需要把JetService强转为JetClientInstanceImpl才能调用。Jar 上传Jar On ClientHazelcastInstance client ...; JetClientInstanceImpl jetService (JetClientInstanceImpl) client.getJet(); SubmitJobParameters submitJobParameters SubmitJobParameters.withJarOnClient() .setJarPath(Paths.get(/path/to/my-job.jar)) .setJobName(my-job) .setMainClass(org.example.Main) .setJobParameters(List.of(--param1, value1)); // 如果出现错误抛出 JetException jetService.submitJobFromJar(submitJobParameters);Jar 执行Jar On MemberHazelcastInstance client ...; JetClientInstanceImpl jetService (JetClientInstanceImpl) client.getJet(); SubmitJobParameters submitJobParameters SubmitJobParameters.withJarOnMember() .setJarPath(Paths.get(/data/jobs/existing-job.jar)) .setSnapshotName(my-snapshot) .setJobName(existing-job); // 如果出现错误抛出 JetException jetService.submitJobFromJar(submitJobParameters);客户端新增方法设计文档明确在JetClientInstanceImpl接口上新增了一个方法void submitJobFromJar(Nonnull SubmitJobParameters submitJobParameters);从 JetClientInstanceImpl.java 的实现可以看到该方法是一个分发入口当submitJobParameters.isJarOnMember()为真时走executeJobFromJar仅发送元数据否则走uploadJobFromJar先发送元数据再发送全部分片。两条路径在发送前都会调用SubmitJobParametersValidator做客户端侧预校验并通过invokeRequestNoRetryOnRandom将消息固定发给会话所选定的同一个成员disallowRetryOnRandom从而保证整个上传会话与单个成员绑定。参数对象 SubmitJobParameters 字段说明参数对象定义在 SubmitJobParameters.java采用链式 setter 的不可变风格通过私有构造器 静态工厂方法创建字段说明jarPathJar 文件的路径java.nio.file.PathjarOnMember是否为成员本地 Jar。withJarOnClient()默认为falsewithJarOnMember()通过setJarOnMember()置为truesnapshotName启动作业时传入的快照名可空jobName启动作业时传入的作业名可空mainClass启动作业时使用的主类规范名如org.example.Main可空为空时要求 Jar 的 manifest 中提供Main-ClassjobParameters传给作业的参数列表默认空列表客户端侧参数校验客户端发送前会先做本地校验逻辑集中在 SubmitJobParametersValidator.javaJar On Client要求jarOnMember为false、jarPath非空、文件扩展名必须以.jar结尾、文件大小不能为 0、jobParameters不能为nullJar On Member要求jarOnMember为true、jarPath非空、扩展名为.jar、jobParameters非空不校验文件大小因为文件在成员本地。任何一条不满足都会抛出JetException。技术设计客户端二进制协议的两个新消息为了让非 Java 客户端也能上传并执行 Jet 作业客户端协议新增了两条消息uploadJobMetaDatauploadJobMultipart两条消息都应发送给集群中的某个成员从实现看是整个会话固定到同一成员。消息一uploadJobMetaData这条消息同时服务于 Jar 上传与 Jar 执行两种场景。Jar On Member 场景该场景只执行 Jar只使用uploadJobMetaData一条消息。消息字段如下字段类型定义sessionIdUUID非空会话 UUIDjarOnMemberBoolean必须为true标志位表示要执行的 Jar 已存在于成员上客户端不会上传 JarfilenameString非空Jar 文件的完整路径sha256HexString非空但被忽略传空字符串即可Jar 文件的十六进制 SHA256snapshotNameString可空启动作业时传入的参数jobNameString可空启动作业时传入的参数mainClassString可空启动作业时传入的参数若为nullJar 的 manifest 中应包含Main-Class值jobParametersListString非空无参数时传空列表启动作业时传入的参数Jar On Client 场景上传流程以uploadJobMetaData开始。消息字段如下字段类型定义sessionIdUUID非空会话 UUID用于关联本次会话中的所有消息jarOnMemberBoolean必须为false标志位表示 Jar 不在成员上需要从客户端上传filenameString非空不带扩展名的 Jar 文件名sha256HexString非空Jar 文件的十六进制 SHA256snapshotNameString可空启动作业时传入的参数jobNameString可空启动作业时传入的参数mainClassString可空启动作业时传入的参数若为nullJar 的 manifest 中应包含Main-Class值jobParametersListString非空无参数时传空列表启动作业时传入的参数与客户端实现对照上传场景由 JetClientInstanceImpl.java 的sendJobMetaDataForUpload编码请求其中filename使用getFileNameWithoutExtension()去扩展名sha256Hex为整个 Jar 的 SHA256执行场景的sendJobMetaDataForExecute则直接使用 Jar 的完整路径作为filename。云环境特别说明上传的文件会存放在临时目录中因此Pod 需要可写的文件系统——例如在部署描述符中通过emptyDir {}提供临时目录路径由TMP环境变量或java.io.tmpdir属性控制。服务端收到元数据后的处理收到uploadJobMetaData后服务端先执行校验任何校验规则失败都会抛出JetException。校验通过后服务端会在JobUploadStore类中为本次会话存入一条新记录。从 JobUploadStore.java 的实现看JobUploadStore内部以sessionId为 key、用ConcurrentHashMapUUID, JobUploadStatus保存进行中的上传状态processJobMetaData会先检查sessionId是否已存在——若已存在则抛出JetException(Session already exists...)防止同一会话重复初始化随后创建JobUploadStatus并调用createNewTemporaryFile()创建临时文件。消息二uploadJobMultipart上传流程继续使用这条消息发送 Jar 的字节内容。消息字段如下字段类型定义sessionIdUUID与上一条消息含义相同标识会话currentPartNumberint从 1 开始表示当前分片序号例如 5 个分片中的第 1 个totalPartNumberint序列的分片总数partDatabyteArray包含 Jar 数据的byte[]partSizeint指明partData中有效的字节数sha256HexString该分片的十六进制 SHA256为什么需要额外的 partSize 字段出于优化考虑假设客户端侧partData缓冲区只分配一次复用同一个数组因此需要额外的字段来指示本次消息中应从该缓冲区读取多少字节。从服务端实现 JobUploadStatus.java 可以看到追加写入时使用的是outputStream.write(partData, 0, partSize)——只有前partSize字节是有效数据。服务端分片校验规则收到uploadJobMultipart后服务端会对消息与会话状态做多种校验任何一条失败都会抛出JetException。设计文档列出的检查项包括校验partData.length为正数校验partSize字段为正数校验partSize partData.length校验currentPart不小于已收到的receivedCurrentPart——即拒绝来自过去的分片校验currentPart 1必须等于receivedCurrentPart——即拒绝来自未来的分片这也意味着重复消息会被拒绝上传操作将失败校验totalPart不为 0 且等于已收到的receivedTotalPart——分片总数中途不允许变化校验校验和checksum。从实际实现 JobUploadStatus.java 可以进一步确认这些规则的细节对单条分片currentPartNumber与totalPartNumber必须为正数、currentPartNumber不能大于totalPartNumber、partData不能为null且长度不能为 0、partSize必须为正数且不能大于partData.length写入时只取前partSize字节、分片sha256Hex不能为null对分片顺序已处理currentPart之后只接受currentPart 1接受前会先校验收到的分片 SHA256 与服务端重新计算的 SHA256 一致对完整性当currentPart totalPart时视为全部完成此时还会对整个 Jar 文件再做一次 SHA256 校验validateJarChecksum与元数据消息中的sha256Hex比对防止分片级校验被绕过。组装与执行临时文件 同 JVM 启动收到第一个分片消息时服务端会创建新的临时文件createJarPath 使用Files.createTempFile若元数据中带上传目录则在该目录下创建否则使用系统默认临时目录同时把文件权限设置为仅属主可读、可写、可执行之后每一条分片消息都把partData中前partSize字节追加CREATEAPPEND到该文件当所有分片全部完成时服务端使用HazelcastBootstrap类启动一个新作业——该类在与成员相同的 JVM内执行 Jar。同 JVM 执行的好处保持了从 hz-client 提交作业时在同一个 JVM 中运行的总体做法可以复用已有的JetService与集群资源效率更高一个 JVM 内可以同时运行多个实例。同 JVM 执行的代价作业内部的任何失败都会直接影响该成员。从 HazelcastBootstrap.java 可以看到它的演进痕迹它原本只为hz-client 命令设计commandLineExecuteJar会设置单例、执行完关闭实例并重置单例本特性将其改造为可在服务端使用memberExecuteJar它仍然是单例并且成员端执行时不允许关闭单例否则成员会被关停。失败清理与会话过期如果服务端抛出异常上传操作即告失败。此时JetServiceBackend中的定时器会负责清理过期的JobUploadStore条目并删除临时文件。对应实现JetServiceBackend中以 30 秒为周期调度jobUploadStore::cleanExpiredUploads见 JetServiceBackend.java 附近周期常量JOB_UPLOAD_STORE_PERIOD 30单位为秒JobUploadStatus中每个条目的过期时间窗口为EXPIRATION_MINUTES 2分钟任何一次分片处理都会刷新lastUpdatedTime以推迟过期见 JobUploadStatus.java。值得注意的运维细节如果上传会话在服务端已因不活跃而超时processJobMultipart找不到 session 时会抛出JetExceptionJobUploadStore 给出的提示信息 会建议网络较慢的客户端改用更小的分片方法有两种设置环境变量hazelcast.client.jobupload.partsize或在ClientConfig中设置ClientProperty.JOB_UPLOAD_PART_SIZE。分片大小调优uploadJobMultipart默认分配10,000,000 字节约 10 MB的缓冲区。缓冲区大小由ClientProperty.JOB_UPLOAD_PART_SIZE属性控制属性名为hazelcast.client.jobupload.partsize默认值10_000_000见 ClientProperty.java。调优含义希望客户端分配更少内存的客户端可以通过调小分片大小来换取更多数量的分片消息即用更多网络往返换取更低的内存占用。反之在带宽充足、网络稳定的内网环境保持默认值或调大可减少分片数量、降低校验与追加开销而在高延迟或慢速网络上调小分片有助于降低单次传输失败与会话超时的概率与上文JobUploadStore超时提示中的建议一致。测试标准设计文档对本次功能要求的测试覆盖如下单元测试在类级别测试各功能逻辑压力测试并行上传多个客户端并行上传 Jar压力测试并行执行多个成员/会话并行执行 Jar。仓库中已落地的客户端侧测试覆盖了成功与失败两条路径可直接作为实现参考JobUploadClientSuccessTest.java 与 JobUploadClientFailureTest.java——覆盖上传场景JobExecuteClientSuccessTest.java 与 JobExecuteClientFailureTest.java——覆盖成员本地执行场景。仓库实现路径速查设计文档docs/design/jet/023-jobupload.md参数对象hazelcast/src/main/java/com/hazelcast/jet/impl/SubmitJobParameters.java客户端入口与协议编码hazelcast/src/main/java/com/hazelcast/jet/impl/JetClientInstanceImpl.java客户端参数校验hazelcast/src/main/java/com/hazelcast/jet/impl/submitjob/clientside/validator/SubmitJobParametersValidator.java成员端上传会话状态hazelcast/src/main/java/com/hazelcast/jet/impl/submitjob/memberside/JobUploadStore.java 与 JobUploadStatus.java同 JVM 执行引导类hazelcast/src/main/java/com/hazelcast/instance/impl/HazelcastBootstrap.java分片大小属性hazelcast/src/main/java/com/hazelcast/client/properties/ClientProperty.java赞分享缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载相关推荐如何快速部署 microservices-demo5分钟搭建微服务演示环境如何快速部署 microservices demo5分钟搭建微服务演示环境 microservices demo 是一个微服务Demo项目展示了微服务架构如从0到1掌握情感分析bert-base-cased-finetuned-sst2完整使用教程从0到1掌握情感分析bert base cased finetuned sst2完整使用教程 bert base cased finetuned sst2是一netfox自定义扩展开发如何为特定需求定制网络调试功能 netfox自定义扩展开发如何为特定需求定制网络调试功能 netfox是一款轻量级、一行代码即可集成的iOS/OSX网络调试库能够帮助开发者轻松捕获和上一篇wgpu Mesh Shader 实战指南基于任务着色器、Payload 与逐图元数据渲染三角形下一篇ctf-wiki 流量包分析指南CTF Misc 中的 PCAP 取證分析方法論與實戰创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考