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

工业AI平台数据血缘落地实践:Datablau方案全解析

中控技术把Datablau数据血缘落地到工业AI平台之后圈子里不少做数据治理的朋友都在讨论这件事。工业AI平台的数据治理一直不好做设备数据、工艺参数、实时告警、历史采样全混在一起链路又长又乱没有血缘关系做底子后面做指标统一、质量监控、模型迭代都容易翻车。这篇我以参与同类工业数据治理项目的视角把整个血缘落地的思路、技术细节、踩坑过程完整拆开来讲给正在做或者准备做数据治理升级的团队提供一个可以直接参考的样本。1. 工业AI平台做数据血缘先解决这4个痛点1.1 数据链路太长数据去哪了没人能说清工业AI平台的数据流向跟互联网业务很不一样。生产现场的DCS集散控制系统、PLC控制器、SCADA系统不断产生点位数据和时序数据这些数据经过采集网关进到Kafka一部分走Flink做实时计算一部分落到Hive数仓做离线分析再往后还有特征工程、模型训练、推理服务、BI报表等多个环节。一条数据从设备端到最终的展示看板中间经过五六层加工处理是家常便饭。问题在于数据每经过一道加工它的来龙去脉就模糊一分。业务方问这个设备健康度得分到底是怎么算出来的你很难当场给出完整答案开发想改一个上游表的字段也说不清楚会影响哪些下游指标和模型。血缘要解决的第一个问题就是把这整条链路用可视化的方式画出来谁是谁的上游、谁是谁的下游一图看清。1.2 数据质量出问题时排查像是大海捞针工业场景里指标口径混乱是常态。同一个设备开机率生产部门、设备部门、管理层看到的是三套数据因为底层取数逻辑分别写在不同团队维护的脚本里。一旦某一个指标数据异常开发人员要翻几十个SQL、查十几个调度任务一个个去找上游依赖运气好半天定位运气不好一两天就耗进去了。有了血缘之后情况完全不同。从异常指标节点反查上游直接能看到它依赖了哪几张表、经过了哪些加工逻辑、最后一次更新是什么时候快速锁定问题SQL把排查时间从小时级压缩到分钟级。这是血缘最立竿见影的价值也是中控技术愿意推动这件事的直接动力。1.3 实时链路越来越多传统血缘覆盖不住工业AI平台里实时计算的比重在快速上升。设备异常预警、工艺参数实时优化、能耗实时监控这些场景都依赖Flink跑实时任务。但大部分数据治理平台对实时计算的血缘支持很弱Flink SQL解析不标准DataStream API的加工逻辑更是完全没法从代码里直接提取。这就形成了一个尴尬的局面离线数仓的血缘建得漂漂亮亮实时链路上全是空白偏偏实时链路又是工业场景里最需要监控、最需要追根溯源的部分。实时血缘怎么采集、怎么和离线血缘统一呈现是这次落地中花精力最多的技术点之一后面我会详细讲。1.4 非结构化数据治理一直是盲区工业现场积累了海量非结构化数据设备运维手册、点检记录表、故障案例文档、工艺参数卡、PLC程序备份这些资料里藏着大量业务知识但传统的血缘体系基本只管结构化数据仓库非结构化数据完全游离在治理范围之外。举个例子一个设备故障分析报告里写了该故障与3号机组润滑油温度升高存在关联如果这句话不跟结构化数据里的润滑油温度字段建立关系那么这份报告永远只是躺在文件服务器里的一份文档无法参与数据资产的统一管理。中控技术这次特别把非结构化数据治理纳入了血缘建设范围思路和做法我在第3部分展开。2. Datablau血缘方案的整体设计思路2.1 为什么不直接用开源工具自己搭很多团队一提到数据血缘第一反应是拿Apache Atlas或者OpenMetadata开源项目自己搭一套。这思路不能说错但放到工业场景里开源工具的短板非常明显。首先是解析能力。OpenMetadata和Atlas的血缘更多停留在表级对字段级血缘的解析精度不够遇到复杂SQL、存储过程、自研脚本的时候解析结果往往缺胳膊少腿。其次是实时计算的支持。这两个开源项目对Flink SQL的血缘提取基本处于半成品状态DataStream API更是无能为力。最后是中文场景和工业场景的适配问题元数据模型偏向互联网业务设备数据、时序点位、工艺参数这类工业资产模型支持不够好。Datablau的血缘引擎封装了完整的SQL解析器和调度依赖解析器对Hive SQL、Spark SQL、Flink SQL的支持比较成熟字段级血缘的解析精度更高还支持通过API方式人工注入血缘。我们内部做了一下对比测试同样一批生产SQL脚本开源工具能解析出七成血缘关系Datablau能稳定覆盖九成以上这个差距在真正用起来的时候会非常明显。提示选型的时候不要光看DEMO效果一定要拿自己生产环境的真实脚本去测试解析准确率。不同行业的SQL写法差异很大工业场景里大量的自定义函数和复杂嵌套逻辑最容易暴露工具的真实水平。2.2 血缘采集的三种路径缺一不可血缘关系不是凭空生成的它来自三条采集路径的汇总路径一是元数据采集。通过Datablau的采集插件把Hive里的表、字段、分区信息Kafka里的Topic时序数据库里的测点模型统一拉取到资产管理中心。这些是血缘的节点没有它们血缘图就无处挂载。路径二是SQL解析。从调度平台拿到SQL脚本用血缘解析引擎生成抽象语法树提取出插入目标表和查询源表之间的关系。离线数仓的血缘主要靠这条路。路径三是任务编排解析。从DolphinScheduler、Airflow这类调度工具读取任务依赖关系把任务之间的时序依赖映射为数据加工的血缘边。三条路径采集的数据在血缘引擎里做融合去重最终形成统一的血缘图谱。实际操作中我发现有个细节特别容易被忽略调度依赖和SQL解析出来的血缘结果可能互相矛盾。比如SQL解析显示表B依赖表A但调度配置里任务B和任务A没有上下游关系。这种情况下新建血缘关系前要有一套收敛规则我们是调度依赖优先SQL解析辅助修正因为调度依赖是实际执行顺序的反映更接近真实数据流向。2.3 血缘分层设计别把技术图和业务图混在一起工业场景里血缘不能只画一张大而全的图得分层。我们在中控技术的落地中把血缘拆成三层。技术血缘解决数据从哪里来的问题描述表、字段、SQL任务、调度任务之间的依赖关系使用对象是数据开发工程师。指标血缘解决这个指标是怎么算出来的的问题从指标反查到它依赖的原子指标、派生指标和底层物理表使用对象是数据分析师和业务方。场景血缘解决这份数据资产对应业务哪个环节的问题把数据资产挂接到设备健康管理、工艺优化、能耗分析等具体业务场景上使用对象是业务管理者和平台运营方。三个层级的血缘在图上用不同的视图呈现底层共享同一套血缘图谱数据。这样既避免了业务方被技术细节淹没也避免了技术人员看不到业务价值。3. 核心细节与实操要点从采集到应用的完整链路3.1 元数据采集的三个细节直接影响血缘质量元数据采集是血缘建设的地基看起来简单实际操作中容易踩坑。第一个坑是采集周期。我们刚开始设置为每24小时全量采集一次结果发现Hive表的分区在一天内频繁变动血缘图上看到的字段信息已经过期了。后来调整为每30分钟增量采集一次元数据变更事件问题才解决。如果你们的调度任务比较密集建议把采集周期压缩到15分钟。第二个坑是表结构变更时的血缘处理。生产环境的表经常加字段、改字段类型如果每次都把旧的血缘关系直接删掉重建历史回溯就完全没法做了。我们的做法是在资产目录里做版本管理血缘关系打上时间戳查询时默认显示最新版本但历史版本全部保留需要排障的时候可以对比某一时刻的血缘快照。第三个坑是采集任务本身的稳定性。元数据采集任务挂了不会影响业务系统运行所以很容易被忽视。我们在监控平台配了专门的血缘采集告警规则采集任务连续失败三次就自动通知数据平台组处理避免血缘静默失效。3.2 SQL血缘解析表级容易字段级才是硬骨头SQL血缘解析里表级血缘是基础能力绝大部分工具都能做真正拉开差距的是字段级血缘。字段级血缘解析有四个老大难。第一是字段别名和表达式嵌套太深比如a.col1 b.col2 * c.col3这种多层计算解析器要正确识别每个输入列和输出列的对应关系。第二是CASE WHEN语句里多个条件分支对应不同来源字段解析器得把每个分支都拆出来。第三是多个字段合并成一个字段的加工逻辑比如把设备状态码的多个枚举值映射成统一的设备状态。第四是存储过程里用临时表做中间计算临时表的字段血缘很容易断掉。Datablau的解析引擎对前面三种情况处理得还算好但临时表导致的血缘断裂问题单靠解析器解决不了。我们配合做了一件事数据开发规范里明确要求存储过程里创建的临时表命名必须包含目标表的名字的关键字比如tmp_fact_device_status血缘引擎在解析时遇到这类临时表会自动做映射归并血缘链条就续上了。实操心得字段级血缘的准确率一半靠工具解析能力一半靠开发规范约束。我们定的规范是一字段一来源禁止一个输出字段同时拼接多个上游字段。刚开始数据开发团队觉得约束太多磨合两个月之后他们发现字段血缘全打通之后排查问题省下的时间远超写规范SQL多花的时间。下面是一段典型的ETL加工语句血缘引擎会解析出ods_device_signal、dim_device_info两张源表输出到dwd_device_status表并建立字段级的加工依赖关系INSERT OVERWRITE TABLE dwd_device_status SELECT d.device_id, d.device_name, s.signal_value, CASE WHEN s.signal_value 80 THEN warning WHEN s.signal_value 100 THEN fault ELSE normal END AS status_code, current_timestamp() AS etl_time FROM ods_device_signal s JOIN dim_device_info d ON s.device_id d.device_id WHERE s.dt ${bizdate};3.3 Flink实时计算血缘怎么处理这里有完整思路实时链路的血缘是这次项目里技术含量最高的部分也是行业里讨论最多的话题。中控技术的实时场景主要分三种形态处理方式各不相同。第一种是Flink SQL任务。这类任务有明确的SQL文本Datablau通过Flink SQL解析器可以直接提取血缘。关键是解析器要能识别Flink SQL特有的语法结构比如CREATE TABLE里定义的connector信息、INSERT INTO的写入逻辑。实测下来只要Flink任务是通过SQL方式提交的血缘准确率能做到和离线Hive SQL持平。第二种是DataStream API任务。这是真正的硬骨头代码里面的map、flatMap、process函数都是逻辑计算没有解析器能从Java代码里直接抽出血缘关系。我们的应对方案是注解标注API注入。开发人员在代码里对数据流的关键处理步骤加注释标签用Python脚本扫描代码注释生成血缘描述文件最后通过Datablau开放API批量注入血缘关系。第三种是实时数据落地数仓的链路。Flink任务把计算结果写到Hive表或者Kafka Topic再被下游任务消费。这个链路中我们需要把Flink任务的产出表和消费Topic建立关联同时标注加工逻辑描述让实时血缘和离线血缘在图谱上自然衔接。# 通过Datablau API批量注入Flink作业血缘的示例代码 import requests import json api_url http://datablau-server:8080/api/v1/lineage/batch headers {Content-Type: application/json, Authorization: Bearer ${API_TOKEN}} payload [{ source_type: kafka_topic, source_name: ods_device_signal_topic, target_type: flink_job, target_name: device_status_etl_job, fields: [ {source_field: device_id, target_field: device_id}, {source_field: signal_value, target_field: signal_value} ], job_type: streaming, description: 实时设备信号清洗入库 }] response requests.post(api_url, headersheaders, datajson.dumps(payload)) print(response.status_code)这里有个关键认知要讲清楚实时血缘没必要追求和离线血缘同样的解析深度。实时任务的调试窗口短、数据时效性强业务团队最关心的是这个指标的上游数据源是谁这条告警链路中间经过了哪些处理能够回答这两个问题实时血缘的核心价值就达成了。3.4 中间加工过程怎么处理很多人在这里卡住你搜openmetadata数据血缘怎么处理中间加工过程会发现一大堆人问这确实是所有血缘体系共同的痛点。所谓中间加工过程指的是数据从源表到目标表之间那些不容易被SQL解析直接捕捉的环节特征工程脚本、算法模型的训练样本生成、Python脚本做数据预处理、自定义Java程序做数据转换都属于这一类。在工业AI平台里中间加工过程尤其多。一个设备故障预测模型输入特征可能来自时序数据的滑动窗口计算、样本均衡处理、缺失值填充等多个步骤这些步骤往往用Python脚本实现SQL解析器完全无能为力。我们采用的组合策略是这样的。能自动解析的就自动解析这是基础能力。解析不到的部分用手动编排接口注入补全。Datablau资产中心允许用户手动在两个资产节点之间创建血缘关系也支持通过API批量导入。我们把所有Python脚本、算法模型、特征处理逻辑都先登记为任务节点再把它们和输入表、输出表挂上血缘边。比如一个特征工程的Python脚本它的血缘关系是输入节点是dwd_equipment_signal、dim_equipment_info输出节点是feature_equipment_health加工逻辑描述字段里写明滑动窗口7天、均值填充、标准化。这样一来血缘图谱上就不仅仅有表和SQL任务还有算法节点、特征节点、模型节点整条数据加工链路的完整性高了很多。注意手动注入血缘的时候一定要在血缘边上带上创建人和更新时间并且定期做审查。否则时间一长人工维护的血缘关系可能跟实际代码逻辑脱节血缘图谱就变成了官方造假。3.5 非结构化数据治理用三步把它拉进血缘体系非结构化数据没有schema结构无法通过SQL解析生成血缘。我们在这次项目里用的是资产登记关系抽取人工确认三步法。资产登记是把所有非结构化数据先变成数据资产节点。运维手册、点检记录表、故障案例文档、PDF图纸、PLC程序在Datablau资产中心里建立对应的文档资产条目录入设备编号、所属系统、文件路径等元数据信息。关系抽取是尝试建立非结构化数据和结构化数据之间的关联。我们首先用NLP工具对文档做关键词抽取识别出其中提到的设备编号、测点名称、指标名称然后把这些名称跟资产中心里已有的结构化字段做模糊匹配。比如一份故障案例文档里频繁出现润滑油温度和振动幅值这两个词组正好对应时序数据库里的两个测点字段系统就会生成一条待确认的血缘关联。人工确认是最后一道把关。NLP抽取结果只是候选关系需要数据治理专员在界面上逐条确认是否有效。确认通过的血缘关系进入正式图谱不通过的进入垃圾回收池。在中控技术的落地过程中人工确认通过率大概在60%到70%之间虽然需要投入一些人力但非结构化数据第一次真正汇入了数据资产版图。-- 在Datablau资产中心登记非结构化数据资产时推荐的标准字段结构 CREATE TABLE metastore.asset_document ( asset_id STRING COMMENT 资产唯一标识, doc_name STRING COMMENT 文档名称, doc_type STRING COMMENT 文档类型: manual/report/drawing/program, device_code STRING COMMENT 关联设备编码, file_path STRING COMMENT 文件存储路径, owner_dept STRING COMMENT 归属部门, extract_status STRING COMMENT 抽取状态: pending/confirmed/rejected, create_time TIMESTAMP COMMENT 登记时间, update_time TIMESTAMP COMMENT 更新时间 );4. 落地实操过程与实施效果4.1 实施节奏试点先行快速迭代别想着一步到位我们这次项目大概分了四个阶段推进。第一阶段是元数据盘点和血缘调研耗了两周。这个阶段不急着部署工具先把中控技术的数仓血缘现状摸排清楚梳理出优先建设的业务域。工业AI平台涉及的模块很多最终我们选了设备健康管理域作为试点因为这个域的数据链路比较长、跨系统多、业务痛点最明显。第二阶段是系统部署和基础采集花了一周。部署Datablau服务端配置Hive、Kafka、DolphinScheduler的采集插件跑通从元数据采集到血缘图呈现的基础链路。这个阶段要求低目标是把有一张图做出来。第三阶段是试点域血缘建设花了三周。把设备健康管理域涉及的所有表、任务、指标、模型的血缘关系全部建立起来重点攻克Flink实时任务和Python特征工程脚本的血缘补全问题。这个阶段结束的时候试点域的血缘覆盖率达到了90%以上。第四阶段是全量推广和运维交接花了两周。把试点阶段积累的血缘采集规范、开发约束、运维流程推广到全平台明确治理专员和数据开发团队各自的职责边界输出运维手册。项目经理要注意血缘建设最大的风险不是技术是业务部门不买账。落地过程中一定要在早期就拉上业务方和数据团队一起看血缘的价值演示让他们亲眼看到点开一个指标就能定位问题的效果。试点选得好推广阻力就小一大半。4.2 两个关键场景的落地实录场景一是设备健康管理的数据追溯。设备健康管理模型依赖大量DCS点位数据、点检记录、维修工单信息。以前模型特征数据出了问题算法工程师说数仓的问题数仓开发说源系统的问题两边互相排查平均要花两天才能定位到根因。血缘落地后从模型推理任务反查血缘直接能看到特征表的上游是哪个点位映射表、哪份点检记录、哪个维度表再加上数据质量监控规则定位问题的时间缩短到小时级。场景二是生产日报指标异常定位。日报里的某个指标突然跌了30%业务方立刻打电话过来问原因。以前只能安排开发逐层排查现在直接在血缘图上点击指标节点向下展开依赖的明细表和加工SQL再结合调度日志找到昨天执行失败的那个任务节点十分钟内定位根因沟通效率提升非常明显。4.3 效果数据与业务价值指标落地前落地后表级血缘覆盖率不足30%92%字段级血缘覆盖率不足10%78%数据问题平均定位时间8小时以上1小时以内指标口径冲突事件每月15起以上每月2-3起数据资产检索效率依赖人工答疑自助检索血缘导航这些数字不是单靠血缘工具本身达成的工具只提供了基础能力真正起作用的是配套的治理机制。我们同步建立了新任务上线前必须完成血缘自检的开发规范没有血缘登记的数据任务不允许发布到生产环境。这条规则一开始推行有阻力但坚持执行三个月后血缘覆盖率自然就维持在高位了。5. 常见问题与排障技巧速查5.1 问题速查表问题现象可能原因排查思路解决方案血缘图里出现幽灵表表已被删除但血缘缓存未更新检查元数据采集任务是否正常手动执行一次全量元数据刷新Flink加工的表血缘始终为空任务提交方式不支持自动解析确认任务是SQL还是DataStream API改用API注入或注解标注方式补全字段血缘解析结果与实际加工逻辑不一致SQL里有复杂表达式或UDF定位到具体SQL做单条解析测试优化SQL写法或手动修正血缘边血缘图查询特别慢血缘数据量过大且未做分区检查血缘存储表的查询计划按业务域分区存储血缘数据手动注入的血缘被自动解析覆盖两种来源的冲突处理策略不当查看血缘冲突日志配置手动优先或自动优先策略5.2 三个典型问题的排查过程问题一血缘图里出现幽灵表。有业务方反馈血缘关系显示某张表依赖了一张已经被归档的老表但那个老表早就下线了。排查发现是归档操作没有触发元数据变更事件让血缘缓存里一直留着旧记录。处理方式是找到这张老表的物理路径手动在资产中心做下线处理血缘图里的关联关系自动变成历史版本状态。问题二Flink加工的表血缘一直为空。阶段三试点时一张设备状态实时分析表在血缘图里死活找不到上游。查调度平台的作业列表确认这张表是由Flink作业产出的但提交方式不是SQL而是纯DataStream API代码。后来通过注解标注API注入的方式把Kafka源Topic、实时计算任务、目标表这条血缘链补齐了。问题三字段血缘解析结果和实际加工逻辑不一致。有一张表的上游字段总是解析成错的单条SQL测试也复现不了。最后是人工翻代码的时候发现的SQL脚本里用了自定义UDFUDF内部还做了三次表的关联查询解析器看不到UDF内部的逻辑只把UDF的入参和出参直接连起来了。这是解析工具的天然局限最终靠开发规范限制UDF内部禁止跨表查询来解决。5.3 避坑经验总结第一血缘覆盖率和血缘准确率要分开管理。覆盖率低说明还有数据链路没纳入采集范围这是量的问题准确率低说明血缘关系画错了这是质的问题。我们每周出一份血缘质量报告分别统计两个率排查的时候先看准确率再看覆盖率。第二新任务的血缘登记一定要卡在发布流程里。我们吃过亏上线了两个月的任务血缘一直是空的等想补的时候已经积累了大量的加工逻辑补起来费时费力。现在发布平台上强制校验血缘登记状态没有血缘确认的不允许发布。第三血缘不是建完就完的一次性工程。表结构变更、调度任务调整、加工逻辑重写每个环节都可能让血缘失真。数据团队要把血缘维护当成日常运维的一部分定期抽查血缘准确率持续优化解析规则和开发规范。写在最后数据血缘要真正用起来我在实际操作中越来越强烈的体会是数据血缘不能只做成一堆看上去很漂亮的关系图它最终要落到审计追踪、问题定位、影响分析这些具体使用场景里。中控技术这次项目能成功关键是把血缘和数据治理的日常流程绑定在一起开发规范、质量监控、发布流程都跟血缘打通了血缘才真正变成了数据平台的活地图。最后再分享一个小技巧让每个人都能在血缘图上找得到自己。我们在血缘图谱里对所有节点增加了责任人的维度每张表、每个任务、每个指标都有明确的owner。这样不仅方便追责更重要的是每个人都觉得血缘是自己的工具而不是治理团队拿来考核他们的枷锁。数据团队的配合度上去了项目就成功了一大半。后续这个体系还可以继续扩展把数据质量规则、数据权限策略、成本分析能力都挂到血缘节点上用一套资产网络把治理动作串起来那是更有意思的下一步。
分享:

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

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