本文围绕 Hazelcast Jet(现已并入 Hazelcast 统一实时数据平台)自 4.1 版本引入的代码部署增强(Code Deployment Improvements)展开:作业提交时如何把内部类、匿名类乃至整个包递归地注入到 Jet 作业的 classpath 中。文章将以设计文档 docs/design/jet/001-code-deployment-improvements.md 为骨架,并结合仓库源码与测试深入讲解JobConfig.addClass()/addPackage()的底层实现、classgraph 库的着色使用方式、资源在集群中的存储与内存开销,以及可验证的测试行为。读完本文,你将能够正确地为 Jet 作业配置类与包级资源,理解“资源以 IMap 存储、主副本 + 备份副本双份驻留内存”的容量规划要点,并能针对常见坑点(如 light job 限制、非现有包名行为)写出可运行的部署代码。
背景:为什么需要“嵌套类”与“包级”部署
在 Hazelcast Jet 中,一个作业(Job)由 Pipeline 或 DAG 描述,作业的代码(类、JAR、classpath 资源)通过JobConfig提交到集群,由 Jet 的专用 classloader 在成员节点上加载执行。早期JobConfig只支持添加单个类、单个 JAR 或单个 classpath 资源,用户面临两个典型的部署痛点:
- 嵌套类无法整体携带:一个业务类往往附带内部类(inner class)与匿名类(anonymous class,例如 Kotlin lambda、
Runnable匿名实现、EntryProcessor 中常见的匿名内部类)。手工逐个addClass()不现实,且容易遗漏。 - 包级资源无法批量携带:用户希望一次把某个包(package)下的所有类与资源一次性注入作业 classpath,而 Java 标准 API 并不支持“列出某个包下的所有类”。
设计文档把目标概括为:
User should be able to add nested (inner & anonymous) classes as well as entire packages to Jet job's classpath.
问题陈述:Java 反射与 ClassLoader 的能力边界
为什么这个需求在技术上并不平凡?设计文档明确指出两个硬性约束:
- Java 反射无法列出匿名类:给定一个类,标准反射 API 只能通过
getDeclaredClasses()拿到声明在源码中的成员类,却无法获知该类在方法体内创建的匿名类(如new Runnable() {...}、Kotlin lambda 生成的Foo$1之类的类)。这些匿名类同样需要被 JVM 加载才能反序列化执行。 - ClassLoader 不支持列出包内容:
ClassLoader.getResources()只能按确定的资源名查找,无法回答“这个包下有哪些类文件”这类目录列举问题。
换句话说,仅靠Class.forName+ 反射,无法穷举一个类“实际依赖的所有字节码文件”。
解决方案:借助 classgraph 扫描包资源
设计文档给出的解决思路非常直接:既然嵌套类与根类位于同一个包(package)下,那么可以列出该包关联的所有资源 URL,再提取并过滤出想要的 class 文件——这一策略同时要处理“列目录、从任意嵌套 JAR 中解压、处理自定义 URL scheme”等场景;同样的策略也可递归用于找出一个包内的全部类。
与其从零实现一套 classpath 扫描器,Jet 选择引入一个轻量级开源库 ——classgraph(文档记载体积约 470KB,Jet 通过 Maven Shade 将其着色打包进自己的发行物中,避免与用户 classpath 冲突)。classgraph 支持广泛的 classpath 规格机制,可以:
- 扫描目录、嵌套 JAR、
file:/jar:等常见 URL scheme; - 提供包级(package)与路径级(path)的接受规则;
- 返回类信息(
ClassInfo)及其内部类关系(getInnerClasses()); - 返回非 class 资源(如包内的
package.properties、静态资源文件)。
在仓库源码中可以看到 Jet 对 classgraph 的实际使用,见 ReflectionUtils.java(io.github.classgraph.ClassGraph、ClassInfo、ScanResult的导入与调用)。
API 变更:JobConfig的两个新能力
设计文档为JobConfig提出两项 API 扩展,两者在 JobConfig.java 中均有完整实现。
1. 增强addClass():递归添加嵌套类
原文档给出的签名(since 4.1):
@Nonnull @SuppressWarnings("rawtypes") public JobConfig addClass(@Nonnull Class... classes)仓库中的实际实现(JobConfig.java):
@Nonnull public JobConfig addClass(@Nonnull Class<?>... classes) { throwIfLocked(); checkNotNull(classes, "Classes cannot be null"); ResourceConfig.fromClass(classes).forEach(cfg -> resourceConfigs.put(cfg.getId(), cfg)); return this; }其语义(Javadoc 原文要点):
- 递归添加给定类及其所有嵌套(内部与匿名)类到 Jet 作业 classpath;
- 这些类只对挂接在该 Pipeline / DAG 上的代码可见,对其他代码不可见——重要示例是
IMap数据源,它只能实例化 Jet 实例 classpath 中的类(即集群侧代码); - 不能用于 light job(
JetService#newLightJob(Pipeline)); - 底层存储是默认备份数为 1 的
IMap,因此添加大文件时需按内存扩容集群:每个文件在集群内会有2 份拷贝(主副本 + 备份副本)。
注意:类被加入后,返回
this以支持 Fluent API 链式调用;addClass与addPackage调用前都会先throwIfLocked()——JobConfig一旦被提交(locked)便不可再修改。
2. 新增addPackage():递归添加整个包
原文档给出的签名:
@Nonnull public JobConfig addPackage(@Nonnull String... packages)仓库中的实际实现(JobConfig.java):
@Nonnull public JobConfig addPackage(@Nonnull String... packages) { checkNotNull(packages, "Packages cannot be null"); Resources resources = ReflectionUtils.resourcesOf(packages); resources.classes().forEach(classResource -> add(classResource.getUrl(), classResource.getId(), CLASS)); resources.nonClasses().forEach(this::addClasspathResource); return this; }语义要点:
- 递归添加给定包内的全部类与资源(非 class 文件,如
package.properties)到 Jet 作业 classpath; - 包内每个 class 文件以
ResourceType.CLASS注册,资源 ID 即“包路径/类名.class”;非类资源则以addClasspathResource方式注册; - 同样不能用于 light job,同样受 IMap 双副本内存开销约束;
- 与
addClass不同,addPackage对传入的包名数组做了checkNotNull但不调用throwIfLocked()——从源码结构看,包内资源是在扫描完成后逐条通过私有add(...)注册的。
配套 API 一览
设计文档提到可参考的addJar与addClasspathResource,仓库中 JobConfig.java 提供了一系列重载:
| 方法 | 入参形式 | 资源 ID 规则 | 说明 |
|---|---|---|---|
addJar(URL) | URL | URL 最后一段路径(filename) | 若 ID 已存在则不替换 |
addJar(File) | File | 文件名 | 同上;要求是文件(ensureIsFile) |
addJar(String) | 路径字符串 | 文件名 | 内部转File再转 URL |
addJarsInZip(URL/File/String) | ZIP 内 JAR | 按 JAR 文件名 | 一次添加 ZIP 包内多个 JAR |
addClasspathResource(URL/File/String) | 单资源 | URL/文件名 | 另有带显式id的重载 |
addClass(Class...) | 类对象 | 包路径/类名.class | 递归嵌套类 |
addPackage(String...) | 包名数组 | 各 class/资源的路径 | 递归包内类与资源 |
所有资源的实际存储形态为ResourceConfig(ResourceConfig.java),它是一个实现IdentifiedDataSerializable的配置对象,携带URL、可空id与ResourceType(CLASS/JAR/CLASSPATH_RESOURCE)。
源码级原理:ReflectionUtils如何用 classgraph 干活
addClass与addPackage的核心逻辑都汇聚在工具类ReflectionUtils(ReflectionUtils.java)。
nestedClassesOf():找出一个类及其全部嵌套类
ResourceConfig.fromClass()的实现(ResourceConfig.java):
public static Stream<ResourceConfig> fromClass(@Nonnull Class<?>... classes) { return ReflectionUtils.nestedClassesOf(classes).stream().map(ResourceConfig::new); }nestedClassesOf()(ReflectionUtils.java)的关键步骤:
- 构造
ClassGraph并开启.enableClassInfo()与.ignoreClassVisibility()——后者保证即使嵌套类是私有的也能被发现; - 收集传入类的classloader(去重后逐个
addClassLoader),并把每个类所在的包名作为acceptPackages规则传入; scan()后先按类名过滤出目标ClassInfo,再对其调用getInnerClasses()取出所有内部/匿名类信息,最后loadClass()得到真实的Class<?>对象;- 返回
传入类本身 + 全部嵌套类的合并集合。
这一步正是设计文档中“先列出包资源、再过滤目标类文件”策略的实现:由于嵌套类与根类同包,acceptPackages加上getInnerClasses()的关系遍历即可补全反射无法枚举的匿名类。
resourcesOf():递归扫描包内的类与资源
addPackage()调用的resourcesOf(String... packages)(ReflectionUtils.java):
- 把包名
com.example.foo转换成路径com/example/foo(.replace('.', '/')); ClassGraph同时acceptPackages(packages)与acceptPaths(paths),ignoreClassVisibility();scan()后getAllClasses()得到包内全部类,包装成ClassResource集合;getAllResources().nonClassFilesOnly()得到包内非 class 资源 URL 集合;- 打包返回
Resources(含 classes 与 nonClasses 两部分)。
回到JobConfig.addPackage():classes 以CLASS类型逐条add(...),nonClasses 交给addClasspathResource,从而实现对“包内一切内容”的完整携带。
测试验证:从单测看行为契约
设计文档特别提到:“In particular created unit test proving Kotlin lambdas (as anonymous classes) are added to the classpath”,即专门有单测证明 Kotlin lambda(本质是匿名类)能被加入 classpath。仓库中的ResourceConfigTest(ResourceConfigTest.java)用四类用例锁定了行为:
when_addClassWithClass(第 62-71 行):addClass(ResourceConfigTest.class)后,资源 ID 等于ReflectionUtils.toClassResourceId(类)(即com/hazelcast/jet/config/ResourceConfigTest.class),资源类型为ResourceType.CLASS;when_addResourcesWithPackage(第 73-92 行):addPackage(本类所在包名)后,资源配置集合中既包含本类的 CLASS 资源,也包含package.properties这类 CLASSPATH_RESOURCE——验证“包内类 + 资源一并携带”;when_addResourcesWithNonExistingPackage(第 94-103 行):对不存在的包名thispackage.does.not.exist调用addPackage返回空集合而不抛异常(空操作语义);- JAR 相关用例(第 105-159 行):
addJar(URL)以 URL 末段为资源 ID;URL 无路径段时抛IllegalArgumentException;相同资源 ID 重复添加抛异常(不静默覆盖)。
此外,JobConfigTest.java 覆盖了addClass/addPackage在JobConfig层面的注册行为;JetTest.java 展示了new JobConfig().addClass(JetTest.class)的典型用法;部署相关测试(如 AbstractDeploymentTest.java、ClientDeployment_StandaloneClusterTest.java)则验证了带用户代码部署配置的客户端如何通过addClass把类随作业提交到集群。
典型用法示例
把上述 API 组合起来,一个典型的 Jet 作业配置如下(Fluent API 风格):
JobConfig jobConfig = new JobConfig() // 递归携带该类的所有内部类与匿名类(含 Kotlin lambda 生成的类) .addClass(MyEntryProcessor.class, MyFunction.class) // 递归携带整个包内的类与资源 .addPackage("com.example.jet.functions", "com.example.jet.mappers") // 携带独立 JAR 与任意 classpath 资源 .addJar("/path/to/lib-extra.jar") .addClasspathResource("/path/to/config.properties"); Pipeline p = Pipeline.create(); p.readFrom(Sources.list("input")) .map(new MyFunction()) // 引用 addClass 加入的类 .writeTo(Sinks.list("output")); JetInstance jet = Jet.newJetClient(); jet.newJob(p, jobConfig).join();使用须知(来自 Javadoc 与源码行为):
- 作用域:这些资源只对附着于本作业 Pipeline / DAG 的代码可见;
IMap数据源等集群侧组件无法使用它们; - light job 限制:
addClass/addPackage/addJar/addClasspathResource均不可用于 light job; - 内存规划:资源以
IMap存储、默认备份数为 1,每个资源在集群中驻留2 份(主副本 + 备份副本),提交大文件前需按此评估成员内存; - ID 唯一性:JAR 类资源以文件名(URL 末段)为 ID,重复 ID 会抛
IllegalArgumentException;addClass/addPackage生成的资源 ID 为包路径形式,重复添加同一类/包会覆盖同 ID 条目(源码中resourceConfigs.put(cfg.getId(), cfg)); - 空包行为:
addPackage对不存在的包名表现为空操作(不抛异常,见测试用例)。
局限与演进方向
设计文档在末尾列出了一项“未来改进”(Improvements):
Maybe instead of 'forcing' user to manually add (
addClass(),addPackage(),addJar()…) all the required resources to the classpath we could scan the classpath and add them automatically - however we would need to figure out the way to filter out unneeded ones so we don’t end up with bloated job’sIMapstate.
即:理想方案是自动扫描 classpath 并推断作业所需资源,免去用户手工addClass/addPackage/addJar的负担;难点在于如何过滤无关类,避免作业的IMap状态被无用资源撑爆——这与本文反复强调的内存双副本约束直接相关,也是自动部署方案必须解决的核心权衡。
从当前仓库源码结构看,addClass的递归嵌套类能力与addPackage的包级扫描能力已经完整落地,自动扫描仍属于文档层面的设计展望而非已实现功能。
总结
- 能力:自 4.1 起,
JobConfig.addClass()递归携带嵌套类(含匿名类与 Kotlin lambda),新增的addPackage()递归携带包内全部类与资源; - 原理:Java 反射无法枚举匿名类、ClassLoader 无法列出包内容,Jet 借助轻量级开源库 classgraph 扫描包级资源并过滤目标类,相关逻辑集中在 ReflectionUtils.java;
- 实现落点:JobConfig.java 的
addClass/addPackage、ResourceConfig.java 的fromClass; - 行为契约:由 ResourceConfigTest.java 等测试锁定(类/包/资源 ID 规则、空包空操作、JAR ID 冲突异常);
- 运维要点:资源经
IMap双副本存储,需按内存扩容;light job 不可用;作用域仅限作业自身 Pipeline / DAG。
在设计文档 001-code-deployment-improvements.md 之外,读者可继续深入 JobConfig.java、ReflectionUtils.java 与 ResourceConfigTest.java 获取一手实现与验证证据。
- 缓存
- KV存储
- 消息队列
- 流处理
- 后端
【免费下载链接】hazelcast
Hazelcast 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
相关推荐
使用 Helm Chart 在 Kubernetes 上部署 Hazelcast Jet 流处理集群的完整指南
使用 Helm Chart 在 Kubernetes 上部署 Hazelcast Jet 流处理集群的完整指南 Hazelcast Jet 是一个构建在 Haz
不只是 Loading!easy-loading-cj Toast 功能实战:图标、偏移、尺寸一次讲透
不只是 Loading!easy loading cj Toast 功能实战:图标、偏移、尺寸一次讲透 easy loading cj 是一款面向 Cangji
缓存KV存储消息队列流处理后端DeepMosaics 预训练模型完全指南:模型分类、下载部署与源码级调用机制解析
DeepMosaics 预训练模型完全指南:模型分类、下载部署与源码级调用机制解析 本篇指南以 DeepMosaics 仓库官方文档 docs/pre trai
人工智能深度学习计算机视觉图像处理视频处理