
DataHub 列级分类框架实战指南配置、自研 Classifier 与源码级原理【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahubDataHub 的分类Classification框架允许摄入源在 ingestion 过程中自动识别列的信息类型info type并将其作为业务术语glossary term挂载到对应字段上从而在数据资产进入目录的第一时间完成敏感信息与业务语义的自动打标。这是一项**显式开启opt-in**的功能默认关闭内置的datahub分类器已被移除但完整的框架——Classifier接口、分类器注册表与各数据源侧的编排逻辑——被完整保留供开发者注册自己的分类器。读完本文你将掌握分类配置的每个参数及其默认值、自研并注册一个Classifier的完整流程以及从采样到术语落库的底层实现原理。分类框架是什么分类框架解决的是这样一个问题当数据源如 Snowflake、BigQuery、Redshift 等的元数据被摄入 DataHub 时框架会先读取每张表若干行的样本数据sample values把每一列的列名、描述、数据类型与样本值交给分类器Classifier由分类器预测该列属于哪种 info type例如email、credit_card、phone_number等最终将预测结果转换为对应的 glossary term 写入该列的SchemaMetadata。整个链路无需人工干预敏感数据一进目录就被自动标注。该功能在 YAML recipe 中以嵌套字段classification声明文档明确规定用.表示嵌套层级配置结构由 ClassificationConfig 定义默认值为enabled: false、sample_size: 100、max_workers: 进程 CPU 核数、table_pattern/column_pattern均放行全部、classifiers默认指向已移除的内置datahub类型。重大变更内置 datahub 分类器已被移除仓库在 classifier_registry.py 的注释中明确记录了本次变更的背景内置的datahub分类器DataHubClassifier依赖未维护的acryl-datahub-classify包该包将numpy钉死在2并携带一套过时的 spaCy 依赖栈阻塞了整个摄入框架的依赖升级因此被移除。随之下架的还有此前从该包导入的数据契约类型现已在 classification_types.py 中原样内置vendored保证第三方分类器不需要该依赖即可使用扩展点。移除之后没有任何分类器在注册表中默认注册。如果 recipe 设置了classification.enabled: true却不注册替代分类器框架会在启动阶段**快速失败fail fast**并给出指引。这条引导信息定义在 BUILTIN_CLASSIFIER_REMOVED_MESSAGE内置的datahub列级分类器已从 acryl-datahub 移除因为它依赖未维护的acryl-datahub-classify包钉死 numpy2。如需继续使用请安装最后一个支持它的版本pip install acryl-datahub1.6.0.5否则请注册你自己的分类器实现或关闭分类功能classification.enabled: false。该快速失败逻辑位于 classification_mixin.py当从注册表查询分类器类型抛出KeyError时若类型正是DEFAULT_CLASSIFIER_TYPE即datahub则抛出携带上述引导信息的ConfigurationError若是其他未知类型则抛出“在注册表中找不到该类型”的通用错误。注意get_classifiers()只在分类开启时才解析分类器——这样即使 recipe 保持默认的classifiers: [{type: datahub}]而未开启分类任何数据源也不会误触发错误。配置详解以下表格完整列出分类功能的全部配置字段字段与默认值均以当前仓库 classifier.py 与官方文档为准字段必填类型说明默认值enabled否boolean是否使用分类来自动识别 glossary termsFalsesample_size否int用于分类的样本值条数100max_workers否int用于分类的 worker 进程数设为 1 可禁用并行进程 CPU 核数os.cpu_count() or 4info_type_to_term否Dict[str, str]可选的 info type 到 glossary term 标识符的映射默认直接用 info type 作为 glossary term 标识符classifiers否Array of object用于自动识别 glossary terms 的分类器列表配置多个分类器时列表中靠后的分类器给出的 info type 预测优先级更高[{type: datahub, config: None}]table_pattern否AllowDenyPattern过滤参与分类的表的正则与父配置中其他 pattern 组合使用正则需匹配database.schema.table格式的完整表名例如匹配 Customer 库 public schema 下所有 customer 开头的表可写Customer.public.customer.*{allow: [.*], deny: [], ignoreCase: True}table_pattern.allow否Array of string纳入摄入的正则列表[.*]table_pattern.deny否Array of string排除出摄入的正则列表[]table_pattern.ignoreCase否boolean模式匹配时是否忽略大小写Truecolumn_pattern否AllowDenyPattern过滤参与分类的列的正则与父配置中其他 pattern 组合使用正则需匹配database.schema.table.column格式{allow: [.*], deny: [], ignoreCase: True}column_pattern.allow否Array of string纳入摄入的正则列表[.*]column_pattern.deny否Array of string排除出摄入的正则列表[]column_pattern.ignoreCase否boolean模式匹配时是否忽略大小写True完整 recipe 示例以 Snowflake 源为例一个开启分类并只针对特定表、特定列进行扫描的 recipe 如下source: type: snowflake config: # ... 数据源连接等基础配置 ... classification: enabled: true sample_size: 200 # 每列取 200 条样本值 max_workers: 4 # 4 个 worker 进程并行分类 info_type_to_term: email: urn:li:glossaryTerm:customer.email # 将 email info type 映射到自定义术语 table_pattern: allow: - Customer.public.customer.* deny: - Customer.public.customer_staging.* ignoreCase: true column_pattern: allow: - .*\\.(email|phone|ssn)$ # 只分类这些后缀的列 deny: [] ignoreCase: true classifiers: - type: my-classifier配置字段的源码级解读enabled 与 classifiers 的联动ClassificationHandler.is_classification_enabled() 要求分类配置存在、enabled为真且classifiers非空三者同时满足才算开启。sample_size 的 1.2 倍放大框架不会只取恰好sample_size行。在 classification_workunit_processor 与 SQL 源的_classify方法中实际请求的行数是sample_size * SAMPLE_SIZE_MULTIPLIER其中SAMPLE_SIZE_MULTIPLIER 1.2见 classification_mixin.py。原因写在 data_reader.py跨列统一取样的查询难免混入 NULL多取 20% 行可以补偿 NULL 对分类质量的影响。max_workers 与并行模型当max_workers 1时框架使用concurrent.futures.ProcessPoolExecutor以**每批 5 列BATCH_SIZE**为单位把列分发给多个 worker 进程并行调用分类器的classify见 async_classify。值得注意的实现细节进程池显式使用multiprocessing.get_context(spawn)启动上下文——因为 Linux 上 Python 默认的fork启动方式在主进程使用线程时不安全Python 3.14 起 Linux 默认也将切换为spawn。如果自定义分类器包含无法序列化的状态可把max_workers设为 1 走串行路径。多分类器的优先级文档与配置注释都强调“列表中靠后的分类器预测优先”。落到实现上get_terms_for_column 对同一列的所有infotype_proposals取confidence_level最大者而多个分类器按顺序逐个执行、update_field_terms以列为键覆盖写入因此后执行的分类器列表中靠后的提案会覆盖先执行者的结果。info_type_to_term 的用途分类器产出的是 info type如email而挂载到字段上的是 glossary term 的标识符。框架默认直接用 info type 字符串作为术语标识符如需映射到既有术语可通过info_type_to_term提供映射未命中的 info type 仍回退为自身。pattern 的组合语义table_pattern 与 column_pattern 是叠加过滤关系——列级判断is_classification_enabled_for_column同时要求表与列都匹配各自 pattern见 classification_mixin.py它们还与父配置如 Snowflake 的schema_pattern/table_pattern组合使用即“父级先过滤分类 pattern 再过滤”。分类的底层工作流完整的分类链路分为四个阶段可通过 ClassificationHandler 串联理解构建分类器列表get_classifiers()根据 recipe 中声明的类型逐一从classifier_registry解析并调用create(config_dict...)实例化未知类型抛出ConfigurationError对内置datahub类型则抛出移除指引。取样通过DataReader.get_sample_data_for_table抽象方法见 data_reader.py获取约sample_size * 1.2行数据返回{列名: 值列表}字典。在 classify_schema_fields 中若取样是惰性回调且执行失败会累加num_tables_fetch_sample_values_failed统计并降级为空字典继续同时提示确认数据集上的 SELECT 权限。逐列构建 ColumnInfo 并调用分类器get_columns_to_classify依据column_pattern过滤列为每个字段构造ColumnInfo其metadata携带Name、Description、DataType、Dataset_Name四项元数据见 classification_types.pyvalues为样本值。随后串行或并行地交给各分类器分类器回填infotype_proposals。写入 glossary termspopulate_terms_in_schema_metadata为每个命中列生成GlossaryTerms方面通过make_term_urn(term)构造术语 URN并保留字段上已有的术语新术语排在前面原术语追加在后审计戳actor固定为datahub见 classification_mixin.py。在 SQL 源侧编排入口位于 sql_common.py 的SQLAlchemySource._classify方法——这是所有 SQLAlchemy 系 SQL 源的公共基类它在生成 schema 元数据后调用ClassificationHandler.classify_schema_fields。非 SQLAlchemy 的源如 Snowflake 的 snowflake_schema_gen.py、BigQuery 的 bigquery_schema_gen.py、Redshift 的 redshift.py则通过classification_workunit_processor以 workunit 装饰器的方式接入该处理器只拦截携带SchemaMetadata方面MCP/MCPW的 workunit其余 workunit 原样透传。自研分类器Bring Your Own Classifier这是当前框架唯一受支持的扩展路径。你需要做三件事实现Classifier接口、注册到注册表、在 recipe 中引用。1. 实现 Classifier 接口Classifier是一个抽象基类见 classifier.py只要求实现classify(columns) - columns一个抽象方法并约定类方法create(config_dict)负责按配置实例化。输入输出都是ColumnInfo列表分类器的职责是遍历输入的列、根据ColumnInfo.metadataName/Description/DataType/Dataset_Name与ColumnInfo.values样本值做预测并把预测结果填回column.infotype_proposals然后原样返回列表。文档给出的完整模板可直接复制使用from typing import Any, Dict, List from datahub.ingestion.glossary.classification_types import ColumnInfo from datahub.ingestion.glossary.classifier import Classifier from datahub.ingestion.glossary.classifier_registry import classifier_registry class MyClassifier(Classifier): def __init__(self, config: Dict[str, Any]) - None: self.config config classmethod def create(cls, config_dict: Dict[str, Any]) - MyClassifier: return cls(config_dict or {}) def classify(self, columns: List[ColumnInfo]) - List[ColumnInfo]: # Populate column.infotype_proposals here. return columns classifier_registry.register(my-classifier, MyClassifier)2. 理解数据契约ColumnInfo及其配套类型定义在 classification_types.pyColumnInfometadataMetadata对象、values样本值列表、infotype_proposals可空预测结果。Metadata从meta_info字典派生出的name、description、datatype、dataset_name四个属性取值键分别为Name、Description、Datatype、Dataset_Name任一字段可能缺失如列无描述因此都是 Optional。InfotypeProposalinfotype信息类型字符串、confidence_level置信度 float、debug_info各特征对预测的贡献占比DebugInfo的name/description/datatype/values四个可选浮点字段。框架只消费InfotypeProposal.infotype与confidence_level同列多个提案取置信度最高者见 get_terms_for_column并把命中列的完整名称数据集名.列名记入报告的info_types_detected。3. 注册并启用classifier_registry是PluginRegistry[Classifier]实例见 classifier_registry.py在模块导入时执行register(my-classifier, MyClassifier)即可完成注册。然后在 recipe 中引用source: type: snowflake config: # ... source config ... classification: enabled: true classifiers: - type: my-classifier若分类器需要自定义参数可通过classifiers[].config传入任意 JSON/YAML 对象框架会原样交给create(config_dict...)config字段的完整校验由你的分类器自己负责见 classifier.py 中DynamicTypedClassifierConfig的注释。4. 常见排查要点启动即报The built-in datahub column-level classifier has been removed...说明 recipe 开了分类但没注册任何分类器且classifiers仍是默认值。按引导信息二选一注册自己的分类器或关闭classification.enabled。报Cannot find classifier class of typexxx in the registry!classifiers里写的类型未注册检查拼写与注册语句是否被执行。某张表没被分类检查table_pattern/column_pattern与父级过滤是否放行检查分类是否开启且classifiers非空见is_classification_enabled_for_table的三重条件classification_mixin.py。取样失败报告中的num_tables_fetch_sample_values_failed会增加日志提示确认对数据集授予了 SELECT 权限classification_mixin.py。并行模式异常ProcessPoolExecutor使用spawn上下文分类器对象需可序列化否则将max_workers设为 1 降级为进程内串行。支持的源文档明确声明所有 SQL 源均支持分类。从仓库代码看支持方式分两类基于 SQLAlchemy 的通用 SQL 源继承 SQLAlchemySource含 Oracle、Teradata、HANA、Doris 等具体实现在公共基类中完成编排各自实现采样与编排的专用源Snowflakesnowflake_schema_gen.py、BigQuerybigquery_schema_gen.py、Redshiftredshift.py等通过ClassificationSourceConfigMixin混入分类配置见各自 config 文件如 snowflake_config.py、bigquery_config.py再经由classification_workunit_processor挂接处理。此外非 SQL 的 DynamoDB 源dynamodb.py同样接入了该框架说明凡是能提供样本数据的源都可以复用这套机制。可以推断只要你的自定义源混入ClassificationSourceConfigMixin并接入ClassificationHandler或classification_workunit_processor就能获得与内置 SQL 源一致的分级能力。小结DataHub 分类框架的价值在于把“数据进入目录”与“语义自动标注”绑定为一次性的流水线动作。当前版本的关键事实是内置分类器已随acryl-datahub-classify依赖的移除而下线如需保留请锁定acryl-datahub1.6.0.5但接口、注册表与编排层完整保留且文档与代码高度自洽——你只需实现一个classify()方法并在注册表中登记即可在全部 SQL 源上获得配置灵活采样量、并行度、表列过滤、术语映射、多分类器优先级的自动术语识别能力。延伸阅读可参考 classification 官方文档、分类配置模型、分类编排实现 与 数据契约类型。【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考