拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Apache DolphinScheduler 的 dolphinscheduler-yarn-aop:用 AspectJ 捕获 YARN ApplicationId 的轻量织入模块

Apache DolphinScheduler 的 dolphinscheduler-yarn-aop用 AspectJ 捕获 YARN ApplicationId 的轻量织入模块【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler本篇技术指南聚焦 Apache DolphinScheduler 仓库中的dolphinscheduler-yarn-aop模块深入讲解它如何通过 AspectJ 在编译期织入 Hadoop YARN 的YarnClientImpl.submitApplication将每个提交的 YARN 应用的ApplicationId捕获并写入appInfo.log。读完本文你将掌握该模块的切面实现细节、编译期/运行期两种织入方式的配置方法以及 Worker 端如何基于appInfo.log完成应用 ID 的追踪与任务关联。模块定位解决“YARN 任务丢了 ApplicationId”的痛点在 DolphinScheduler 中MRMapReduce、Spark-on-YARN、Sqoop 等任务最终都是通过yarn jar ...方式把 Jar 提交到 YARN 集群。任务提交后Worker 或运维人员需要拿到 YARN 分配的ApplicationId才能完成追踪任务在 YARN 上的运行状态通过getApplicationReport在任务失败或需要取消时按ApplicationId精确 kill 对应的 YARN 作业将远程应用的 ID 与 DolphinScheduler 的任务实例task instance关联回传给 Master。但通过外部进程方式提交 YARN 任务时ApplicationId往往只出现在任务日志里解析并不稳定。dolphinscheduler-yarn-aop正是为此而生——它是仓库中的一个微型 AspectJ 模块直接在YarnClientImpl.submitApplication返回的瞬间截获ApplicationId并将其追加写入工作目录下的appInfo.log为下游跟踪/终止 YARN 作业提供可靠依据。该模块的官方定位在 dolphinscheduler-yarn-aop/CLAUDE.md 中描述为 “Tiny AspectJ module”其pom.xml的描述也印证了用途aop 4 YarnClient to get application id when submitting jars using yarn jar mainClass args。模块结构与依赖关系模块的源码组织非常精简主包org.apache.dolphinscheduler.aop仅包含一个切面类 YarnClientAspect.java测试包org.apache.dolphinscheduler.pocYarnClientMoc、YarnClientAspectMoc与org.apache.dolphinscheduler.YarnClientAspectMocTest。从 dolphinscheduler-yarn-aop/pom.xml 可以看出其依赖设计org.aspectj:aspectjweaver与org.aspectj:aspectjrtAspectJ 的织入器与运行时org.apache.hadoop:hadoop-yarn-client与org.apache.hadoop:hadoop-common提供YarnClientImpl、ApplicationId、ApplicationSubmissionContext、ApplicationReport等 YARN 类型通过dolphinscheduler-bom统一管理依赖版本aspectj.version当前固定为 1.9.7见 dolphinscheduler-bom/pom.xml。在构建层面模块作为子模块登记在根 pom.xml 中并在根 POM 的pluginManagement里统一配置了aspectj-maven-plugin版本 1.14.0设置了complianceLevel/source/target对齐 Java 版本、开启showWeaveInfo与verbose见 pom.xml子模块无需重复声明插件版本。YarnClientAspect 源码级解析YarnClientAspect是整个模块的核心声明为Aspect只包含两个AfterReturning通知切面表面刻意保持最小化。通知一捕获 submitApplication 的返回值AfterReturning(pointcut execution(ApplicationId org.apache.hadoop.yarn.client.api.impl.YarnClientImpl. submitApplication(ApplicationSubmissionContext)) args(appContext), returning submittedAppId, argNames appContext,submittedAppId) public void registerApplicationInfo(ApplicationSubmissionContext appContext, ApplicationId submittedAppId) { try { Files.write(Paths.get(appInfoFilePath), Collections.singletonList(submittedAppId.toString()), StandardOpenOption.CREATE, StandardOpenOption.WRITE, StandardOpenOption.APPEND); } catch (IOException ioException) { logger.error( YarnClientAspect[registerAppInfo]: cant output current application information, because {}, ioException.getMessage()); } logger.info(YarnClientAspect[submitApplication]: current application context {}, appContext); logger.info(YarnClientAspect[submitApplication]: submitted application id {}, submittedAppId); logger.info( YarnClientAspect[submitApplication]: current application report {}, currentApplicationReport); }要点对应源码 YarnClientAspect.java切点精确匹配YarnClientImpl.submitApplication(ApplicationSubmissionContext)方法并用args(appContext)绑定入参返回值通过returning submittedAppId绑定通知逻辑将ApplicationId以追加StandardOpenOption.APPEND方式写入appInfoFilePath文件不存在时自动创建StandardOpenOption.CREATE。由于是追加写一个任务多次提交或多个应用提交都会完整保留记录写入失败不会抛出导致任务失败而是记录 error 日志体现了“尽量不影响主流程”的降级设计同时以logger.info记录应用上下文、提交的应用 ID 与当前应用报告便于排查。appInfoFilePath在构造函数中确定为System.getProperty(user.dir) /appInfo.log见 YarnClientAspect.java这正是后续“Gotchas”中关于工作目录约定的由来。通知二cflow 限定的 getApplicationReport 捕获AfterReturning(pointcut cflow(execution(ApplicationId org.apache.hadoop.yarn.client.api.impl.YarnClientImpl.submitApplication(ApplicationSubmissionContext))) !within(YarnClientAspect) execution(ApplicationReport org.apache.hadoop.yarn.client.api.impl.YarnClientImpl.getApplicationReport(ApplicationId)), returning appReport, argNames appReport) public void registerApplicationReport(ApplicationReport appReport) { currentApplicationReport appReport; }要点对应 YarnClientAspect.javagetApplicationReport是 YARN 客户端在提交后查询应用状态的高频调用如果无差别织入会产生大量噪音这里用cflow(...)谓词限定仅当getApplicationReport发生在submitApplication的执行流control flow之内时才触发从而把捕获范围收敛到“提交过程中 YARN 内部自发的状态查询”再叠加!within(YarnClientAspect)排除切面自身调用避免递归每次命中都会把最新的ApplicationReport赋给currentApplicationReport字段——正如源码注释所述该方法可能被调用多次最后一次的 Report 实例会被保留供通知一记录日志时引用。织入方式编译期为主运行期为辅编译期织入CTW模块默认采用编译期织入由aspectj-maven-plugin在 Maven 生命周期内完成它同时绑定compile与test-compile两个 goal见根 pom.xml因此业务代码与测试代码都会被织入。产出是一个普通 jarpackagingjar/packaging。这个 jar 被打入 Worker 的 classpath在 dolphinscheduler-worker/pom.xml 中可以确认 Worker 直接依赖dolphinscheduler-yarn-aop。当 MR、Spark、Sqoop 等基于 YARN 的任务插件在 Worker 进程内直接使用YarnClientImpl时切面就会在编译期生效。运行期织入LTW对于第三方代码例如由 Spark 的 classloader 加载的 YARN 客户端编译期织入无法覆盖此时可以改用加载期织入load-time weavingjava -javaagent:aspectjweaver.jar ...通过-javaagent挂载aspectjweaver.jarJVM 在加载目标类时动态织入切面。原文档明确指出“operators may need to enable it depending on the task plugin”——具体哪些插件需要 LTW 取决于插件加载 YARN 客户端的 classloader 结构属于运维侧需要按任务插件实际行为验证的开关。appInfo.log 的消费链路从文件到任务关联appInfo.log只是切入点真正让“捕获的 ID 发挥作用”的是 Task API 侧的消费代码。两条取 ID 的路径LogUtils.java 的getAppIds(logPath, appInfoPath, fetchWay)提供了两种获取方式fetchWay aop直接从appInfoPath指向的文件逐行读取应用 IDgetAppIdsFromAppInfoFile见 LogUtils.java文件不存在或读取异常时返回空列表并记录警告/错误日志默认方式从任务日志文件logPath中用正则application_\d_\d匹配提取getAppIdsFromLogFile。该正则常量定义在 TaskConstants.java即YARN_APPLICATION_REGEX application_\\d_\\d对应 YARN 应用 ID 的标准格式application_startTime_id。appInfoPath本身由TaskExecutionContext携带字段定义见 TaskExecutionContext.java其默认文件名常量APPINFO_PATH appInfo.log与路径拼接String.format(%s/%s, execPath, APPINFO_PATH)定义在 FileUtils.java——注意这里默认把appInfo.log放在任务的执行目录execPath下与切面写入的“进程当前工作目录”需要在部署上保持一致。终止 YARN 作业时的使用在 ProcessUtils.java 中终止远程 YARN 作业的逻辑会取出taskExecutionContext.getAppInfoPath()调用LogUtils.getAppIds(...)解析应用 ID 列表再据此向对应节点发起 kill。这就是“Worker或运维人员可以追踪/终止与任务实例绑定的 YARN 作业”这条链路的具体落地。回传 Master任务侧通过AbstractRemoteTask在获取应用 ID 后调用taskCallBack.updateRemoteApplicationInfo(taskRequest.getTaskInstanceId(), new ApplicationInfo(getAppIds()))见 AbstractRemoteTask.java回调接口见 TaskCallBack.java将应用 ID 关联到任务实例并上报供前端监控页展示。测试验证用 AspectJ 织入的 Mock 模拟 YARN由于真实 YARN 集群不可得测试侧采用“用 AspectJ 语法 mock YARN”的策略三个关键类均在 src/test/java 下YarnClientMoc.java模仿真实YarnClientImpl的形态——提供一个createAppId()内部用ApplicationId.newInstance(System.currentTimeMillis(), random.nextInt())生成 ID和一个submitApplication(ApplicationSubmissionContext)方法返回createAppId()的结果模拟“提交即返回 ApplicationId”的行为YarnClientAspectMoc.java复刻真实切面的两个AfterReturning切点结构——一个绑定submitApplication的返回值另一个用cflow(...) !within(...)限定在提交流内捕获createAppId的返回值YarnClientAspectMocTest.javaJUnit 5 测试构造一个ApplicationSubmissionContext依次调用createAppId()与submitApplication(...)并把标准输出重定向到ByteArrayOutputStream最后断言输出中同时包含YarnClientAspectMoc[submitApplication]与YarnClientAspectMoc[createAppId]:以此证明“在 submitApplication 的调用流中两个通知都正确触发”。由于aspectj-maven-plugin绑定了test-compilegoal测试类本身同样被织入这正是“AspectJ-woven mock classes verify the aspect fires correctly”这句描述的由来。运维注意事项Gotchas原文档与源码共同指向以下三点实践约束部署与二次开发时务必遵守appInfo.log写在工作目录多 Worker 同主机同 cwd 会冲突。切面使用System.getProperty(user.dir)决定落盘位置YarnClientAspect.java且文件是追加写。若同一台主机上多个 Worker 进程共享同一个工作目录各任务实例的ApplicationId会混写在同一个文件中导致追踪与 kill 串扰。每个 Worker 应拥有独立的cwd并与TaskExecutionContext.appInfoPath的执行目录约定对齐。不要在此处添加更多切点。该切面刻意保持最小表面因为对 Hadoop 代码添加宽泛的 AspectJ 切点在不同 YARN 版本之间非常脆弱pointcut 匹配的是YarnClientImpl的具体方法签名YARN 升级可能改变实现类结构。AspectJ 版本升级敏感。aspectj.version固定在 dolphinscheduler-bom/pom.xml当前 1.9.7配套的aspectj-maven-plugin为 1.14.0。升级 AspectJ 后必须逐个回归验证所有基于 YARN 的任务插件MR、Spark、Sqoop、HiveCLI 等的织入效果。相关模块与排障指引dolphinscheduler-yarn-aop处于一条完整的 YARN 任务链上dolphinscheduler-worker运行时消费者该 jar 位于 Worker classpathdolphinscheduler-worker/pom.xml。dolphinscheduler-worker/CLAUDE.md 明确给出排障建议如果基于 YARN 的任务插件MR、Spark-on-YARN丢失了 application ID请检查 AspectJ 织入是否生效dolphinscheduler-task-plugin下的task-mr、task-spark、task-sqoop、task-hivecli等YARN 提交型任务插件是切面能力的主要受益方dolphinscheduler-yarn-aop/CLAUDE.md 列出dolphinscheduler-task-api的LogUtils/ProcessUtils/TaskExecutionContext读取appInfo.log、解析应用 ID、执行 kill 与回传 Master 的消费端。日常排查可遵循以下顺序先确认 Worker 进程工作目录下是否存在appInfo.log且内容包含application_\d_\d格式的 ID若为空检查 AspectJ 织入是否生效构建时showWeaveInfo与verbose已开启可查看 weave 日志若任务插件通过独立 classloader 加载 YARN 客户端如 Spark则需启用-javaagent:aspectjweaver.jar的运行期织入最后核对 Worker 的cwd是否与任务执行目录/appInfoPath指向一致避免多 Worker 共写同一文件。【免费下载链接】dolphinschedulerApache DolphinScheduler is the modern data orchestration platform. Agile to create high performance workflow with low-code项目地址: https://gitcode.com/GitHub_Trending/dol/dolphinscheduler创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门