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

📅 2026/7/20 16:28:26
javaagent-lineage-flink:基于 Java Agent 的 Flink 作业级血缘采集工具
javaagent-lineage-flink基于 Java Agent 的 Flink 作业级血缘采集工具项目地址https://github.com/TKilome/javaagent-lineage-flinkjavaagent-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 AgentFlink 作业可能来自 DataStream、Flink SQL也可能来自不同团队封装后的提交框架。如果在每种 API 或每套业务框架里单独埋点入口会越来越多维护成本也会越来越高。javaagent-lineage-flink选择拦截更靠近 Flink 提交流程核心的位置PipelineExecutorUtils#getJobGraph(...)当 Flink 生成JobGraph时作业的拓扑已经基本成型。Agent 可以从StreamGraph和JobGraph中读取作业元信息再结合版本匹配的 connector parser 解析外部读写端点。这样既能覆盖 DataStream也能覆盖 Flink SQL 场景。当前支持范围Flink 版本ConnectorConnector 版本支持能力1.19.3Kafka3.3.0-1.19DataStream / Flink SQL Kafka source 和 sink1.20.0Kafka3.4.0-1.20DataStream / Flink SQL Kafka source 和 sink1.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