当前位置: 首页 > news >正文

javaagent-lineage-flink:基于 Java Agent 的 Flink 作业级血缘采集工具

javaagent-lineage-flink:基于 Java Agent 的 Flink 作业级血缘采集工具

项目地址:https://github.com/TKilome/javaagent-lineage-flink

javaagent-lineage-flink是一个面向 Apache Flink 的 Java Agent 血缘采集项目。它可以在 Flink 作业提交前拦截 JobGraph 生成流程,解析 DataStream / Flink SQL 作业中的 Source 和 Sink,并输出统一的LineageEvent血缘事件。

简单说,它现在能做这些事:

  • 支持 Flink DataStream 作业级数据血缘采集。
  • 支持 Flink SQL 作业级数据血缘采集。
  • 支持 Kafka source / sink 血缘解析。
  • 支持 Paimon source、sink 和 CDC combined dynamic sink 元数据解析。
  • 支持 Logging Reporter 输出单行 JSON。
  • 支持 HTTP Reporter 将血缘事件 POST 到外部元数据平台、数据地图或治理系统。
  • 支持按 Flink 版本和 connector 版本拆包适配,让兼容性边界更清楚。

在实时数仓和流式计算平台里,Apache Flink 往往承载着大量关键链路:订单、支付、履约、风控、营销、埋点、用户画像。随着作业数量增长,一个问题会越来越明显:我们知道作业在跑,但很难稳定、自动、低侵入地知道它到底读了哪些数据、写到了哪些数据。

这个项目就是为这个问题设计的:不要求每个业务作业改代码埋点,也不依赖作业运行后再从日志或外部系统反推,而是在作业真正提交运行前拿到更早、更明确的血缘事件。

为什么选择 Java Agent

Flink 作业可能来自 DataStream、Flink SQL,也可能来自不同团队封装后的提交框架。如果在每种 API 或每套业务框架里单独埋点,入口会越来越多,维护成本也会越来越高。

javaagent-lineage-flink选择拦截更靠近 Flink 提交流程核心的位置:

PipelineExecutorUtils#getJobGraph(...)

当 Flink 生成JobGraph时,作业的拓扑已经基本成型。Agent 可以从StreamGraphJobGraph中读取作业元信息,再结合版本匹配的 connector parser 解析外部读写端点。这样既能覆盖 DataStream,也能覆盖 Flink SQL 场景。

当前支持范围

Flink 版本ConnectorConnector 版本支持能力
1.19.3Kafka3.3.0-1.19DataStream / Flink SQL Kafka source 和 sink
1.20.0Kafka3.4.0-1.20DataStream / Flink SQL Kafka source 和 sink
1.20.0Paimon1.4.xPaimon source、精确表 sink、CDC combined dynamic sink 元数据

血缘事件可以通过 Reporter 输出到不同位置:

Reporter能力
Logging Reporter输出单行 JSON,适合本地调试、日志采集和快速验证
HTTP Reporter同步 POSTLineageEvent到外部 HTTP 服务,适合集成元数据平台、数据地图或数据治理系统

输出事件示例:

{"engineType":"flink","jobId":"...","jobName":"lineage-agent-kafka-debug","timestamp":1784357204468,"sources":[{"connector":"kafka","namespace":"broker-a:9092,broker-b:9092","name":"orders-input","properties":{"topic":"orders-input","bootstrap.servers":"broker-a:9092,broker-b:9092"}}],"sinks":[{"connector":"kafka","namespace":"broker-a:9092,broker-b:9092","name":"orders-output","properties":{"topic":"orders-output","bootstrap.servers":"broker-a:9092,broker-b:9092"}}]}

架构设计

项目采用模块化设计,把通用核心、Flink 版本适配、connector parser、reporter 分开打包。

javaagent-lineage-flink/ ├── lineage-core/ ├── lineage-flink/ │ ├── lineage-flink-1.19/ │ └── lineage-flink-1.20/ ├── lineage-reporter/ └── lineage-dist/

运行时只需要把lineage-core配置为-javaagent。对应 Flink 版本的 instrumentation、connector parser 和 reporter jar 放到 Flink classpath 中,通过 JavaServiceLoader自动发现。

处理链路很直接:

LineageAgent.premain() -> 发现 LineageFactory 实现 -> 安装 Flink instrumentation -> 拦截 PipelineExecutorUtils#getJobGraph(...) -> 提取 jobId、jobName、StreamNode -> parser registry 解析 source/sink dataset -> coverage validator 校验血缘完整性 -> reporter registry 上报 LineageEvent

设计原则

这个项目有几个明确取舍:

  • 不做运行时 Flink 或 connector 版本自动猜测。
  • 用户自行放入与运行环境匹配的 lineage jar。
  • lineage-core是唯一通过-javaagent指定的 jar。
  • instrumentation、parser、reporter 通过 Flink classpath 和 SPI 发现。
  • 解析、校验、上报失败会直接阻止作业提交。

这些取舍让系统更适合生产环境。血缘系统最怕“看起来成功,实际没拿到可信结果”。如果用户启用了 Agent,作业提交前就应该拿到明确、可信的血缘事件;拿不到就快速失败。

这个项目适合谁

如果你的 Flink 平台正在补数据治理、元数据采集或作业血缘能力,这个项目可以作为一个轻量、清晰、可扩展的起点。它尤其适合这些场景:

  • 已经有大量 Flink DataStream / Flink SQL 作业,不希望逐个改业务代码。
  • 希望在作业提交前就拿到 source、sink 和 job 维度的血缘事件。
  • 希望把血缘事件上报到内部元数据平台、数据地图或治理系统。
  • 希望以低侵入方式接入现有 Flink 集群。
  • 希望 connector 适配按版本显式管理,避免一个大包里混杂多套不兼容逻辑。
  • 希望基于 SPI 继续扩展 Hive、Iceberg、JDBC、OpenLineage 或其他上报方式。

javaagent-lineage-flink目前还处在持续演进阶段,但核心链路已经打通:Java Agent 插桩、SPI 扩展、Kafka/Paimon parser、Logging/HTTP reporter、发行包和 quickstart 文档都已具备。后续可以继续扩展 Hive、Iceberg、JDBC 等 connector,也可以演进到更多计算引擎。

项目地址:https://github.com/TKilome/javaagent-lineage-flink

http://www.cnnetsun.cn/news/3540936.html

相关文章:

  • 5步掌握SGLang多模态AI处理:从图像理解到视频分析实战指南
  • ddpo-pytorch核心功能解析:prompt_fn与reward_fn如何塑造生成式AI的创造力
  • 小程序毕设项目:用户行为驱动的智能音乐推荐系统实现 在线音乐资源聚合与智能推荐管理系统 (源码+文档,讲解、调试运行,定制等)
  • 电科网安保序加密检索技术解析与应用
  • MusicFreeDesktop:打造你的专属音乐空间,插件化播放器终极指南
  • 10个你不知道的Signature PDF实用技巧:让PDF处理更简单
  • 4大架构挑战深度解析:VPet虚拟桌宠核心系统设计与扩展方案
  • 告别卡文断更,10款爆火的 AI 写小说工具实测合集【7月最新指南】
  • 2026年国内外最新10款AI写小说软件(持续更新!)
  • 为什么用AI写小说还卡文?10款AI写作软件实测(内含工作流)
  • CefFlashBrowser终极指南:如何在Windows上完美运行经典Flash游戏
  • 如何快速部署轻量级AI模型:3步搞定跨平台推理
  • 录音修音一体的软件有哪些:从录音到导出的AI工具怎么选
  • 考证含金量高工商管理专业证书
  • AI编程团队每日站会失效的9种信号,及用LLM自动生成协作洞察报告的实操路径
  • Asp.net core Controller传值到视图的几种方式
  • Konado视觉小说框架:30分钟创建你的第一个互动故事游戏
  • Carnac系统托盘集成:Windows桌面应用的最佳实践
  • 终极终端输入法切换指南:告别手动切换的烦恼![特殊字符]
  • 小白程序员必备:收藏这份AI大模型学习地图,轻松入门20个核心概念!
  • 2026毕业必看|90%学生论文翻车的5个真相!避开这些坑,免费稳过双检
  • HarmonyOS应用开发实战:小事记 - Scroll 滚动容器深度剖析:滚动机制、edgeEffect、scrollBar 与嵌套滚动
  • WEEX API 接入指南:从创建 API Key 到完成首个行情请求
  • Checkra1n Windows版越狱工具使用指南与A12设备支持解析
  • 【AI字体适配紧急补丁】:从Figma Auto Layout到Canva Magic Design,3步强制锁定视觉节奏一致性
  • OpenNFS项目深度解析:从NFS1到NFS6的完整逆向工程之旅
  • MAA明日方舟助手:5分钟完成日常任务的终极解决方案
  • 移动端实时语义分割:在智能手机上实现精准图像识别
  • 如何在5分钟内搭建TonWeb开发环境:从安装到第一个TON应用
  • 【Linux 系统篇(四)】权限详解(一)