尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
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
RELATED

相关推荐

如何从SQL快速创建数据库图表:drawDB导入功能完整指南

如何从SQL快速创建数据库图表:drawDB导入功能完整指南

如何从SQL快速创建数据库图表:drawDB导入功能完整指南 【免费下载链接】drawdb Free, simple, and intuitive online database diagram editor and SQL generator. 项目地址: https://gitcode.com/GitHub_Trending/dr/drawdb 你是否曾面对复杂的SQL脚本感到无…

📅 2026/9/15 4:34:11
2026年7月国内外小程序开发选择-主流小程序制作工具客观对比

2026年7月国内外小程序开发选择-主流小程序制作工具客观对比

一、汇总表 工具更适合谁价格开发方式核心特点餐宝盈适合所有行业的商家,尤其是拥有自己实体门店的商家,如餐饮、茶饮、烘焙、便利店、生鲜、社区零售门店,尤其适合先把点单、会员、发券和复购做起来的老板。99/年模板SAAS先点单、先会员、先…

📅 2026/9/15 4:34:50
10分钟掌握Taichi:为Python注入GPU加速魔法的终极指南

10分钟掌握Taichi:为Python注入GPU加速魔法的终极指南

10分钟掌握Taichi:为Python注入GPU加速魔法的终极指南 【免费下载链接】taichi Productive, portable, and performant GPU programming in Python. 项目地址: https://gitcode.com/GitHub_Trending/ta/taichi 还在为Python数值计算性能瓶颈而苦恼吗&#xf…

📅 2026/9/12 6:11:43
MORE NEWS

更多资讯

📰

Text2SQL 多租户行级权限隔离(RLS):AST 动态注入与数据越权防御

Text2SQL 多租户行级权限隔离(RLS):AST 动态注入与数据越权防御在将 Text2SQL 智能取数助手开放给企业内部数万名运营人员、或者对外开放给数千家入驻商户进行自助数据分析时,整个系统面临的最严峻、最高危的合规挑战莫过于——数…

📰

别只连蓝牙了!车机互联方案选择与连接教程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

Bootstrap 5 加载效果实战:从 spinner 到进度条,打造不焦虑的等待体验

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

📰

简历里的“精通 MySQL”:当我面试了 50 位资深工程师后的冷面复盘

简历里的“精通 MySQL”:当我面试了 50 位资深工程师后的冷面复盘九月秋招与社招旺季,作为大厂存储团队的技术专家,我的日历上排满了密集的技术面试。 在过去一个月里,我连续面试了 50 多位对标高级架构师(P7/P8&#…

📰

Redis版本升级实战:从5.0到6.2的性能优化与安全加固

1. Redis升级的必要性与场景分析在Linux服务器运维中,Redis作为关键的内存数据库,版本升级是每个运维人员必须掌握的技能。我经历过数十次生产环境的Redis升级,从早期的2.8版本到现在的7.0版本,每次升级都伴随着性能提升和新特性支…

📰

为什么不是最强的大模型反而成了工程师主力?

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬