十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

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

DataHub Elasticsearch 数据源接入指南:索引元数据、Schema 字段类型与 Index Template 的完整摄取方案

DataHub Elasticsearch 数据源接入指南:索引元数据、Schema 字段类型与 Index Template 的完整摄取方案 DataHub Elasticsearch 数据源接入指南索引元数据、Schema 字段类型与 Index Template 的完整摄取方案【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubElasticsearch 是 DataHub 元数据生态中的常用平台之一elasticsearch数据源模块负责把 Elasticsearch 集群中的索引Index、索引字段类型Column Types以及可选的 Index Template 摄取为 DataHub 中统一的 Dataset 实体供数据发现、治理与检索使用。本文基于该模块的官方文档与源码实现elastic_search.py完整讲解其能力边界、配置参数、Recipe 示例、底层摄取流程与故障排查方法读者阅读后可独立完成 Elasticsearch 到 DataHub 的生产级元数据摄取配置。模块概览与核心能力elasticsearch模块是 DataHub 针对 Elasticsearch 平台的元数据集成官方将其定位为面向元数据实体的实用型集成。根据模块说明文档metadata-ingestion/docs/sources/elasticsearch/elasticsearch_pre.md该插件面向生产环境摄取工作流具体提取以下两类信息索引的元数据Metadata for indexes将每个被允许的索引映射为 DataHub 中的一个 Dataset 实体索引字段关联的列类型Column types associated with each index field将索引 Mapping 中的每个字段转换为 DataHub SchemaField并携带 DataHub 标准化的 Schema 类型。从源码装饰器elastic_search.py可以看到该模块的官方支持状态与能力标注平台名称为Elasticsearch支持状态为GAGeneral Availability即正式可用已声明能力PLATFORM_INSTANCE平台实例默认启用。此外模块说明metadata-ingestion/docs/sources/elasticsearch/README.md还指出该集成支持有状态删除检测stateful deletion detection——摄取器会记录已扫描的实体状态从而在重跑时识别并移除源端已不存在的实体。该机制在源码中通过继承StatefulIngestionSourceBase与StaleEntityRemovalSourceReport实现elastic_search.py。概念映射在 DataHub 的通用元数据模型下Elasticsearch 源概念与 DataHub 实体之间的映射关系如下源自 README.md源概念DataHub 概念说明平台 / 账户 / 项目范围Platform Instance、Container在平台上下文中组织资产核心技术资产如表 / 视图 / Topic / 文件Dataset主要摄取的实体此处即索引与索引模板Schema 字段 / 列SchemaField支持 Schema 提取时包含所有权与协作主体CorpUser、CorpGroup由支持所有权与身份元数据的模块发出依赖与处理关系Lineage 边当血缘提取受支持并启用时可用前置条件运行摄取前需满足以下前提elasticsearch_pre.md网络连通性确保执行摄取的主机可以访问 Elasticsearch 集群的 HTTP 端口默认9200有效认证凭据具备用户名 / 密码或 API Key元数据 API 读取权限拥有读取该模块所需元数据 API 的权限包括索引别名、索引 Mapping、模板等接口。从源码看模块通过opensearchpyOpenSearch 官方 Python 客户端兼容 Elasticsearch建立连接底层依赖的接口主要包括indices.get_alias()、indices.get()、indices.get_template()、indices.get_index_template()与cat.indices()elastic_search.py。因此只读权限如viewer/read角色通常即可满足摄取需要。完整配置参数解析模块的配置模型定义在 elastic_search.py 的ElasticsearchSourceConfig中官方示例 Recipe 见 elasticsearch_recipe.yml。各参数说明如下连接坐标与认证参数默认值说明hostlocalhost:9200Elasticsearch 集群地址支持host:port格式也支持逗号分隔的多地址。源码中的host_colon_port_comma校验器会移除协议前缀http:///https://与末尾/并逐一校验端口格式username无基本认证用户名可选password无基本认证密码可选属于敏感配置使用TransparentSecretStr类型存储api_key无API Key 认证接受两种形式由(id, api_key)组成的元组或已编码的 Base64 字符串。源码_api_key_authorizationelastic_search.py会将其组装为Authorization: ApiKey token请求头因为 opensearchpy 客户端没有api_key参数且会静默丢弃未知关键字说明http_auth属性elastic_search.py仅在username非空时才启用基本认证若只配置了api_key则通过请求头传递。二者可按需选用其一。SSL 与 TLS 配置参数默认值说明use_sslFalse是否启用 SSL 连接verify_certsFalse是否校验 SSL 证书ca_certs无CA 根证书文件路径如./path/ca.certclient_cert无客户端证书文件路径若与私钥分离则仅填证书否则为同时包含私钥与证书的文件client_key无客户端私钥文件路径当证书与私钥分离时使用ssl_assert_hostnameFalse是否进行主机名校验ssl_assert_fingerprint无校验提供的证书指纹如./path/cert.fingerprint作用域与过滤参数默认值说明url_prefix企业多集群场景下若所有集群共用同一端点并通过 URL 前缀路由则在此填写前缀envPRODDataset 所属环境EnvConfigMixin 提供会写入 Dataset URN 的 env 部分platform_instance无平台实例名PlatformInstanceConfigMixin 提供启用后会额外发出DataPlatformInstanceClass方面elastic_search.pyindex_pattern.allow[.*]索引包含正则决定摄取哪些索引index_pattern.deny[^_.*, ^ilm-history.*]索引排除正则默认排除系统隐藏索引_开头与 ILM 历史索引ilm-history.*ingest_index_templatesFalse是否同时摄取 Index Template 作为 Datasetindex_template_pattern.allow[.*]模板包含正则index_template_pattern.deny[^_.*]模板排除正则默认排除以_开头的模板index_pattern与index_template_pattern均为AllowDenyPattern类型规则为先匹配 allow 列表命中后若同时命中 deny 列表则拒绝默认值从源码elastic_search.py可直接确认。数据概要Profiling与 URN 折叠参数默认值说明profiling.enabledFalse是否启用数据概要摄取。启用后通过cat.indices获取每个索引的文档数docs.count与存储大小store.size发出DatasetProfileClass行数、列数、字节大小见 elastic_search.pyprofiling.operation_config默认 OperationConfig概要运行的实验性操作配置。注意若启用lower_freq_profile_enabled等参数组合不正确配置校验会直接报错见单元测试 test_elasticsearch_source.pycollapse_urns.urns_suffix_regex[]从索引名中剥离后缀的正则列表用于将时间分区索引如log-2025-01-01、metrics-1755000000折叠合并为同一个 Dataset。正则按顺序依次应用当日志索引同时存在-YYYY-MM-DD与-epochtime两种后缀格式时需要配置多个正则源码中的collapse_name/collapse_urnelastic_search.py会在生成 URN、SchemaMetadata 名称以及 Profiling 匹配时统一应用折叠逻辑保证同一逻辑数据集被视作一个实体。Recipe 配置示例以下为官方提供的完整 Recipe源自 elasticsearch_recipe.yml可直接作为基础模板使用source: type: elasticsearch config: # 连接坐标 host: localhost:9200 # 凭据 username: user # 可选 password: pass # 可选 # SSL 支持 use_ssl: False verify_certs: False ca_certs: ./path/ca.cert client_cert: ./path/client.cert client_key: ./path/client.key ssl_assert_hostname: False ssl_assert_fingerprint: ./path/cert.fingerprint # 选项 url_prefix: # 可选 url_prefix env: PROD index_pattern: allow: [.*some_index_name_pattern*] deny: [.*skip_index_name_pattern*] ingest_index_templates: False index_template_pattern: allow: [.*some_index_template_name_pattern*] sink: # sink 配置运行摄取的标准命令DataHub CLI 方式# 使用 datahub CLI 执行摄取 datahub ingest -c elasticsearch_recipe.yml底层摄取流程源码级实现原理理解摄取流程有助于针对性地调整配置与排查问题。以下流程来自ElasticsearchSource的get_workunits_internal与_extract_mcps实现elastic_search.py。步骤 1枚举索引并按模式过滤摄取首先调用client.indices.get_alias()获取集群中全部索引含别名对每个索引命中index_pattern的 allow 且未命中 deny 的索引进入后续处理未命中的索引调用report_dropped记录到报告LossyList保存避免内存膨胀供后续查看被过滤明细每个索引调用report_index_scanned记录扫描计数。步骤 2Schema 提取与类型映射对每个索引摄取indices.get(index...)返回的 Mapping通过ElasticToSchemaFieldConverter转换为 DataHub SchemaField。该转换器elastic_search.py维护了 Elasticsearch 原生类型到 DataHub 标准类型的映射表Elasticsearch 类型DataHub SchemaFieldDataTypeElasticsearch 类型DataHub SchemaFieldDataTypebooleanBooleanTypeClasskeyword/constant_keyword/wildcardStringTypeClassbinaryBytesTypeClasstext/match_only_text/completion/search_as_you_type/ipStringTypeClassbyte/integer/long/short/double/float/half_float/scaled_float/unsigned_long/token_countNumberTypeClassobject/flattened/nested/geo_pointRecordTypeClassdate/date_nanosDateTypeClasshistogram/aggregate_metric_doubleArrayTypeClass未在映射表中的类型会记录警告并使用NullTypeClass兜底elastic_search.py。字段路径采用版本化格式例如[version2.0].[typekeyword].browserId [version2.0].[typetext].actorUrn [version2.0].[typelong].height嵌套对象含properties的字段会以[typeproperties].字段名形式递归展开为层级路径_get_schema_fields的实现elastic_search.py会为嵌套父字段生成 Record 类型的 SchemaField再递归处理其子字段若字段既无type也无properties则记录警告并跳过。字段路径的唯一性由单元测试保证见 test_elasticsearch_source.py。每个索引最终发出以下方面aspectSchemaMetadataschemaName为索引名折叠后platform为elasticsearch平台 URNhash为 Mapping JSON 的 MD5platformSchema.rawSchema保留完整原始 Mapping便于回溯Status标记实体为未删除removedFalseSubTypes索引标记为ElasticIndex、数据流标记为ElasticDatastream、模板标记为ElasticIndexTemplateDatasetProperties提取aliases别名列表、index_patterns匹配模式、num_shards、num_replicas等自定义属性DataPlatformInstance仅配置platform_instance时。步骤 3数据流Data Stream去重当索引元数据中包含data_stream字段时该索引属于某个数据流。源码会按数据流名称累计data_stream_partition_count同一数据流的第二个及后续分区索引跳过重复处理最终为每个数据流额外发出携带numPartitions自定义属性的DatasetProperties方面elastic_search.py便于统计分区数量。步骤 4Index Template 摄取可选当ingest_index_templates: True时模块同时处理两类模板Legacy旧式模板通过indices.get_template()获取Composable可组合模板通过indices.get_index_template()获取ES 7.8 / OpenSearch。若该接口调用失败如版本过旧仅记录警告不影响整体流程elastic_search.py。模板同样受index_template_pattern过滤Mapping 提取位置有所区别可组合模板的 Mapping 位于template.mappings之下elastic_search.py。模板的 Dataset 属性aliases、index_patterns、num_shards、num_replicas由_extract_template_custom_properties提取elastic_search.py。步骤 5数据概要可选启用profiling.enabled后模块通过cat.indices一次性拉取全部索引的docs.count与store.size字节按折叠后的索引名匹配汇总发出DatasetProfileClasstimestampMillis、rowCount、columnCount、sizeInBytes并消费掉已匹配条目避免重复统计elastic_search.py。步骤 6有状态删除检测由于模块继承自StatefulIngestionSourceBase摄取器会持久化已处理实体的状态再次运行时源端已删除的索引对应的 DataHub 实体会被标记删除从而保持两侧元数据一致。相关报告类ElasticsearchSourceReport同时维护index_scanned扫描计数与filtered被过滤索引明细可在摄取报告中查看elastic_search.py。能力边界Limitations模块行为受源平台 API、权限配置以及平台暴露的元数据范围约束elasticsearch_post.md。结合源码需要注意以下边界不摄取文档级数据模块只摄取索引结构Mapping与概要统计不会拉取索引中的业务文档内容不提取血缘与所有权当前实现未发出 Lineage 边或 Ownership 元数据对应 README.md 概念映射表中仅当支持时启用的说明系统索引默认被排除默认 deny 规则会过滤_开头如.kibana与ilm-history.*索引如需摄取需显式调整index_pattern权限不足时 API 调用受限若账号缺少读取别名 / 模板 / cat API 的权限对应能力将静默缺失或失败能力确认以Important Capabilities表为准每个数据源文档页上方的能力表Important Capabilities是判断某项功能是否受支持、是否需要额外配置的唯一权威来源本文所列能力与上述源码证据一致。故障排查Troubleshooting按官方文档建议elasticsearch_post.md当摄取失败时遵循以下排查顺序验证凭据与权限确认username/password或api_key正确且账号具备读取索引 Mapping、别名、模板及若启用 Profilingcat API 的权限验证网络连通性确认主机可访问host指定的地址与端口若启用了 SSL 请核对use_ssl、ca_certs、client_cert、client_key、ssl_assert_hostname、ssl_assert_fingerprint是否与集群 TLS 配置一致核对作用域过滤确认index_pattern/index_template_pattern的 allow / deny 正则是否意外过滤掉了目标索引^_.*等默认规则尤其容易误伤检查摄取日志结合日志中的 source-specific 错误信息定位问题。常见可依据日志的现象包括Cannot map type to SchemaFieldDataType警告说明存在未映射的 Elasticsearch 类型已用 NullTypeClass 兜底不影响整体摄取Unable to fetch composable index templates警告说明目标版本不支持可组合模板接口配置校验异常如profiling.operation_config参数错误在摄取启动阶段即被 pydantic 拒绝对应测试 test_elasticsearch_source.py。调整配置后重新执行datahub ingest -c recipe.yml即可若涉及 Schema 变更也可借助 DataHub 的删除 / 重跑流程配合有状态删除检测重建元数据。延伸阅读模块源码与全部配置模型metadata-ingestion/src/datahub/ingestion/source/elastic_search.py官方 Recipe 模板metadata-ingestion/docs/sources/elasticsearch/elasticsearch_recipe.yml模块概览与概念映射metadata-ingestion/docs/sources/elasticsearch/README.md摄取前置条件说明metadata-ingestion/docs/sources/elasticsearch/elasticsearch_pre.md单元测试类型映射、URN 折叠、配置校验metadata-ingestion/tests/unit/test_elasticsearch_source.pyDataHub 摄取通用入门metadata-ingestion/README.md【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表