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

使用 DataHub Java SDK V2 管理 DataFlow 实体:从 URN 建模到多编排器流水线实践

使用 DataHub Java SDK V2 管理 DataFlow 实体从 URN 建模到多编排器流水线实践【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本文以 DataHub Java SDK V2datahub-client中的DataFlow实体为核心讲解如何在 Airflow、Spark、dbt、Flink 等编排平台之上以统一的数据模型表达数据流水线pipeline/workflow。读完本文你将掌握 DataFlow 的 URN 结构、Builder 创建方式、属性与治理操作Owner/Tag/Glossary Term/Domain、Fluent 链式 API以及 DataFlow 与 DataJob 的父子层级关系并能在真实项目中使用DataHubClientV2直接落地代码。什么是 DataFlow 实体在 DataHub 的元数据模型中DataFlow实体用于表示数据加工流水线或工作流data processing pipeline / workflow。它抽象了来自 Apache Airflow、Apache Spark、dbt、Apache Flink 等数据编排与处理平台中的编排工作负载定时调度工作流如 Airflow DAG批处理作业如 Spark Application转换项目如 dbt project流式管道如 Flink job。DataFlow 在概念上是DataJob实体的父容器DataFlow 代表整条流水线而流水线中的具体任务task / step则由DataJob表示。这样的两级模型可以同时承载流水线级与任务级的治理元数据与血缘信息。在源码层面DataFlow 实体的定义位于 DataFlow.java其类声明为public class DataFlow extends Entity implements HasTagsDataFlow, HasGlossaryTermsDataFlow, HasOwnersDataFlow, HasDomainsDataFlow, HasSubTypesDataFlow, HasStructuredPropertiesDataFlow, HasDocumentationDataFlow它实现了标签、术语、Owner、Domain、子类型、结构化属性与文档等多组能力接口并通过getDefaultAspects()见 DataFlow.java声明了实体默认加载的 aspect包括Ownership、GlobalTags、GlossaryTerms、Domains、Status、InstitutionalMemory、DataFlowInfo与EditableDataFlowProperties。DataFlow 的 URN 结构每个 DataFlow 都由一个唯一 URN 标识格式如下urn:li:dataFlow:(orchestrator,flowId,cluster)三个组成分量分量说明示例orchestrator运行工作流的平台或工具airflow、spark、dbt、flinkflowId流水线在编排器内的唯一标识DAG ID、作业名、项目名my_dag_id、my_job_namecluster流水线运行的集群或环境prod、prod-us-west-2、emr-cluster常见示例urn:li:dataFlow:(airflow,customer_etl_daily,prod) urn:li:dataFlow:(spark,ml_feature_generation,emr-prod-cluster) urn:li:dataFlow:(dbt,marketing_analytics,prod)从源码看URN 的实现类为 DataFlowUrn.java它继承自通用Urn实体类型固定为dataFlow并强校验 TupleKey 必须恰好包含 3 个元素orchestrator / flowId / cluster。createFromString与createFromUrnDataFlowUrn.java会在解析时校验命名空间必须为li、实体类型必须为dataFlow否则抛出URISyntaxException。因此从字符串构造 DataFlow URN 时必须严格遵守上述三元组格式。创建 DataFlow基础创建使用DataFlow.builder()创建实体最少只需指定 URN 三要素随后调用client.entities().upsert(dataflow)将实体写入 DataHubDataFlow dataflow DataFlow.builder() .orchestrator(airflow) .flowId(my_dag_id) .cluster(prod) .displayName(My ETL Pipeline) .description(Daily ETL pipeline for customer data) .build(); client.entities().upsert(dataflow);带自定义属性创建可以通过customProperties(Map)在构建阶段直接注入键值对元数据MapString, String properties new HashMap(); properties.put(schedule, 0 2 * * *); properties.put(team, data-engineering); properties.put(sla_hours, 4); DataFlow dataflow DataFlow.builder() .orchestrator(airflow) .flowId(customer_pipeline) .cluster(prod) .customProperties(properties) .build();Builder 的源码级约束见 DataFlow.javaorchestrator、flowId、cluster三个字段必填缺失时build()会抛出IllegalArgumentException与 DataFlowTest.java 中的testDataFlowBuilderMissingOrchestrator/MissingFlowId/MissingCluster三个测试一一对应若未显式设置displayNameDataFlowInfo.name会回退为flowIdBuilder 也提供了urn(DataFlowUrn)重载可从现成 URN 反解出三个分量DataFlow.java底层 URN 由new DataFlowUrn(orchestrator, flowId, cluster)构造与前面所述的 URN 结构完全一致。DataFlow 属性详解核心属性属性类型说明示例orchestratorString运行流水线的平台必填airflow、spark、dbtflowIdString流水线唯一标识必填my_dag_id、my_job_nameclusterString集群/环境必填prod、dev、prod-us-west-2displayNameString人类可读名称Customer ETL PipelinedescriptionString流水线描述Processes customer data daily附加属性属性类型说明externalUrlString指向编排工具中该流水线的链接projectString关联的项目或命名空间customPropertiesMapString, String键值对元数据createdLong创建时间戳毫秒lastModifiedLong最后修改时间戳毫秒这些字段在元数据模型中对应DataFlowInfoaspectPDL 定义见 DataFlowInfo.pdl值得注意的细节name、description、project、created、lastModified、env均带有Searchable注解其中name使用WORD_GRAM分词、支持自动补全并作为_entityName参与搜索description为TEXT类型并生成hasDescription字段envFabricType会作为Environment 过滤器加入搜索筛选条件created与lastModified的语义是在源数据平台上的创建/修改时间而非 DataHub 上的时间对应 PDL 注释not on DataHubaspect 还继承了CustomProperties自定义属性与ExternalReference外部引用即externalUrl的载体。DataFlow 常用操作DataFlow 的所有变更操作都走patch 机制每次调用addTag、addOwner、setDescription等方法时并不会直接修改内存中的 aspect而是将操作累积进对应的PatchBuilder直到client.entities().upsert(dataflow)时才统一以 Metadata Change ProposalMCP的形式提交见 DataFlow.java 中getDataFlowInfoPatchBuilder()的实现。这意味着多次修改最终会合并为一次高效写入。Ownership所有者// 添加所有者 dataflow.addOwner(urn:li:corpuser:johndoe, OwnershipType.TECHNICAL_OWNER); dataflow.addOwner(urn:li:corpuser:analytics_team, OwnershipType.BUSINESS_OWNER); // 移除所有者 dataflow.removeOwner(urn:li:corpuser:johndoe);对应的测试 DataFlowTest.java 验证了addOwner后会生成 aspect 名为ownership的待提交 patch。Tags标签// 添加标签带或不带 urn:li:tag: 前缀均可 dataflow.addTag(etl); dataflow.addTag(production); dataflow.addTag(urn:li:tag:pii); // 移除标签 dataflow.removeTag(etl);Glossary Terms业务术语// 添加术语 dataflow.addTerm(urn:li:glossaryTerm:ETL); dataflow.addTerm(urn:li:glossaryTerm:DataPipeline); // 移除术语 dataflow.removeTerm(urn:li:glossaryTerm:ETL);Domain数据域// 设置 Domain dataflow.setDomain(urn:li:domain:DataEngineering); // 移除指定 Domain dataflow.removeDomain(urn:li:domain:DataEngineering); // 或清空所有 Domain dataflow.clearDomains();Custom Properties自定义属性// 添加单个属性 dataflow.addCustomProperty(schedule, 0 2 * * *); dataflow.addCustomProperty(team, data-engineering); // 移除属性 dataflow.removeCustomProperty(schedule); // 全量设置覆盖已有属性 MapString, String props new HashMap(); props.put(key1, value1); props.put(key2, value2); dataflow.setCustomProperties(props);描述与显示名// 设置描述 dataflow.setDescription(Daily ETL pipeline for customer data); // 设置显示名 dataflow.setDisplayName(Customer ETL Pipeline); // 获取描述 String description dataflow.getDescription(); // 获取显示名 String displayName dataflow.getDisplayName();时间戳与外部链接// 设置外部 URL dataflow.setExternalUrl(https://airflow.example.com/dags/my_dag); // 设置项目 dataflow.setProject(customer_analytics); // 设置时间戳 dataflow.setCreated(System.currentTimeMillis() - 86400000L); // 1 天前 dataflow.setLastModified(System.currentTimeMillis());需要说明的是externalUrl与project在 SDK 中通过直接修改DataFlowInfoaspect 缓存实现见 DataFlow.java而setCreated/setLastModified同样累积进dataFlowInfopatch所有 setter 均标注RequiresMutable对只读实体如从服务端加载后未调用mutable()调用会抛出异常。编排器专属示例Apache Airflow定时 DAGDataFlow airflowFlow DataFlow.builder() .orchestrator(airflow) .flowId(customer_etl_daily) .cluster(prod) .displayName(Customer ETL Pipeline) .description(Daily pipeline processing customer data from MySQL to Snowflake) .build(); airflowFlow .addTag(etl) .addTag(production) .addCustomProperty(schedule, 0 2 * * *) .addCustomProperty(catchup, false) .addCustomProperty(max_active_runs, 1) .setExternalUrl(https://airflow.company.com/dags/customer_etl_daily);Apache Spark批处理作业DataFlow sparkFlow DataFlow.builder() .orchestrator(spark) .flowId(ml_feature_generation) .cluster(emr-prod-cluster) .displayName(ML Feature Generation Job) .description(Large-scale Spark job generating ML features) .build(); sparkFlow .addTag(spark) .addTag(machine-learning) .addCustomProperty(spark.executor.memory, 8g) .addCustomProperty(spark.driver.memory, 4g) .addCustomProperty(spark.executor.cores, 4) .setDomain(urn:li:domain:MachineLearning);dbt转换项目DataFlow dbtFlow DataFlow.builder() .orchestrator(dbt) .flowId(marketing_analytics) .cluster(prod) .displayName(Marketing Analytics Models) .description(dbt transformations for marketing data) .build(); dbtFlow .addTag(dbt) .addTag(transformation) .addCustomProperty(dbt_version, 1.5.0) .addCustomProperty(target, production) .addCustomProperty(models_count, 87) .setProject(marketing) .setExternalUrl(https://github.com/company/dbt-marketing);Apache Flink流式管道DataFlow flinkFlow DataFlow.builder() .orchestrator(flink) .flowId(real_time_fraud_detection) .cluster(prod-flink-cluster) .displayName(Real-time Fraud Detection) .description(Real-time streaming pipeline for fraud detection) .build(); flinkFlow .addTag(streaming) .addTag(real-time) .addTag(fraud-detection) .addCustomProperty(parallelism, 16) .addCustomProperty(checkpoint_interval, 60000) .setDomain(urn:li:domain:Security);Fluent 链式 APIDataFlow的所有变更方法均返回this因此可以连续链式调用后一次性提交代码更紧凑、意图更清晰DataFlow dataflow DataFlow.builder() .orchestrator(airflow) .flowId(sales_pipeline) .cluster(prod) .build(); dataflow .addTag(etl) .addTag(production) .addOwner(urn:li:corpuser:owner1, OwnershipType.TECHNICAL_OWNER) .addOwner(urn:li:corpuser:owner2, OwnershipType.BUSINESS_OWNER) .addTerm(urn:li:glossaryTerm:Sales) .setDomain(urn:li:domain:Sales) .setDescription(Sales data pipeline) .addCustomProperty(schedule, 0 2 * * *) .addCustomProperty(team, sales-analytics); client.entities().upsert(dataflow);DataFlow 与 DataJob 的父子关系DataFlow 是 DataJob 的父实体DataJob 表示 DataFlow 内部的一个具体任务或步骤。创建时只需在 DataJob 的 Builder 中通过flow(dataflow.getUrn())引用父 DataFlow 的 URN// 创建父 DataFlow DataFlow dataflow DataFlow.builder() .orchestrator(airflow) .flowId(customer_etl) .cluster(prod) .build(); client.entities().upsert(dataflow); // 创建引用父 Flow 的子 DataJob DataJob extractJob DataJob.builder() .flow(dataflow.getUrn()) // 引用父 DataFlow .jobId(extract_customers) .build(); DataJob transformJob DataJob.builder() .flow(dataflow.getUrn()) .jobId(transform_customers) .build(); client.entities().upsert(extractJob); client.entities().upsert(transformJob);借助这一层级你可以以 DataFlow 建模整条流水线以 DataJob 建模流水线内的具体任务在任务级别追踪血缘与依赖关系在流水线与任务两级分别组织治理元数据。最佳实践保持命名一致组织内部统一 orchestrator 命名始终使用小写如airflow避免Airflow或AIRFLOW混用选择有意义的 clustercluster 名称应能体现环境与地域如prod-us-west-2、staging-eu-central-1补充调度信息批处理工作流请将调度表达式写入 custom properties回链源系统始终设置externalUrl以便跳转到编排工具的原生 UI尽早设置 Owner创建流水线时即指定技术负责人与业务负责人用标签做分类按类型etl、streaming、ml、环境production、staging与关键程度打标签记录 SLA通过 custom properties 记录 SLA 要求与告警渠道追踪版本对带版本的工作流如 dbt在 custom properties 中记录版本信息。完整示例一站式创建并治理 DataFlow下面是一个端到端示例覆盖客户端初始化、实体构建、治理元数据注入与提交对应 SDK 的完整用法客户端配置细节可参考 client.md 与 getting-started.md// 初始化客户端 DataHubClientConfigV2 config DataHubClientConfigV2.builder() .server(http://localhost:8080) .token(System.getenv(DATAHUB_TOKEN)) .build(); try (DataHubClientV2 client new DataHubClientV2(config)) { // 创建完整的 DataFlow MapString, String customProps new HashMap(); customProps.put(schedule, 0 2 * * *); customProps.put(catchup, false); customProps.put(team, data-engineering); customProps.put(sla_hours, 4); customProps.put(alert_channel, #data-alerts); DataFlow dataflow DataFlow.builder() .orchestrator(airflow) .flowId(production_etl_pipeline) .cluster(prod-us-east-1) .displayName(Production ETL Pipeline) .description(Main ETL pipeline for customer data processing) .customProperties(customProps) .build(); dataflow .addTag(etl) .addTag(production) .addTag(pii) .addOwner(urn:li:corpuser:data_eng_team, OwnershipType.TECHNICAL_OWNER) .addOwner(urn:li:corpuser:product_owner, OwnershipType.BUSINESS_OWNER) .addTerm(urn:li:glossaryTerm:ETL) .addTerm(urn:li:glossaryTerm:CustomerData) .setDomain(urn:li:domain:DataEngineering) .setProject(customer_analytics) .setExternalUrl(https://airflow.company.com/dags/production_etl_pipeline) .setCreated(System.currentTimeMillis() - 86400000L * 30) .setLastModified(System.currentTimeMillis()); // Upsert 到 DataHub client.entities().upsert(dataflow); System.out.println(Created DataFlow: dataflow.getUrn()); }延伸阅读DataJob Entity —— DataFlow 内部的子任务Dataset Entity —— DataFlow 的数据源与数据目标entities-overview.md —— SDK V2 实体体系总览patch-operations.md —— patch 更新机制详解DataFlowTest.java —— DataFlow 单元测试Builder 约束、Owner/Tag/Term/Domain 等操作验证DataFlowIntegrationTest.java —— DataFlow 集成测试【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
分享:

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

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