DataHub Flink 连接器完全指南:作业元数据、算子拓扑与跨平台血缘提取实战
DataHub Flink 连接器完全指南作业元数据、算子拓扑与跨平台血缘提取实战【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubApache Flink 是业界广泛使用的分布式流批一体处理框架。在数据治理场景中我们不仅需要知道数据长什么样更需要回答这份数据从哪来、被哪个作业加工、又流向哪里的问题。本指南以 DataHub 仓库中 Flink 数据源文档 为核心结合 源码实现 与 集成测试系统讲解 DataHub Flink 连接器的连接方式、概念映射、SQL/Table API 与 DataStream 两条血缘提取路径、平台解析机制、运行历史跟踪以及完整配置实战。读完本文你将能够独立编写一份可运行的 Flink 摄入配方Recipe并掌握排查作业可见但血缘为空等典型问题的完整思路。连接器概述它能从 Flink 提取什么DataHub 的 Flink 集成通过连接Flink JobManager REST API默认端口 8081来提取作业元数据、执行计划Execution Plan与运行历史当同时配置SQL Gateway默认端口 8083时它还能进一步通过 catalog 内省catalog introspection将 SQL/Table API 中的表引用解析到其真实平台Kafka、Postgres、Iceberg、Paimon 等并借助 DataProcessInstance 跟踪作业执行历史。同时支持基于状态摄入stateful ingestion的陈旧实体清理。从 FlinkSource 类定义 可以看到连接器的能力声明platform_name(Flink, idflink) config_class(FlinkSourceConfig) support_status(SupportStatus.BETA) capability(SourceCapability.PLATFORM_INSTANCE, Enabled by default) capability( SourceCapability.LINEAGE_COARSE, Table-level lineage from Kafka DataStream sources/sinks and SQL/Table API via catalog resolution, ) capability(SourceCapability.DELETION_DETECTION, Via stateful ingestion)即默认支持平台实例Platform Instance、表级粗粒度血缘并可通过状态摄入检测并删除陈旧元数据。当前连接器处于BETA支持状态使用时请留意其已知限制见下文限制清单。概念映射Flink 概念如何对应 DataHub 实体原文档给出了 Flink 概念到 DataHub 概念的映射关系这是理解后续所有配置项的基础Source ConceptDataHub ConceptNotesFlink JobDataFlow每个 Flink 作业对应一个 DataFlowFlink OperatorDataJob粒度取决于operator_granularity配置Job ExecutionDataProcessInstance当include_run_history开启时生成Kafka TopicDataset通过血缘解析DataStream 或 SQL/Table APIJDBC TableDataset通过 SQL Gateway catalog 内省解析Iceberg TableDataset通过 SQL Gateway 或catalog_platform_map配置解析在实体构建层面entities.py 中的FlinkEntityBuilder将作业构造成带丰富自定义属性的 DataFlow包括flink_job_id、job_state、flink_version、job_type、max_parallelism、start_time、duration_ms以及来自检查点配置的state_backend、checkpoint_interval_ms、checkpoint_mode、externalized_checkpoints等属性同时通过external_url指向 JobManager Web UI 的作业详情页格式为{rest_api_url}/#/jobs/{jid}。血缘涉及的 Dataset 实体也会被物化materialize确保血缘图谱中的数据集实体在 DataHub 中真实存在。前置条件与所需权限要接入 Flink 元数据你需要能够访问启用了JobManager REST API的 Flink 集群默认端口 8081Flink 版本 1.16官方以 1.19 做过测试基于DESCRIBE CATALOG的平台解析需要Flink 1.20若要解析 SQL/Table API 作业的跨平台血缘还需要访问Flink SQL Gateway默认端口 8083。所需权限如下表能力API所需访问权限作业元数据、运行历史JobManager REST API/v1/jobsREST API 的读访问权限平台解析血缘SQL Gateway REST API/v1/sessions会话创建与 SQL 执行权限如果 Flink 集群启用了认证bearer token 或 basic auth请在connection配置中提供凭据同一组凭据会同时用于 JobManager 与 SQL Gateway 两个 API。完整配方Recipe实战以仓库自带的 flink_recipe.yml 为蓝本一份完整的摄入配方如下source: type: flink config: connection: rest_api_url: http://localhost:8081 # SQL Gateway 用于解析 SQL/Table API 血缘的平台。 # 不配置时只能提取 DataStream Kafka 血缘。 # sql_gateway_url: http://localhost:8083 # 认证取消注释其中一种方式 # token: ${FLINK_API_TOKEN} # username: admin # password: ${FLINK_PASSWORD} # 高级连接调优 # timeout_seconds: 30 # max_retries: 3 # verify_ssl: true # 按作业名过滤 # job_name_pattern: # allow: # - ^prod_.* # 按作业状态过滤默认为 RUNNING、FINISHED、FAILED、CANCELED # include_job_states: # - RUNNING # - FINISHED # 按 catalog 的 platform_instance 覆盖。大多数 catalog 类型可通过 # SQL Gateway 自动检测平台仅在自动检测不可用时才需显式指定平台 # 详见文档。 # catalog_platform_map: # pg_catalog: # platform_instance: prod-postgres # kafka_catalog: # platform_instance: prod-kafka # DataJob 粒度 - job默认或 vertex每个算子一个 # operator_granularity: job # include_lineage: true # include_run_history: true # 平台级 platform_instance 兜底当 catalog_platform_map 中没有 # 该 catalog 的条目时使用。 # platform_instance_map: # kafka: prod-kafka-cluster # 并行作业处理 # max_workers: 10 # 陈旧实体清理 # stateful_ingestion: # enabled: true # remove_stale_metadata: true env: PROD sink: type: datahub-rest config: server: http://localhost:8080运行该配方datahub ingest -c flink_recipe.yml连接配置详解Connection Config在 config.py 中FlinkConnectionConfig定义了以下连接参数参数默认值说明rest_api_url必填JobManager REST API 端点例如http://localhost:8081sql_gateway_urlNoneSQL Gateway REST API 端点例如http://localhost:8083。提供后才会解析 SQL/Table API 表引用到真实平台kafka、postgres、iceberg 等tokenNoneBearer token 认证与 username/password 互斥username/passwordNoneHTTP Basic 认证必须成对提供timeout_seconds30HTTP 请求超时秒最小 1max_retries3失败请求的最大总尝试次数初始 1 次 最多重试 2 次带指数退避verify_ssltrue对 HTTPS 连接校验 SSL 证书sql_gateway_operation_timeout_seconds60等待 SQL Gateway 操作SHOW CATALOGS、DESCRIBE CATALOG 等完成的最大秒数最小 5对慢速 catalog 后端可调大值得注意的底层行为来自config.py的校验器URL 尾部斜杠自动剥离rest_api_url与sql_gateway_url会被自动执行rstrip(/)因此两种写法等价认证方式互斥同时配置token与username/password会抛出ValueErrorusername与password必须成对出现否则同样报错作业状态白名单校验include_job_states必须非空且每个状态必须是 Flink 合法状态之一。完整合法集合为INITIALIZING、CREATED、RUNNING、FAILING、FAILED、CANCELLING、CANCELED、FINISHED、RESTARTING、SUSPENDED、RECONCILING。配置值会被统一转换为大写未配置 SQL Gateway 的提示当include_lineage开启但sql_gateway_url未配置时连接器会记录一条日志提示 SQL/Table API 作业的血缘将仅限于 DataStream Kafka 的 source/sink。血缘提取的两条路径连接器通过分析 Flink 执行计划来提取表级血缘其核心实现在 lineage.py。它内部使用多个血缘提取器LineageExtractor分别处理不同形态的算子描述目前覆盖两种截然不同的场景场景一DataStream API仅 Kafka连接器识别算子描述中的KafkaSource-{topic}与KafkaSink-{topic}模式。由于 Kafka 连接器类名本身包含在算子名中平台固定为kafkatopic 名直接从描述中提取无需 SQL Gateway。对应源码是DataStreamKafkaExtractor其正则为KAFKA_DATASTREAM_SOURCE re.compile(rKafkaSource-([^\s])) # 兼容旧格式 KafkaSink-topic - ... 与 Flink 1.19 Kafka Sink V2 的 # KafkaSink-topic: Writer 和 KafkaSink-topic[N]: Writer 格式 KAFKA_DATASTREAM_SINK re.compile(rKafkaSink-([^:\s\[]))场景二SQL/Table API所有连接器连接器解析TableSourceScan(table[[catalog, db, table]])与Sink(table[[catalog, db, table]])模式。这些是通用的 Flink 计划格式——对 Kafka、JDBC、Iceberg、Paimon 等所有连接器都相同。真正的难点在于解析出 catalog 路径后还需要知道这个 catalog 对应哪个平台。这一步通过 SQL Gateway catalog 内省完成解析优先级如下catalog_platform_map配置—— 用户提供的覆盖项优先级最高DESCRIBE CATALOGFlink 1.20—— 判定 catalog 类型jdbc、iceberg、paimon、hive 等SHOW CREATE TABLE—— 从表 DDL 中读取connector属性适用于 hive / generic_in_memory 这类混合连接器类型的 catalog。从源码看SqlTableExtractor处理四种计划描述格式双括号逗号分隔的 sourceTableSourceScan(table[[catalog, db, table]])、双括号 sink旧式 SinkFunction 连接器、单括号点分隔的 sinkSink(table[catalog.db.table])用于 JDBC、filesystem 等 Sink V2 连接器、以及算子链产生的tableName[N]: Writer模式。其中最后一种不含 catalog/db 信息无法解析为平台详见限制清单第 4 条。平台解析机制与实例平台自动解析示例原文档给出了两个典型的平台解析链路可以直观理解表引用 → 平台 → Dataset URN的转换过程。Flink 作业读取 Postgres JDBC catalog 表pg_catalog.mydb.public.usersPlan: TableSourceScan(table[[pg_catalog, mydb, public.users]]) → SQL Gateway: DESCRIBE CATALOG → typejdbc, base-urljdbc:postgresql:// (Flink 1.20) → URN: urn:li:dataset:(urn:li:dataPlatform:postgres, mydb.public.users, PROD)Flink 作业读取 Iceberg catalog 表ice_catalog.lake.eventsPlan: TableSourceScan(table[[ice_catalog, lake, events]]) → SQL Gateway: DESCRIBE CATALOG → typeiceberg (Flink 1.20) → URN: urn:li:dataset:(urn:li:dataPlatform:iceberg, lake.events, PROD)JDBC 平台推断在 lineage.py 中有明确映射表jdbc:postgresql→postgres、jdbc:mysql→mysql、jdbc:sqlserver→mssql、jdbc:oracle→oracle、jdbc:db2→db2、jdbc:mariadb→mariadb、jdbc:redshift→redshift、jdbc:snowflake→snowflake。同时对于 hive / generic_in_memory 这类可能混合多种连接器的 catalog连接器会通过SHOW CREATE TABLE读取 DDL 中的connector属性再根据连接器平台映射表kafka、jdbc、iceberg、paimon、elasticsearch、mongodb、kinesis、clickhouse、doris、starrocks、pulsar 等二十余种解析出平台。平台实例映射Platform Instance Mapping如果你的数据集属于特定的平台实例例如某个具体的 Kafka 集群或 Postgres 部署可以使用catalog_platform_map做按 catalog 的映射或用platform_instance_map做平台级的兜底source: type: flink config: connection: rest_api_url: http://localhost:8081 sql_gateway_url: http://localhost:8083 # 按 catalog优先级更高 catalog_platform_map: pg_us: platform_instance: us-postgres pg_eu: platform_instance: eu-postgres # 平台级兜底 platform_instance_map: kafka: prod-kafka-clustercatalog_platform_map中的每个条目支持两个字段对应CatalogPlatformDetail见 config.pyplatform该 catalog 下数据集对应的 DataHub 平台名如iceberg、postgres、kafka。在 Flink 1.20 上通常由 DESCRIBE CATALOG 自动检测仅当自动检测失败时才需要显式指定在 Flink 1.20 上Iceberg/Paimon catalog 必须显式指定platform_instance该 catalog 下数据集的平台实例名如prod-postgres、us-east-kafka用于区分同一平台的多个部署。在 entities.py 的compute_dataset_urns中可以看到平台实例的解析优先级先查catalog_platform_map中对应 catalog 的platform_instance若未命中再回退到platform_instance_map按平台名取值最终通过make_dataset_urn_with_platform_instance构造带平台实例的 URN。Iceberg / Paimon 在 Flink 1.20 上的处理在 Flink 1.20 之前的版本上DESCRIBE CATALOG不可用。连接器会回退到SHOW CREATE TABLE但 Iceberg 和 Paimon 表的 DDL 中没有connector属性因此必须通过catalog_platform_map显式指定平台source: type: flink config: connection: rest_api_url: http://localhost:8081 sql_gateway_url: http://localhost:8083 catalog_platform_map: ice_catalog: platform: iceberg paimon_catalog: platform: paimon而在 Flink 1.20 上平台会根据 catalog 类型自动检测此配置不再必需。作业过滤与并行处理连接器在 source.py 中实现了两层作业过滤按名称过滤job_name_pattern使用 DataHub 标准的AllowDenyPattern支持allow/deny正则列表按状态过滤include_job_states默认值为[RUNNING, FINISHED, FAILED, CANCELED]可自行调整例如只摄入已完成的作业。随后会对按名称状态筛选出的作业做去重同一名称存在多个运行实例不同 job id时保留start_time最新的那个并记录警告日志。作业详情的拉取采用ThreadPoolExecutor并行进行max_workers控制并行线程数默认 10范围 1~50这在集群作业数量较多时可显著加快摄入速度。每个作业的处理流程_process_job为拉取作业详情 → 拉取检查点配置 → 从执行计划提取血缘 → 构建 DataFlow/DataJob → 计算数据集 URN 并物化 → 按需构建 DataProcessInstance。算子粒度Operator Granularity默认配置下operator_granularity: job连接器为每个 Flink 作业发出一个 DataJob所有 source/sink 血缘都合并coalesced到这个 DataJob 上。DataJob 的自定义属性包含flink_job_id与operator_count执行计划中的算子节点数。将operator_granularity设为vertex则为执行计划中的每个算子/顶点各发出一个 DataJob血缘粒度更细但实体数量也随之增多。选择哪种粒度取决于你的治理诉求若只想回答作业读了哪些表、写了哪些表job粒度足够若需要逐算子排查数据流经的每一步可选用vertex。运行历史Run History与 DataProcessInstance当include_run_history开启默认开启时连接器会为作业执行发出DataProcessInstance实体来跟踪单次执行开始与结束时间戳取自 Flink 作业时间线timeline运行结果映射FINISHED→ SUCCESSFAILED→ FAILURECANCELED→ SKIPPED处理类型STREAMING 或 BATCH取决于 Flink 作业类型。状态为RUNNING的作业只发出 start 事件已完成的作业同时发出 start 与 end 事件。DPIs 会关联到对应 DataJob 以及其 inlets/outlets 数据集从而在 DataHub 中形成数据集 ← 作业执行 → 数据集的完整数据流链路。状态摄入与陈旧实体清理stateful_ingestion配置继承自StatefulStaleMetadataRemovalConfig可开启状态摄入用于软删除soft-delete陈旧实体。典型配置stateful_ingestion: enabled: true remove_stale_metadata: true当某个 Flink 作业从集群中下线后再次摄入时连接器会依据上次摄入的检查点checkpoint状态将不再存在的作业及关联实体标记为删除。摄入报告中也会通过StaleEntityRemovalSourceReport记录清理统计。连接测试Test Connection连接器实现了TestableSource可通过 DataHub CLI 的test connection能力验证连通性。其逻辑source.py分三步校验配置合法性不合规则报告 basic_connectivity 失败调用client.test_connectivity()探测 Flink REST API 连通性随后调用get_jobs_overview()验证能否列出作业对应 LINEAGE_COARSE 能力若配置了sql_gateway_url则额外通过FlinkSQLGatewayClient.test_connection()探测 SQL Gateway并把结果作为单独的sql_gateway能力项报告。摄入报告Report指标FlinkSourceReportreport.py提供了丰富的诊断指标排查问题时非常有用jobs_discovered扫描到的作业数、jobs_filtered_by_name/jobs_filtered_by_state被过滤数、jobs_processed处理成功数、jobs_failed失败作业列表、lineage_extracted/lineage_failed、lineage_sources_found/lineage_sinks_found、lineage_unclassified_nodes无法分类的血缘节点描述排查血缘为空的第一入口、dpis_emitted以及flink_version。限制清单原文档明确列出了 7 项限制使用前务必了解SQL/Table API 血缘必须有 SQL Gateway。未配置sql_gateway_url时连接器无法将TableSourceScan(table[[catalog, db, table]])引用解析到实际平台。DataStream Kafka 血缘KafkaSource-{topic}不需要 SQL Gateway 即可工作。catalog 必须对 SQL Gateway 会话可见。在作业代码中通过编程方式注册的 catalog、临时 SQL 客户端会话或独立 FileCatalogStore 中的 catalog 对连接器不可见。生产部署应使用持久化 catalog例如由 Hive Metastore 支撑的 HiveCatalog确保表定义跨会话可见。Flink 1.20 上的 Iceberg/Paimon 需要显式配置。DESCRIBE CATALOG自 Flink 1.20 引入。在更早版本上由于 Iceberg/Paimon 表在SHOW CREATE TABLE中没有connector属性无法自动检测平台需用catalog_platform_map手动指定。算子链上的 sink 无 catalog 信息。Flink 算子链产生的tableName[N]: Writer模式不包含 catalog 和 database 信息只有裸表名。此类 sink 无法解析到平台会被报告为 unclassified未分类。临时表对 SQL Gateway 不可见。CREATE TEMPORARY TABLE的定义是会话级别的不会持久化到任何 catalog。SQL Gateway 无法查到其定义因此临时表无法解析到平台同样被报告为 unclassified。不支持 DataStream 的非 Kafka 连接器。只识别KafkaSource-{topic}和KafkaSink-{topic}两种 DataStream 模式。其他 DataStream 连接器Kinesis、Pulsar、RabbitMQ、自定义等只会产生用户提供的名称没有平台信息。无列级血缘。目前只从执行计划提取表级粗粒度血缘。常见问题排查TroubleshootingFailed to connect to Flink cluster检查rest_api_url是否正确且可达。可手动验证curl http://host:8081/v1/config。源码中该错误对应get_workunits_internal里client.get_cluster_config()抛出的异常会被记录到摄入报告的 failure 中。作业可见但没有血缘检查摄入报告中的 unclassified 节点。常见原因SQL/Table API 作业未配置sql_gateway_url—— 添加 SQL Gateway URL表位于临时会话中创建的default_catalogGenericInMemoryCatalog—— 改用 HiveCatalog 等持久化 catalogDataStream 作业使用非 Kafka 连接器 —— 当前不支持。SQL Gateway 已配置但平台未解析在 Flink 1.20 上DESCRIBE CATALOG不可用。检查表的SHOW CREATE TABLE输出是否包含connector属性对于 Iceberg/Paimon catalog添加catalog_platform_map配置。此外还可以通过调大sql_gateway_operation_timeout_seconds默认 60 秒应对慢速 catalog 后端。血缘 URN 与其他连接器如 Kafka 连接器不匹配确保platform_instance_map或catalog_platform_map产出的平台实例与其他摄入源一致。例如 Kafka 连接器使用了platform_instance: prod-cluster则这里需配置platform_instance_map: kafka: prod-cluster否则同一 topic 会在 DataHub 中形成两套不同的 Dataset URN导致血缘无法打通。深入阅读连接器全部源码位于 metadata-ingestion/src/datahub/ingestion/source/flink/其中 config.py 是全部配置项的权威定义lineage.py 是血缘提取与平台解析的核心source.py 是作业处理主流程单元测试覆盖了配置校验、血缘提取、平台解析、实体构建与 SQL Gateway 客户端metadata-ingestion/tests/unit/flink/集成测试通过 docker-compose 拉起真实 Flink SQL Gateway 环境验证了 Kafka、Postgres、Iceberg 三类血缘拼接、无 SQL Gateway 场景、运行历史与平台实例、连接测试等场景metadata-ingestion/tests/integration/flink/test_flink.py 及其 golden 文件 flink_mces_golden.json带 SQL Gateway与 flink_no_gw_mces_golden.json无 SQL Gateway摄入配方模板见 metadata-ingestion/docs/sources/flink/flink_recipe.yml涉及的数据模型文档DataFlow、DataJob、DataProcessInstance、Dataset。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →