DataHub Excel 连接器(Source Connector)实战指南:从工作表到数据集的元数据摄取、Schema 推断与 Profiling
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
导读
本文以 DataHub 仓库中 Excel 元数据摄取连接器(Source Connector)为对象,系统讲解如何将 Excel 工作簿中的工作表(Worksheet)作为 Dataset 摄取进 DataHub,涵盖支持的文件类型、三种数据来源(本地文件系统、AWS S3、Azure Blob Storage)、工作表识别与表头探测算法、Schema 推断、数据 Profiling 以及有状态删除检测等能力。读完本文,你将能够编写可运行的 Ingestion Recipe、理解每个配置项的行为边界,并掌握该连接器底层(metadata-ingestion/src/datahub/ingestion/source/excel/)的核心实现原理。
一、连接器概览:Excel 与 DataHub 的概念映射
DataHub 的 Excel 连接器将 Excel 文件中的元数据摄入 DataHub,覆盖文件/湖仓元数据实体(数据集、路径、容器),并支持数据画像(Data Profiling)与有状态删除检测(Stateful Deletion Detection)。官方文档的说明位于 metadata-ingestion/docs/sources/excel/README.md。
两者的实体映射关系如下表:
| Excel 实体 | DataHub 实体 | 说明 |
|---|---|---|
| Excel 工作表(Worksheet) | Dataset | 每个工作表成为一个 Dataset,URN 模式为urn:li:dataset:(urn:li:dataPlatform:excel,{path}/[{filename}]{sheet_name},PROD) |
| 文件/目录结构 | Container | 目录层级会创建 URN 被混淆处理的 Container,用于组织数据集 |
需要注意的是:Excel 工作簿(Workbook 文件)本身不会成为独立的 DataHub 实体——只有工作簿内部的各个工作表会被作为 Dataset 摄取。这一设计与源码实现一致:在 source.py 中,process_file遍历工作簿内的每个工作表生成 Dataset,而 Container 层级则交由ContainerWUCreator依据文件的目录相对路径创建。
在 source.py 的类装饰器中声明了该连接器的能力清单与支持状态:
@support_status(SupportStatus.ALPHA):连接器当前处于 Alpha 支持阶段;@capability(SourceCapability.CONTAINERS, "Enabled by default"):容器(目录层级)能力默认开启;@capability(SourceCapability.SCHEMA_METADATA, "Enabled by default"):Schema 元数据默认开启;@capability(SourceCapability.DATA_PROFILING, "Optionally enabled via configuration"):数据画像需通过配置启用;@capability(SourceCapability.DELETION_DETECTION, "Optionally enabled via stateful_ingestion.remove_stale_metadata"):删除检测需开启有状态摄取。
二、支持的文件类型与工作表识别机制
2.1 支持的文件类型
根据 excel_pre.md 的说明,连接器支持以下文件类型:
- Excel 工作簿(
*.xlsx) - 启用宏的 Excel 工作簿(
*.xlsm)
在源码层面,source.py 中定义了更完整的扩展名白名单ALLOWED_EXTENSIONS = [".xlsx", ".xlsm", ".xltx", ".xltm"],is_excel_file方法据此过滤文件;同时check_file_is_valid还会先用path_pattern(AllowDenyPattern)对文件路径进行正则过滤。
2.2 工作表的"表格"识别机制
连接器会尝试识别工作表中哪些单元格构成"表格数据"。一张表的定义是:一个用于派生列名的表头行(header row)加上其后的数据行(data rows)。Schema 则根据列中实际出现的数据类型推断得出。
这一逻辑实现在 excel_file.py 的ExcelFile类中,其核心流程为:
get_tables():遍历工作表(或仅活动工作表),对每个 sheet 调用get_table();find_header_row():采用**打分制(scoring)**寻找最可能的表头行。算法先跳过工作表开头连续的空行,然后对每个候选行计算得分,综合以下多个维度的启发式信号:_score_row_with_numeric_cells:含数值单元格的行会被扣分(表头通常不是数值);_score_non_empty_cells:比前一行非空单元格更多的行得分更高;_score_header_like_text:匹配^[A-Z][a-zA-Z\s]*$模式的文本(如Name、Sales Amount)每个加 1 分;_score_text_followed_by_numeric:文本型表头后紧跟数值型数据行的模式加分(最高贡献 6 分以上);_score_column_type_consistency:后续多行列类型保持一致的候选行加分;_score_metadata_patterns:对疑似元数据行进行扣分惩罚;- 若最终最高得分
max_score <= 0,则判定该工作表不包含表格(对应get_tables中的report_worksheet_dropped与 "Worksheet does not contain a table" 警告);
find_footer_start():从表头后的数据区继续扫描,识别表格的页脚(footer)边界,判定依据包括:空行后跟内容显著变少或长文本的行、出现total/sum/average/mean/source/note/footnote等页脚指示词、单行填充单元格数远低于数据行平均值、长文本单元格(>50 字符)以及列数据类型大面积不一致等;- 表头行中末尾空列会被截断,空列名自动命名为
Unnamed_{i},重名列则追加序号(如Name_1)以保证列名唯一; - 最终以表头行 + 数据行构建 pandas DataFrame,
ExcelTable数据类(df、header_row、footer_row、row_count、column_count、metadata、sheet_name)成为后续 Dataset 与 Schema 生成的输入。
2.3 元数据的提取与自定义属性
紧邻表格上方或下方、且只有前两列有值的行会被认为是元数据行:第一列为键(key),第二列为值(value)。这类行会转换为 Dataset 的自定义属性(custom properties)——键中结尾的:或=会被去除。此外,工作簿自身的标准属性(standard properties)和自定义属性(custom properties)也会作为 Dataset 自定义属性一并导入。
这一点在源码中对应 excel_file.py 的extract_metadata()与read_excel_properties():后者从openpyxl的wb.properties(DocumentProperties)读取title、author、subject、description、keywords、category、last_modified_by、created、modified、status、revision、version、language、identifier等标准属性,并通过wb.custom_doc_props读取自定义属性;与已存在键冲突的自定义属性会以custom.{name}前缀命名。头尾元数据中的键若与已有属性冲突,则追加_1后缀。
三、三种数据来源与配置
工作簿可以从本地文件系统、S3 桶、Azure Blob Storage三个位置摄取。path_list中支持使用通配符*代替一个目录层级或作为文件名的一部分,从而用一条路径规则匹配多个目录或文件。
URI 类型判定在 source.py 的uri_type()静态方法中完成:http/https且 host 含.blob.core.windows.net判为 Azure Blob(ABS),s3/s3a判为 S3,file:///与绝对/相对路径判为本地文件。
3.1 本地文件
source: type: excel config: path_list: - "/data/path/reporting/excel/*.xlsx" profiling: enabled: false本地文件通过local_browser()(基于glob.glob递归匹配)发现文件,get_local_file()读取为二进制流。metadata-ingestion/docs/sources/excel/excel_recipe.yml给出了同样结构的完整示例(含sink占位)。
3.2 AWS S3
source: type: excel config: path_list: - "s3://bucket/data/excel/*/*.xlsx" aws_config: aws_access_key_id: ... aws_secret_access_key: ... aws_region: us-east-1 profiling: enabled: falseS3 场景下,s3_browser()使用aws_config.get_s3_resource()以指定前缀分页(每页 1000 个对象)扫描桶内对象;get_s3_file()负责拉取对象内容。若设置了use_s3_bucket_tags或use_s3_object_tags,还会将桶/对象标签转换为 DataHub 的 GlobalTags。
S3 前置权限要求
在 excel_pre.md 中明确给出了 S3 摄取的最小权限 IAM 策略:
{ "Version": "2012-10-17", "Statement": [ { "Sid": "VisualEditor0", "Effect": "Allow", "Action": ["s3:ListBucket", "s3:GetBucketLocation", "s3:GetObject"], "Resource": [ "arn:aws:s3:::your-bucket-name", "arn:aws:s3:::your-bucket-name/*" ] } ] }权限说明:
s3:ListBucket:允许列出桶中的对象,是 S3 摄取源了解可读取对象清单的前提;s3:GetBucketLocation:允许获取桶的所在区域;s3:GetObject:允许读取对象实际内容,用于从样本文件推断 Schema。
随后将该策略关联到实际执行摄取任务的 IAM 用户或角色,并在 Ingestion Recipe 的aws_config中指定该身份。
3.3 Azure Blob Storage
source: type: excel config: path_list: - "https://storageaccountname.blob.core.windows.net/abs-data/excel/*/*.xlsx" azure_config: account_name: storageaccountname sas_token: sv=2022-11-02&ss=b&srt=sco&sp=rwdlacx&se=2025-06-07T21:00:00Z&st=2025-05-07T13:00:00Z&spr=https&sig=a1B2c3D4%3D container_name: abs-data profiling: enabled: falseAzure 前置条件(见 excel_pre.md):
- Azure 存储账户(Storage Account):为数据提供唯一的命名空间;
- 认证凭据,支持以下任一方式:
- 账户密钥(Account Key):使用存储账户的访问密钥;
- 客户端密钥(Client Secret):使用服务主体(Service Principal)的 Client ID 与 Client Secret 进行 Microsoft Entra ID 认证;
- SAS Token:提供受限、限时的共享访问签名;
- 容器(Container):组织 Blob 的容器(类似文件系统目录);
- 访问权限:账户密钥方式需存储账户全量访问权限;客户端密钥方式需合适的 Azure 角色分配(如
Storage Blob Data Contributor);SAS Token 方式的权限由 Token 自身定义。
在实现上,abs_browser()通过azure_config.get_blob_service_client()获取容器客户端,按前缀分页列出 Blob;get_abs_file()下载 Blob 内容;设置use_abs_blob_tags后会将 Blob 标签同步为 DataHub 标签。
四、完整配置项详解
连接器的全部配置项定义在 config.py 的ExcelSourceConfig中(基于 Pydantic),下表汇总了各字段及其行为:
| 配置项 | 类型 | 默认值 | 说明 |
|---|---|---|---|
path_list | List[str] | 必填 | 要摄取的 Excel 文件或文件夹路径列表,支持*通配符 |
path_pattern | AllowDenyPattern | 全部允许 | 用于过滤文件路径的正则(allow/deny)模式 |
aws_config | AwsConnectionConfig | None | AWS 连接配置(S3 来源必需) |
use_s3_bucket_tags | bool | False | 是否从 S3 桶创建 DataHub 标签 |
use_s3_object_tags | bool | False | 是否从 S3 对象创建 DataHub 标签 |
verify_ssl | bool或str | True | 布尔值控制是否校验 TLS 证书;字符串则为 CA bundle 路径 |
azure_config | AzureConnectionConfig | None | Azure 连接配置(ABS 来源必需) |
use_abs_blob_tags | bool | False | 是否从 ABS Blob 标签创建 DataHub 标签 |
convert_urns_to_lowercase | bool | False | 是否将 Excel 资产的 URN 转换为小写 |
active_sheet_only | bool | False | 仅摄取工作簿的活动工作表;未设置时摄取全部工作表 |
worksheet_pattern | AllowDenyPattern | 全部允许 | 按工作表过滤的正则模式,工作表以文件名(不含扩展名).工作表名标识,例如匹配report.xlsx的Sheet1写作report.Sheet1 |
profile_pattern | AllowDenyPattern | 全部允许 | 按工作表过滤 Profiling 范围,标识方式同worksheet_pattern |
profiling | GEProfilingConfig | 默认配置 | 数据画像配置(详见下文) |
stateful_ingestion | StatefulStaleMetadataRemovalConfig | None | 有状态摄取与陈旧元数据删除配置 |
从源码结构看,worksheet_pattern的匹配对象正是gen_dataset_name()生成的名称(形如path/[filename]sheet_name),source.py的process_file中通过self.config.worksheet_pattern.allowed(dataset_name)决定是否跳过某个工作表;profile_pattern则在 profiling.py 的generate_profile()中做同样判断。
五、从工作表到 Dataset 的落地流程
source.py 的get_workunits_internal()是摄取主流程,整体链路为:
path_list -> uri_type() 判定来源 -> retrieve_file_data() 遍历文件 -> check_file_is_valid()(path_pattern + 扩展名校验) -> get_*_file() 读取二进制流 -> process_file() 加载工作簿 -> get_tables() 识别各工作表 -> process_dataset() 生成 MCP WorkUnitprocess_dataset()中依次完成:
- Dataset 命名与 URN:调用
gen_dataset_name()(见 util.py)拼出{directory}/[{filename}]{sheet_name}形式的名称(支持小写转换),再经make_dataset_urn_with_platform_instance生成urn:li:dataset:(urn:li:dataPlatform:excel,{path}/[{filename}]{sheet_name},PROD)格式的 URN,其中PROD为默认环境,受env与platform_instance配置影响; - Schema 元数据:
construct_schema_metadata()读取 DataFrame 的dtypes,通过field_type_mapping将 numpy 数据类型映射为 DataHub 的 SchemaField 类型——整型/浮点/复数/布尔映射为NumberTypeClass/BooleanTypeClass,object/string映射为StringTypeClass,datetime64等时间类型映射为DateTypeClass,category/interval/sparse映射为RecordTypeClass,未知类型回退为NullTypeClass; - DatasetProperties:将工作簿属性与表头/表尾元数据合并结果写入
customProperties,并将created/modified转换为毫秒时间戳(TimeStamp); - Container 层级:调用
container_WU_creator.create_container_hierarchy(relative_path, dataset_urn)依据目录结构创建 Container 实体; - 标签同步:S3/ABS 场景按配置同步桶/对象/Blob 标签;
- Profiling:若
is_profiling_enabled()为真,则构建ExcelProfiler产出 DatasetProfile。
六、数据画像(Profiling)
Profiling 由 profiling.py 的ExcelProfiler实现,输出DatasetProfileClass(行数、列数)与逐字段的DatasetFieldProfileClass。其统计能力包括:
- 表级:
rowCount、columnCount; - 字段级:
uniqueCount(去重计数)、nullCount、min、max、mean、median、stdev、25%/75% 分位数(QuantileClass)以及前 N 个样本值(sampleValues,数量由field_sample_values_limit控制)。
这些指标由compute_field_statistics()计算,且遵循GEProfilingConfig的开关控制(如include_field_null_count、include_field_distinct_count、include_field_min_value、include_field_max_value、include_field_mean_value、include_field_stddev_value、include_field_median_value、include_field_quantiles、include_field_sample_values)。数值型统计仅对数值列(int/uint/float/complex 系列 dtype)计算;profile_table_level_only可只做表级画像;max_number_of_fields_to_profile可限制参与画像的字段数,超出的字段会被丢弃并在报告中记录。开启方式即 Recipe 中的:
source: type: excel config: profiling: enabled: true profile_table_level_only: false此外use_sampling/sample_size控制采样策略(sample_size=0表示全量)。
七、有状态摄取与删除检测
通过配置stateful_ingestion.remove_stale_metadata可启用删除检测:当某个工作表/文件在源端消失时,DataHub 中对应的陈旧元数据会被自动清理。配置示例如下:
source: type: excel config: path_list: - "/data/path/reporting/excel/*.xlsx" stateful_ingestion: enabled: true remove_stale_metadata: true这一能力基于StatefulIngestionSourceBase基类(stateful_ingestion_base.py 体系),在 source.py 的类装饰器中声明为supported=True。
八、故障排查与验证
根据 excel_post.md 的指引,当摄取失败时,应优先检查:凭据、权限、连通性、范围过滤(path_list/path_pattern/worksheet_pattern),随后查看摄取日志中的来源相关错误并调整配置。连接器的ExcelSourceReport(report.py)会记录文件扫描/处理、工作表扫描/处理/丢弃、Profiling 跳过等计数,便于定位问题。
仓库中的测试用例可作为行为验证的参考:
- 单元测试:test_excel_file.py(表头探测、页脚识别、元数据提取等)、test_excel_samples.py;
- 集成测试:test_excel_local.py、test_excel_s3.py、test_excel_abs.py,并配有对应的 golden 输出文件(
excel_file_test_golden.json、excel_s3_test_golden.json、excel_abs_test_golden.json),可直接对照查看各类工作表的预期摄取结果。
九、使用注意与限制
- 连接器当前为Alpha支持状态(
SupportStatus.ALPHA),行为受平台暴露的 API、权限与元数据约束,建议在生产环境先在小范围路径上验证; - 每个工作表独立成为 Dataset,工作簿文件本身不会生成实体,不要期待以文件为粒度管理 Excel 元数据;
- 表头/页脚识别依赖启发式打分,若工作表的表头行为异常(如全数值、多级表头、无表头),可能被判定为"不含表格"而跳过,此时可检查报告中的
worksheet dropped计数; - S3/ABS 摄取必须提供对应的
aws_config/azure_config,否则s3_browser/abs_browser会直接抛出ValueError; - 更多通用摄取(Recipe 编写与运行)方法可参考仓库 metadata-ingestion 根文档 与 recipe_overview.md。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考