ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

DataHub Rest Emitter 深度指南:通过 REST 协议向 DataHub 推送元数据的完整实战

DataHub Rest Emitter 深度指南:通过 REST 协议向 DataHub 推送元数据的完整实战 DataHub Rest Emitter 深度指南通过 REST 协议向 DataHub 推送元数据的完整实战【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub导读Rest Emitter 是 DataHub Python SDK 提供的核心元数据发射器Emitter它把 Metadata Change ProposalMCP、Metadata Change EventMCE等元数据对象序列化后通过 HTTP 直接推送到 DataHub GMSGraph Metadata Service是无需部署 Kafka、即可快速接入 DataHub 元数据的首选通道。本文以仓库中 rest-emitter.rst 文档为骨架结合 rest_emitter.py 的完整源码实现深入讲解 Rest Emitter 的初始化参数、emit 模式、REST 与 OpenAPI 两种传输协议、批量发送与自动分块、重试机制、鉴权与会话管理等细节帮助你写出生产可用的元数据推送代码。一、什么是 Rest Emitter在 DataHub 元数据体系中写入元数据主要有两条路径Kafka 异步链路通过DatahubKafkaEmitter写入 Kafka topic由 MCE/Mae Consumer 消费和REST 直连链路。Rest Emitter 属于后者它基于requests库构建 HTTP 会话把MetadataChangeProposalMCP、MetadataChangeProposalWrapperMCPW或MetadataChangeEventMCE以 JSON 形式 POST 到 GMS 的 REST 接口。从源码看DataHubRestEmitter同时实现了两个抽象契约generic_emitter.pyEmitterProtocol定义emit()与flush()两个方法flush()在 REST 场景下是空操作pass因为每次请求同步完成Closeable通过close()关闭底层requests.Session释放连接池。DataHubRestEmitter还同时扮演了上层组件的基础角色DataHubGraph继承自 Rest Emitter源码注释中标注DataHubGraph inherits from the rest emitter即 Graph Client 底层仍走 REST元数据 ingestion 的RestSink底层也复用 Rest Emitter源码注释中标注the rest sink uses the rest emitter under the hood。因此理解 Rest Emitter 就等于理解了 DataHub 大部分写入路径的底层机制。二、快速上手最小可用示例仓库的 examples/library 与 examples/structured_properties 中提供了大量可运行示例。最小的使用方式如下摘自 dataset_replace_properties.pyfrom datahub.emitter.mce_builder import make_dataset_urn from datahub.emitter.rest_emitter import DataHubRestEmitter from datahub.specific.dataset import DatasetPatchBuilder gms_endpoint http://localhost:8080 emitter DataHubRestEmitter(gms_servergms_endpoint) dataset_urn make_dataset_urn(platformhive, namefct_users_created, envPROD) property_map_to_set { cluster_name: datahubproject.acryl.io, retention_time: 2 years, } with emitter: for patch_mcp in ( DatasetPatchBuilder(dataset_urn) .set_custom_properties(property_map_to_set) .build() ): emitter.emit(patch_mcp)再如 update_structured_property.py 展示了向 URN 为io.acryl.dataManagement.dataSteward的结构化属性打补丁的用法。示例中with语句块会自动调用close()关闭会话这是推荐的资源管理方式。三、构造函数参数详解DataHubRestEmitter.__init__支持一组丰富的参数见 rest_emitter.py下面按功能分组说明其中标注的默认值均来自源码常量。3.1 服务地址与鉴权参数默认值说明gms_server必填GMS 地址形如http://localhost:8080也可传特殊值__from_env__从环境变量解析。源码会调用fixup_gms_url对 URL 做归一化tokenNoneBearer Token设置后会在请求头写入Authorization: Bearer tokenauthNonerequests.auth.AuthBase实例用于 OAuth Token Provider 等动态鉴权优先级高于token且不会把 Authorization 头固化到 headersclient_certificate_path/client_key_path/ca_certificate_pathNone双向 TLS 客户端证书与 CA 证书路径disable_ssl_verificationFalse是否关闭 SSL 证书校验鉴权解析的完整顺序见构造函数源码 rest_emitter.py显式传入auth实例每次请求由 AuthBase 动态生成 Authorization否则若传入token写入Authorization: Bearer token否则若系统环境变量配置了 system auth如DATAHUB_GMS_TOKEN自动注入。此外gms_server__from_env__时构造函数会通过config_utils.require_config_from_env()解析环境变量中的服务地址与 token并支持基于环境变量的 OAuth 配置DATAHUB_AUTH_TYPE。3.2 超时与重试参数默认值说明timeout_sec30统一超时连接 读取对应_DEFAULT_TIMEOUT_SEC 30connect_timeout_sec/read_timeout_secNone分别指定连接与读取超时只要其一被设置就会构成(connect, read)元组retry_status_codes[429, 500, 502, 503, 504]触发重试的 HTTP 状态码_DEFAULT_RETRY_STATUS_CODESretry_methods[HEAD,GET,POST,PUT,DELETE,OPTIONS,TRACE]允许重试的 HTTP 方法retry_max_times4最大重试次数环境变量DATAHUB_REST_EMITTER_DEFAULT_RETRY_MAX_TIMES可覆盖源码中还实现了一个有趣的细节_WeightedRetryrest_emitter.py针对 429限流做了加权处理——每次遇到 429 都会按_429_RETRY_MULTIPLIER默认 2环境变量DATAHUB_REST_EMITTER_429_RETRY_MULTIPLIER可调扩充剩余重试次数即限流时的总重试预算约为retry_max_times * 2并会尊重服务器返回的Retry-After头。重试退避因子为2backoff_factor2。3.3 连接池与 TCP 保活参数默认值说明pool_connections100连接池缓存的最大连接数DATAHUB_REST_EMITTER_DEFAULT_POOL_CONNECTIONSpool_maxsize100单主机连接池最大连接数DATAHUB_REST_EMITTER_DEFAULT_POOL_MAXSIZEtcp_keepalive由环境变量DATAHUB_REST_SINK_DEFAULT_TCP_KEEPALIVE决定是否启用 TCP Keepalivetcp_keepaliveTrue时使用自定义的_KeepAliveHTTPAdapterrest_emitter.py它会给池化连接设置SO_KEEPALIVE并在 Linux/macOS 上额外设置TCP_KEEPIDLE60、TCP_KEEPINTVL10、TCP_KEEPCNT5用于防止空闲连接被服务端关闭后触发SSLEOFError。如果平台不支持这些 socket 选项会自动优雅降级为普通HTTPAdapter。当连接池中的 SSL 连接报错时send()会先关闭旧连接再用全新连接重试一次。3.4 传输协议与 emit 模式参数默认值说明openapi_ingestionNone是否使用 OpenAPI 端点/openapi/v3而非传统 RestLi 端点。为None时由服务端能力协商决定SDK 客户端且服务端支持OPEN_API_SDK特性则自动启用见_post_fetch_server_configrest_emitter.pydefault_emit_mode全局默认受环境变量DATAHUB_EMIT_MODE影响兜底为SYNC_PRIMARY实例级默认 emit 模式respect_mcp_sync_markerFalse是否尊重 MCP 系统元数据中的emitModeMarkersync标记开启后只要批次中任一 MCP 带 sync 标记整批自动升级为同步发送server_config_refresh_intervalNone服务端配置缓存刷新间隔秒超过该间隔会重新拉取/configemit_mode参数与async_flag已废弃二选一async_flagTrue强制ASYNCasync_flagFalse强制SYNC_PRIMARY见 rest_emitter.py。四、核心方法emit 家族DataHubRestEmitter对外暴露的核心方法可以分成四类全部定义在 rest_emitter.py 中。4.1 emit统一入口emit(item, callbackNone, emit_modeNone)rest_emitter.py是Emitter协议约定的统一入口根据 item 类型自动路由UsageAggregation→emit_usage()已废弃改用datasetUsageStatisticsaspectMetadataChangeProposal/MetadataChangeProposalWrapper→emit_mcp()MetadataChangeEvent→emit_mce()。callback形如Callable[[Exception, str], None]成功时以(None, success)调用失败时以(exception, str(exception))调用随后重新抛出异常。4.2 emit_mcp / emit_mce单条写入emit_mcp(mcp, emit_modeNone, wait_timeouttimedelta(seconds3600))返回Optional[TraceData]异步追踪信息是现在最常用的写入方法。其内部处理逻辑rest_emitter.py解析 emit 模式兼容废弃的async_flag调用ensure_has_system_metadata确保 MCP 带系统元数据若走 OpenAPI把 MCP 转成OpenApiRequest见 request_helper.py发起 POST若走 RestLiDELETE类型且 aspect 非 Key aspect 时直接抛OperationalErrorRestLi 只允许删除 Key aspect否则 POST 到/aspects?actioningestProposal请求体为{proposal: ..., async: true/false}响应体经过extract_trace_data_from_mcps提取 TraceData若当前 emit 模式需要追踪ASYNC_WAIT且服务端支持 API tracing则调用_await_status阻塞直到写入确认。emit_mce(mce)则 POST 到/entities?actioningestpayload 结构为{entity: {value: {snapshot_fqn: mce_obj}}, systemMetadata: ...}rest_emitter.py并对超过INGEST_MAX_PAYLOAD_BYTES的负载打印警告。4.3 emit_mcps批量写入emit_mcps(mcps, emit_modeNone, wait_timeout...)用于批量提交多个 MCP返回List[TraceData]rest_emitter.py。它会按协议自动路由OpenAPI 路径_emit_openapi_mcpsrest_emitter.py先把 MCP 按(HTTP method, URL)分组再按字节大小INGEST_MAX_PAYLOAD_BYTES默认 15MB与条数BATCH_INGEST_MAX_PAYLOAD_LENGTH默认 200 条两个维度切分 chunk每个 chunk 只序列化一次后拼接为 JSON 数组发送RestLi 路径_emit_restli_mcpsrest_emitter.pyPOST 到/aspects?actioningestProposalBatch同样按大小与条数自动分块请求体为{proposals: [...], async: ...}。分块上限常量定义在 rest_emitter.pyINGEST_MAX_PAYLOAD_BYTES默认15 * 1024 * 102415MB环境变量DATAHUB_REST_EMITTER_BATCH_MAX_PAYLOAD_BYTES可调BATCH_INGEST_MAX_PAYLOAD_LENGTH默认 200环境变量DATAHUB_REST_EMITTER_BATCH_MAX_PAYLOAD_LENGTH可调。之所以 15MB 而非 GMS 的 16MB 上限源码注释解释为为请求头等开销预留空间限制条数则是为了避免单次请求处理时间过长导致 GMS 超时返回 500。批量发送还有一个细节_is_batch_asyncrest_emitter.py会在respect_mcp_sync_markerTrue时检查批次中是否含有 sync 标记的 MCP若有则整批升级为同步只会变得更同步、绝不会更异步。4.4 emit_usage使用统计已废弃emit_usage(usageStats)POST 到/usageStats?actionbatchIngest方法标注了deprecated(Use emit with a datasetUsageStatistics aspect instead)。五、Emit Mode四种写入一致性模式EmitMode枚举定义在 emit_mode.py它取代了旧的async_flag布尔开关把一致性/性能权衡显式化为四种模式模式行为一致性适用场景SYNC_WAIT同步写入主存储SQL并同步更新搜索存储Elasticsearch后才返回最强关键操作要求写入后立即可检索SYNC_PRIMARY同步写主存储SQL搜索索引异步更新中等需要实体立即可直接读取、搜索可稍延迟ASYNC入队后立即返回异步处理最终一致高吞吐、可接受最终一致ASYNC_WAIT异步入队但阻塞直到确认写入已持久化强且高吞吐需要持久化确认又不牺牲性能其中ASYNC与ASYNC_WAIT的is_async属性为TrueSYNC_WAIT/SYNC_PRIMARY为False。ASYNC_WAIT模式依赖 OpenAPI 协议与服务端的 API Tracing 能力_should_tracerest_emitter.py会校验_openapi_ingestion与服务端ServiceFeature.API_TRACING特性不满足时发出警告并降级。ASYNC_WAIT的确认机制由_await_status实现rest_emitter.py它轮询{gms}/openapi/v1/trace/write/{trace_id}?onlyIncludeErrorsfalsedetailedtrue检查每个 URN/aspect 的primaryStorage.writeStatus与searchStorage.writeStatus直到二者都不再是PENDING轮询采用指数退避初始 1 秒、翻倍增长、上限 5 分钟超过wait_timeout默认 1 小时抛TraceTimeoutError发现写入失败则抛TraceValidationError。对应地get_trace_status(trace, only_include_errors, detailed)允许手动查询某次写入的追踪状态。六、会话管理与请求发送底层所有请求都经由_session_config.build_session()构造的requests.Session发出rest_emitter.py统一带上如下基础请求头User-Agent格式为DataHub-Client/1.0 (mode; component/caller; version)其中component来自datahub_component参数或DATAHUB_COMPONENT_ENVX-DataHub-Client-Mode当前ClientModeX-DataHub-Py-Cli-VersionPython SDK 版本RestLi 场景额外带X-RestLi-Protocol-Version: 2.0.0与Content-Type: application/json。客户端证书与 CA 证书通过session.cert/session.verify注入disable_ssl_verification会把session.verify置为False。底层发送统一走_emit_generic(url, payload, methodPOST)rest_emitter.py关键行为发送前检查 payload 大小并告警调用response.raise_for_status()触发 HTTP 错误出错时解析 GMS 返回的 JSON抛OperationalError若报错信息含unrecognized field found but not allowed会附加提示服务端版本可能过旧若响应非 JSON则以原始错误信息构造OperationalError网络层异常RequestException统一包装为OperationalError。另外session属性rest_emitter.py被刻意设计为公开可访问SDK 中任何自定义 REST 调用如轮询 SDK 未封装端点都应通过该 session 发起否则会绕过session.auth中的 OAuth Token Provider 导致鉴权失败。七、连接校验与服务端能力协商test_connection()与fetch_server_config()rest_emitter.py会在首次使用时 GET{gms}/config做三件事校验服务类型响应的noCode必须为true否则提示你连到了 Frontend 而不是 GMS。Rest Emitter 应连接 DataHub GMS通常datahub-gms-host:8080或 Frontend 的 GMS API通常frontend:9002/api/gms加载服务端能力解析RestServiceConfig判断是否支持OPEN_API_SDK、API_TRACING等特性协商传输协议决定 OpenAPI 还是 RestLi 路径并缓存配置可用invalidate_config_cache()手动失效。fetch_server_config在收到 401 时会给出鉴权错误提示配置缓存可通过server_config_refresh_interval控制刷新频率。八、与其他 Emitter 的关系与选型DataHubRestEmitter并非唯一的 Emitter 实现仓库 emitter 目录 中还包含DatahubKafkaEmitterkafka_emitter.py写入 Kafka适用于大吞吐、对 Kafka 基础设施有依赖的场景CompositeEmittercomposite_emitter.py实验性组件把多个 Emitter 组合为一个向每个子 Emitter 依次转发且回调只绑定到第一个 EmitterSynchronizedFileEmittersynchronized_file_emitter.py落盘式发射器。选型建议元数据量不大、无需额外 Kafka 依赖时优先 Rest Emitter已有 Kafka 链路或需要削峰时用 Kafka Emitter需要同时写多个目标时用CompositeEmitter组合。九、常见问题与最佳实践9.1 连接地址选错报错connected to the frontend service instead of the GMS endpoint时请确认gms_server指向 GMS 本身如http://localhost:8080或 Frontend 的 GMS APIhttp://frontend:9002/api/gms而不是前端页面地址。9.2 超大负载单条 MCP 超过 15MB 时日志会告警likely fail to be emitted应拆分批次批量发送会自动按 15MB / 200 条分块无需手动处理。9.3 限流与重试生产环境可适当调大retry_max_times或通过环境变量DATAHUB_REST_EMITTER_429_RETRY_MULTIPLIER增强 429 场景的重试预算不建议把超时设为低于 1 秒源码会对低于_TIMEOUT_LOWER_BOUND_SEC 1的配置打印警告。9.4 空闲连接 SSL 错误长期运行的进程建议保持tcp_keepaliveTrue默认由DATAHUB_REST_SINK_DEFAULT_TCP_KEEPALIVE控制避免空闲连接被回收后产生SSLEOFError。9.5 使用with上下文DataHubRestEmitter实现了上下文协议Closeable推荐用with DataHubRestEmitter(...) as emitter:确保退出时调用close()释放连接池。十、深入阅读API 文档骨架docs-website/sphinx/apidocs/clients/rest-emitter.rst核心实现metadata-ingestion/src/datahub/emitter/rest_emitter.pyEmit 模式定义metadata-ingestion/src/datahub/emitter/emit_mode.pyMCP 封装metadata-ingestion/src/datahub/emitter/mcp.pyEmitter 抽象metadata-ingestion/src/datahub/emitter/generic_emitter.py组合发射器metadata-ingestion/src/datahub/emitter/composite_emitter.py可运行示例metadata-ingestion/examples/library/dataset_replace_properties.py、metadata-ingestion/examples/structured_properties/update_structured_property.py【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表