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 可以从StreamGraph和JobGraph中读取作业元信息,再结合版本匹配的 connector parser 解析外部读写端点。这样既能覆盖 DataStream,也能覆盖 Flink SQL 场景。
当前支持范围
| Flink 版本 | Connector | Connector 版本 | 支持能力 |
|---|---|---|---|
1.19.3 | Kafka | 3.3.0-1.19 | DataStream / Flink SQL Kafka source 和 sink |
1.20.0 | Kafka | 3.4.0-1.20 | DataStream / Flink SQL Kafka source 和 sink |
1.20.0 | Paimon | 1.4.x | Paimon 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