ARTICLE DETAIL

资讯详情

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

DataHub Athena 元数据摄取指南:从配置、IAM 权限到血缘与分区剖析

DataHub Athena 元数据摄取指南:从配置、IAM 权限到血缘与分区剖析 DataHub Athena 元数据摄取指南从配置、IAM 权限到血缘与分区剖析【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本文围绕 DataHub 提供的 Athena 元数据摄取模块athenasource展开说明如何将 AWS Athena 中的表、视图、Schema 字段、容器等核心元数据同步进 DataHub并支持表级/字段级血缘、数据剖析Profiling与有状态删除检测。读完本文你将掌握 Athena 连接器的完整 recipe 配置、所需 IAM 权限清单、各配置项的含义与默认值以及源码级的分区提取、复杂类型映射与血缘拼接原理。一、概述DataHub 如何对接 AthenaAthena 是 AWS 上用于存储与查询分析型/操作型数据的 Serverless 数据平台。DataHub 的 Athena 集成metadata-ingestion中的datahub.ingestion.source.sql.athena覆盖以下核心元数据实体数据集实体表Table、视图View以及构成这些实体的Schema 字段列容器实体Athena 中的数据库/Schema 被建模为 DataHub 的 Database/Container数据血缘表级与字段级 Lineage数据剖析可选的 Profiling有状态删除检测Stateful Deletion识别源端已删除的实体。从源码能力声明看athena.py该模块当前为GA 支持状态默认启用平台实例Platform Instance与描述Descriptions摄取支持通过配置开启数据剖析与血缘。二、概念映射Athena 概念到 DataHub 概念的对应关系原文档指出具体映射细节仍可能随版本演进但通用的概念映射关系如下表所示源概念Source ConceptDataHub 概念说明平台 / 账号 / 项目范围Platform Instance、Container在平台上下文中组织资产。核心技术资产如表 / 视图 / Topic / 文件Dataset主要的被摄取技术资产。Schema 字段 / 列SchemaField在支持 Schema 提取时包含。所有与协作相关的参与者CorpUser、CorpGroup由支持所有者和身份元数据的模块发出。依赖关系与处理关系Lineage edges在支持并启用血缘提取时可用。在 Athena 场景下容器层级有特殊实现由于 Athena 中Schema 即数据库不存在独立的 database 层源码中gen_database_containers直接返回空、gen_schema_containers将 Schema 容器作为 Database 容器生成见 athena.py。这一点对理解 DataHub 界面中看到的容器树形结构很重要。三、前置条件IAM 权限配置摄取前需要创建一个 IAM 策略并附加到 recipe 中使用的 AWS 角色或凭证上。以下策略覆盖 Athena 元数据读取、Glue Catalog 查询、查询结果读写与 S3 访问所需的全部权限策略模板来自 athena_pre.md{ Version: 2012-10-17, Statement: [ { Sid: VisualEditor0, Effect: Allow, Action: [ athena:GetTableMetadata, athena:StartQueryExecution, athena:GetQueryResults, athena:GetDatabase, athena:GetDataCatalog, athena:ListQueryExecutions, athena:GetWorkGroup, athena:StopQueryExecution, athena:GetQueryResultsStream, athena:ListDatabases, athena:GetQueryExecution, athena:ListTableMetadata, athena:BatchGetQueryExecution, glue:GetTables, glue:GetDatabases, glue:GetTable, glue:GetDatabase, glue:SearchTables, glue:GetTableVersions, glue:GetTableVersion, glue:GetPartition, glue:GetPartitions, s3:GetObject, s3:ListBucket, s3:GetBucketLocation ], Resource: [ arn:aws:athena:${region-id}:${account-id}:datacatalog/*, arn:aws:athena:${region-id}:${account-id}:workgroup/*, arn:aws:glue:${region-id}:${account-id}:tableVersion/*/*/*, arn:aws:glue:${region-id}:${account-id}:table/*/*, arn:aws:glue:${region-id}:${account-id}:catalog, arn:aws:glue:${region-id}:${account-id}:database/*, arn:aws:s3:::${datasets-bucket}, arn:aws:s3:::${datasets-bucket}/* ] }, { Sid: VisualEditor1, Effect: Allow, Action: [ s3:PutObject, s3:GetObject, s3:ListBucketMultipartUploads, s3:AbortMultipartUpload, s3:ListBucket, s3:GetBucketLocation, s3:ListMultipartUploadParts ], Resource: [ arn:aws:s3:::${athena-query-result-bucket}/*, arn:aws:s3:::${athena-query-result-bucket} ] }, { Sid: VisualEditor2, Effect: Allow, Action: [athena:ListDataCatalogs], Resource: [*] } ] }部署时请将${region-id}、${account-id}、${datasets-bucket}、${athena-query-result-bucket}替换为实际值。其中athena:StartQueryExecution、athena:GetQueryResults等查询类权限用于分区提取与数据剖析执行 SQLs3:PutObject等权限用于将 DataHub 执行的 Athena 查询结果写入query_result_location指向的 S3 桶。四、最小可用 Recipe 与连接原理一个最小的 Athena 摄取 recipe 如下来自 athena_recipe.ymlsource: type: athena config: # AWS Keys可选——仅当本地未配置 AWS 凭证时需要 username: my_aws_access_key_id password: my_aws_secret_access_key # 坐标信息 aws_region: my_aws_region work_group: primary # 选项 query_result_location: s3://my_staging_athena_results_bucket/results/ sink: # sink configs4.1 连接串是如何拼出来的源码中AthenaConfig.get_sql_alchemy_url()athena.py会基于配置生成形如awsathenarest://...的 SQLAlchemy URL并将以下配置项作为 URI 选项透传给 PyAthenas3_staging_dir←query_result_location数据剖析与分区查询的结果落地位置work_group使用的 Athena 工作组catalog_name目标 Data Catalogrole_arn与duration_seconds需要切换角色时使用的 ARN 与有效期。连接建立后get_inspectors()会创建一个CustomAthenaRestDialect实例并注入 inspector用自定义方言来修正 PyAthena 默认行为见下文「分区与复杂类型」。五、配置参数详解下表汇总AthenaConfig的完整配置项字段定义与默认值均来自 athena.py配置项必填默认值说明aws_region是无Athena 数据库所在的 AWS 区域。work_group是无Amazon Athena 工作组名称。query_result_location是无DataHub 执行的 Athena 查询结果要写入的 S3 路径。username否自动检测访问密钥 ID不填则按 boto3 标准规则检测凭证。password否自动检测秘密访问密钥检测规则同上。database否自动检测指定要摄取的 Athena 数据库未设置时自动检测全部 Schema。aws_role_arn否无供 PyAthena 在连接时扮演的 AWS 角色 ARN。aws_role_assumption_duration否3600角色扮演时长秒最大4320012 小时。catalog_name否awsdatacatalogAthena Catalog 名称。当值为 S3 Tables catalog以s3tablescatalog/开头时会自动将其用作platform_instance见 6.3 节。glue_platform_instance否无血缘上游 Glue 数据集的平台实例。仅当你的 Glue 连接器配置了非默认 platform_instance 时才需设置。iceberg_platform_instance否无血缘上游 Iceberg 数据集的平台实例对 S3 Tables catalog 会默认取 catalog 名以保证 URN 拼接。extract_partitions否true是否提取表分区。分区提取会对表执行select * from table$partitions查询若不想授予 SELECT 权限可关闭。extract_partitions_using_create_statements否false是否改用SHOW CREATE TABLE提取分区实验特性。Iceberg 分区需要开启失败时自动回退到默认提取方式。emit_schema_fieldpaths_as_v1否false将简单列路径转换为 DataHub field path v1 格式不含嵌套字段的列适用。profiling否见下数据剖析配置继承通用GEProfilingConfig。s3_staging_dir否无已废弃改用query_result_location配置后会自动迁移并打印警告pydantic_renamed_field。5.1 Profiling 配置要点Athena 的剖析配置AthenaProfilingConfig继承自通用剖析配置并覆写了一个关键默认值athena.pypartition_profiling_enabled默认false通用配置默认开启此处特意关闭。开启后将对表的最新分区执行剖析。需要特别注意的是数据剖析会对整表执行 SQL 查询成本可能较高。能力声明中对此的说明是可选通过配置开启剖析对整表执行 SQL属于代价较高的操作athena.py。生产环境建议先在非核心表上验证开销再全量启用。5.2 Schema 过滤行为由于连接串中的 database/schema 过滤在 Athena 场景下不生效源码覆写了get_schema_names()athena.py当配置了database时仅摄取与该名称完全匹配的 Schema。六、分区提取与复杂类型映射源码原理6.1 分区提取的两种路径get_partitions()athena.py依据配置选择提取方式默认路径SQLAlchemy通过 PyAthena 的get_table_metadata返回的partition_keys获取分区列名列表SHOW CREATE TABLE路径当开启extract_partitions_using_create_statements时先执行SHOW CREATE TABLE再交给AthenaPropertiesExtractorathena_properties_extractor.py解析。该解析器基于 sqlglot 的 Athena 方言能够识别简单分区列PARTITIONED BY (ds string)时间变换day、month、year、hour函数Iceberg 变换bucket(n, col)、truncate(len, col)等。若SHOW CREATE TABLE解析失败如视图不支持该语句会打印警告并自动回退到 SQLAlchemy 路径。开启分区剖析profiling.partition_profiling_enabledtrue时还会额外执行最大分区查询select partitions from schema.table$partitions where ...并将结果缓存到table_partition_cache供后续生成仅剖析最新分区的查询使用generate_partition_profiler_queryathena.py。6.2 复杂数据类型的深度支持PyAthena 默认无法识别array、map、struct等复杂类型CustomAthenaRestDialect._get_column_type()athena.py覆写了类型推断arrayT/listT→ SQLAlchemyARRAYstruct.../record/row→sqlalchemy_bigquery.STRUCT通过 Hive DDL 解析工具递归展开嵌套字段mapK,V→ 自定义MapTypeTupleType的包装key/value 分别递归推断decimal(p,s)→DECIMAL解析精度与标度还兼容了剖析管线中可能泄漏的 Pandas nullable dtype如Int64Dtype→BIGINT。顶层通过register_custom_type(STRUCT, RecordTypeClass)、register_custom_type(MapType, MapTypeClass)将这些类型映射为 DataHub SchemaField 的复杂类型从而保证嵌套列在 DataHub 中以正确的层级展示。6.3 表级属性与血缘拼接get_table_properties()athena.py从 Athena Table Metadata 中提取描述comment、分区键列表、create_time、last_access_time、table_type以及所有 TBLProperties 参数写入 Dataset 的自定义属性custom properties。上游血缘的拼接逻辑集中在_get_upstream_location()athena.py按 Catalog 类型分四种情况S3 Tables Catalogcatalog_name以s3tablescatalog/开头表一定是 Iceberg 格式直接生成iceberg平台的 Dataset URN避免把 Athena 内部存储路径当作血缘Glue Catalog默认awsdatacatalog或经ListDataCatalogs判定为GLUE生成glue平台的上游 URN并与你的 Glue 摄取配置相互关联非 Glue 但表类型为 ICEBERG生成iceberg平台上游 URN其他情况读取 TBLProperties 中的location生成s3://...对应的 S3 Dataset URN非s3://的 location 会跳过血缘并记入报告。如果is_glue_catalog无法判定Catalog 类型未知且名称不是awsdatacatalog会跳过 Glue/Iceberg 血缘并在报告中告警避免静默丢失。视图血缘则基于SHOW CREATE VIEW提取的定义做字段级/表级拼接对应能力声明中的LINEAGE_COARSE与LINEAGE_FINE见 athena.py。6.4 标识符安全防护源码对分区查询、最大分区过滤等动态 SQL 中的所有标识符执行白名单校验_sanitize_identifier仅允许字母、数字、下划线与点号并对分区值中的单引号做转义athena.py。若遇到不安全的标识符会跳过该表的剖析而不是让整个摄取任务崩溃。七、能力清单速览结合能力装饰器athena.pyAthena 模块当前支持的能力如下能力状态说明Platform Instance默认启用支持多环境/多账号资产隔离。Domains通过domain配置支持给数据集打域标签。Data Profiling可选开启对整表执行 SQL代价较高谨慎启用。表级血缘Lineage Coarse支持上游指向 Glue、Iceberg 或 S3 实体依 Catalog 类型而定覆盖表与视图。字段级血缘Lineage Fine支持同上覆盖表与视图。Descriptions默认启用表/列描述自动同步。八、测试验证仓库内置了 Athena 连接器的集成测试 test_athena_source.py通过 mock inspector 模拟 Schema、表、视图、复杂列MapType、STRUCT、ARRAY与表属性执行完整Pipeline.create(...).run()流程并将输出与黄金文件 athena_mce_golden.json 比对。它验证了表与视图含视图间级联引用被正确摄取嵌套复杂类型能生成正确的 SchemaField表描述、自定义属性与 S3 上游血缘 URN 正确落盘。如需在本机复现可在metadata-ingestion目录下运行pytest tests/integration/athena/test_athena_source.py按仓库开发文档配置环境后执行。九、限制与排障受源端 API、权限与平台暴露的元数据约束模块行为存在以下边界数据剖析成本较高剖析查询作用于整表务必评估成本后启用分区剖析默认关闭需显式打开partition_profiling_enabledIceberg 分区提取为实验特性extract_partitions_using_create_statements标记为 experimental失败时会回退默认路径非 Glue/Iceberg 血缘仅支持s3://location其余 location 会被跳过并产生告警SHOW CREATE TABLE不适用于视图解析失败属于预期行为会自动回退。故障排查建议按以下顺序进行来自 athena_post.md校验凭证与权限确认 IAM 策略已附加、角色可被扮演、凭证在 boto3 检测规则下可见校验连通性确认网络可达 Athena 端点、S3 查询结果桶可读写校验范围过滤确认database、catalog_name、platform_instance等设置与实际环境一致查看摄取日志关注report.warning中暴露的源端错误如 Catalog 类型无法判定、S3 Tables 列表失败、血缘 location 不支持等据此调整配置。十、总结DataHub 的 Athena 连接器通过 PyAthena 自定义 SQLAlchemy 方言实现了对 Athena 表、视图、Schema 容器、复杂类型字段、分区信息的完整摄取并基于 Catalog 类型自动选择 Glue、Iceberg 或 S3 上游血缘目标配合可选的整表/分区剖析与有状态删除检测构成一条生产可用的 Athena 元数据同步链路。开始使用前请务必按 athena_pre.md 配好 IAM 策略并按本文第五节核对query_result_location、work_group、catalog_name三个关键坐标即可用文中的 recipe 快速跑通第一条摄取管线。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表