ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

DataHub Kafka 元数据连接器完全指南:从 Topic、Schema Registry 到 Dataset 的实战接入

DataHub Kafka 元数据连接器完全指南:从 Topic、Schema Registry 到 Dataset 的实战接入 DataHub Kafka 元数据连接器完全指南从 Topic、Schema Registry 到 Dataset 的实战接入【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub本篇技术指南围绕 DataHub 元数据摄入框架metadata-ingestion内置的kafka连接器展开详细讲解如何将 Apache Kafka 集群中的 Topic、Schema RegistryConfluent Schema Registry中的 Avro / Protobuf / JSON 模式以及 Confluent Cloud Stream Catalog 中的标签与业务元数据批量同步为 DataHub 中的 Dataset、Schema Metadata、Dataset Fields 等标准资产。读完本文你将掌握最小可用 Recipe 的编写、Confluent Cloud 专用配置、自定义 Schema Registry 接入、OAuth 认证回调、Avro 元数据自动映射、多阶段 Schema 解析与消息级数据画像Data Profiling等一整套生产可用的接入方案。连接器能做什么从 Kafka 到 DataHub 的元数据管道kafka模块位于 metadata-ingestion/src/datahub/ingestion/source/kafka/用于将 Apache Kafka 的元数据摄入 DataHub专为生产级摄入工作流设计。它是一个开源实现支持状态 GASupportStatus.GA主要提取三类内容Topic 元数据通过 Kafka Admin / Consumer 客户端枚举集群中的 Topic并额外抓取每个 Topic 的分区数、副本因子、retention.ms、cleanup.policy、max.message.bytes等配置作为 Dataset 的自定义属性custom propertiesSchema 元数据从 Schema Registry 获取每个 Topic 关联的 key/value 模式支持 Avro、Protobuf 与 JSON 三种模式类型源码中以SCHEMA_TYPE_AVRO/SCHEMA_TYPE_PROTOBUF/SCHEMA_TYPE_JSON定义见 kafka_constants.py并转换为 DataHub 的 SchemaFieldConfluent Cloud Stream Catalog 元数据可选当目标集群是 Confluent Cloud 时可同步 Stream Governance 中为 Topic 整理的标签tags与业务元数据business metadata。在源码的能力声明中kafka.py该连接器通过capability装饰器明确标注了支持范围能力支持情况说明SCHEMA_METADATA✅从 Schema Registry 提取每个 Topic 关联的 SchemaAvro / Protobuf 为认证certified支持JSON 为孵化incubating支持且支持 Schema 引用referencesDATA_PROFILING✅可选通过profiling.enabled开启消息内容画像TAGS✅有条件需要 Confluent Cloud Stream Governance通过confluent_catalog开启PLATFORM_INSTANCE✅多 Kafka 集群场景使用platform_instance配置DESCRIPTIONS✅将 Avro Schema 顶层的doc字段映射为 Dataset 描述LINEAGE_COARSE/LINEAGE_FINE❌Topic 间血缘不支持如使用 Kafka Connect请改用 kafka-connect 连接器TEST_CONNECTION✅连接测试默认启用概念映射Kafka 世界与 DataHub 数据模型的对应关系要理解连接器如何组织元数据首先需要掌握两者的概念映射。以下是文档给出的标准映射关系原始表格见 README.mdKafka 概念DataHub 概念说明TopicDataset子类型为Topic源码中对应DatasetSubTypes.TOPICSchemaSubjectSchema Metadata支持 Avro、Protobuf、JSON 模式Message FieldsDataset Fields从模式中提取或在开启 Schema 解析时由消息内容推断Kafka ClusterData Platform Instance当配置了platform_instance时生效Schema 元数据Tags、Glossary Terms、OwnersCorpUser / CorpGroup可选仅 Avro当开启enable_meta_mapping并配置meta_mapping/field_meta_mapping指令时从 Schema 属性推导从源码看Dataset 的生成集中在_emit_dataset方法kafka.pyTopic 被映射为Dataset实体subtypeDatasetSubTypes.TOPIC若开启ingest_schemas_as_entitiesSchema Registry 的 Subject 还会以DatasetSubTypes.SCHEMA子类型单独摄入。同时Topic 的Partitions、Replication Factor以及min.insync.replicas、retention.bytes等配置枚举定义见 kafka.py 的KafkaTopicConfigKeys会被写入 Dataset 的自定义属性为运维侧排查 Topic 配置问题提供了直接依据。快速上手最小可用 Recipe官方在 kafka_recipe.yml 中给出了一个最小可用的 Recipe 骨架source: type: kafka config: platform_instance: YOUR_CLUSTER_ID connection: bootstrap: broker:9092 schema_registry_url: http://localhost:8081 # Optional: Enable data profiling profiling: enabled: true sample_size: 1000 max_sample_time_seconds: 60 sampling_strategy: latest sink: # sink configs其中connection.bootstrap是 Kafka broker 地址connection.schema_registry_url是 Schema Registry 端点。platform_instance用于区分同一平台下的多个 Kafka 集群详见 docs/platform-instances.md。执行摄入的命令为datahub ingest -c kafka_recipe.yml摄入前建议先做连接测试该能力默认开启capability(SourceCapability.TEST_CONNECTION, Enabled by default)。从 kafka.py 的KafkaConnectionTest实现可见它通过consumer.list_topics(timeout10)验证 broker 连通性通过SchemaRegistryClient(...).get_subjects()验证 Schema Registry 连通性并将两项结果以basic_connectivity与schema_metadata能力报告的形式返回datahub ingest -c kafka_recipe.yml --dry-run连接 Confluent CloudAPI Key、ACL 与完整 Recipe前置条件为 API Key 配置最小 ACL使用 Confluent Cloud 时consumer_config.sasl.username和consumer_config.sasl.password使用集群页面的Data Integration - API Keys中创建的 API 凭证schema_registry_config.basic.auth.user.info使用 Schema Registry 的 API 凭证位于Schema Registry - API credentials。创建集群 API Key 时必须为 Key 关联如下 ACLDataHub 才能读取 Confluent Cloud 上 Topic 的元数据Topic Name * Permission ALLOW Operation DESCRIBE Pattern Type LITERAL完整 Recipesource: type: kafka config: platform_instance: YOUR_CLUSTER_ID connection: bootstrap: abc-defg.eu-west-1.aws.confluent.cloud:9092 consumer_config: security.protocol: SASL_SSL sasl.mechanism: PLAIN sasl.username: ${CLUSTER_API_KEY_ID} sasl.password: ${CLUSTER_API_KEY_SECRET} schema_registry_url: https://abc-defgh.us-east-2.aws.confluent.cloud schema_registry_config: basic.auth.user.info: ${REGISTRY_API_KEY_ID}:${REGISTRY_API_KEY_SECRET} sink: # sink configs${...}形式的占位符支持从环境变量或 DataHub 的 Secret 机制中取值避免凭证明文落盘仓库同时提供datahub.configuration.kafka下的连接配置基类与 Secret 脱敏机制见 kafka_config.py。为 Topic 分配 Domain如需在摄入时为 Topic 自动归属 Domain可配置domain映射。domain中的键既可以是完整 URN也可以是裸 Domain ID如13ae4d85-d955-49fc-8474-9004c663a810其值使用allow/deny正则模式匹配 Topic 名source: type: kafka config: # ...connection block domain: urn:li:domain:13ae4d85-d955-49fc-8474-9004c663a810: allow: - .* urn:li:domain:d6ec9868-6736-4b1f-8aa6-fee4c5948f17: deny: - .*注意目标 Domain 需要先存在于 DataHub 实例中可参考 docs/domains.md 创建 Domain摄入时连接器会通过DomainRegistry解析 URN 并写入 Dataset 的 domain 属性见 kafka.py。非默认 Subject 命名策略topic_subject_map如果 Schema Registry 使用了非默认的 Subject 命名策略例如RecordNameStrategy默认的topic-key/topic-value查找会失败。此时必须通过topic_subject_map显式声明 Topic 的 key/value 模式与 Subject 名的对应关系source: type: kafka config: # ...connection block # Defines the mapping for the key value schemas associated with a topic the subject name registered with the # kafka schema registry. topic_subject_map: # Defines both key value schema for topic my_topic_1 my_topic_1-key: io.acryl.Schema1 my_topic_1-value: io.acryl.Schema2 # Defines only the value schema for topic my_topic_2 (the topic doesnt have a key schema). my_topic_2-value: io.acryl.Schema3在 kafka_config.py 中topic_subject_map的语义被进一步明确一旦提供它将覆盖默认的 Subject 解析即使使用的是TopicNameStrategy或TopicRecordNameStrategy。Confluent Cloud Stream Catalog同步标签与业务元数据在 Confluent Cloud 上Stream Governance 中维护的标签与业务元数据存放在Stream Catalog中而不是直接挂在 Topic 上。开启confluent_catalog配置块即可将它们同步到对应的 DataHub Topic Dataset 上。Catalog 由 Schema Registry 端点提供且复用同一个 API Key因此一个已能访问 Schema Registry 的 Recipe 只需增加一行enabled: truesource: type: kafka config: connection: bootstrap: abc-defg.eu-west-1.aws.confluent.cloud:9092 schema_registry_url: https://abc-defgh.us-east-2.aws.confluent.cloud schema_registry_config: basic.auth.user.info: ${REGISTRY_API_KEY_ID}:${REGISTRY_API_KEY_SECRET} confluent_catalog: enabled: true同步规则Confluent 标签tags映射为 DataHub 标签业务元数据属性business metadata attributes映射为 Topic 的自定义属性。可以通过include_tags: false或include_business_metadata: false只取其中一类。对应的完整配置类为KafkaConfluentCatalogConfigkafka_config.py其默认值include_tagsTrue、include_business_metadataTrue。前提条件与限制仅支持 Confluent Cloud且环境需要购买Stream Governance Advanced套餐。标签与业务元数据属于 Advanced 功能在 Essentials 套餐上Catalog API 即使读取也会返回403自建 Kafka 不存在 Catalog配置块会被忽略。Schema Registry API Key 的角色需要具备 Catalog 读取权限实践中为环境的DataSteward角色。仅EnvironmentAdmin不够——它可以管理环境但不具备 Catalog 读取权限。如果 Key 无法读取 Catalog摄入不会失败而是跳过 Catalog 元数据并记录一条警告。Catalog 覆盖整个环境如果该环境包含多个 Kafka 集群同一 Topic 名可能重复出现。这些重名 Topic 会被跳过伴随警告除非设置confluent_catalog.cluster_id为当前摄入的集群 ID如lkc-xxxxx。Catalog不提供 Topic 之间的血缘。Confluent Stream Lineage UI 中展示的生产者/消费者关系图无法通过 API 获取。需要 Connector 到 Topic 血缘的场景可改用 kafka-connect 连接器。从源码实现看Catalog 元数据的应用逻辑位于_apply_catalog_metadatakafka.py当include_tags开启时Catalog 标签会追加到 Topic 的标签列表当include_business_metadata开启时业务元数据通过non_colliding_business_metadata与 broker 侧 Topic 属性做冲突规避后合并进自定义属性避免覆盖Partitions、Replication Factor等关键属性。值得注意的是全局标签是替换型replacementaspect如果 Catalog 只被部分读取is_complete()为假被遗漏的 Topic 再次摄入时可能丢失原有 Catalog 标签此时源码会记录一次警告提示待 Catalog 可完整读取后重跑恢复。自定义 Schema Registry实现 KafkaSchemaRegistryBaseKafka 连接器默认使用 Confluent 的 Kafka Schema Registry 来解析 Topic 的 key/value 模式原生支持AVRO与PROTOBUF两种模式类型。如果使用的是自定义 Schema Registry或者模式类型不是 Avro / Protobuf可以自行实现KafkaSchemaRegistryBase抽象类并提供get_schema_metadata(topic, platform_urn)方法——该方法接收 Topic 名返回包含该 Topic 模式的SchemaMetadata对象class KafkaSchemaRegistryBase(ABC): abstractmethod def get_schema_metadata( self, topic: str, platform_urn: str ) - Optional[SchemaMetadata]: pass接口定义见 kafka_schema_registry_base.py其中还包含get_subjects()、_get_subject_for_topic()、get_schema_registry_client()等抽象方法以及默认的批量模式拉取实现get_schema_and_fields_batch()与build_schema_metadata_with_key()。官方默认实现参考datahub.ingestion.source.confluent_schema_registry::ConfluentSchemaRegistry。自定义类通过schema_registry_class配置项指定连接器会使用import_path动态加载kafka.pysource: type: kafka config: # Set the custom schema registry implementation class schema_registry_class: datahub.ingestion.source.confluent_schema_registry.ConfluentSchemaRegistry # Coordinates connection: bootstrap: broker:9092 schema_registry_url: http://localhost:8081OAuth 认证回调为 Source 与 Sink 配置 OAuth BearerKafka 连接器为 Source消费者与 Sink生产者均提供了 OAuth 回调支持Sourceconfig.connection.consumer_config.oauth_cbSinkconfig.connection.producer_config.oauth_cb回调以python-module:function-name格式引用 Python 函数。例如oauth:create_token表示create_token定义在oauth.py中且oauth.py必须在PYTHONPATH中可被导入。内置回调推荐DataHub 内置了常见场景的 OAuth 回调AWS MSK IAMdatahub_actions.utils.kafka_msk_iam:oauth_cbAzure Event Hubsdatahub_actions.utils.kafka_eventhubs_auth:oauth_cb使用内置回调需安装acryl-datahub-actions包pip install acryl-datahub-actions1.3.1.2自定义回调自定义回调模块需确保 DataHub 进程可访问例如通过PYTHONPATH/path/to/your/module:$PYTHONPATH或pip install my-oauth-package。Kafka Source 示例source: type: kafka config: # Set the custom schema registry implementation class schema_registry_class: datahub.ingestion.source.confluent_schema_registry.ConfluentSchemaRegistry # Coordinates connection: bootstrap: broker:9092 schema_registry_url: http://localhost:8081 consumer_config: security.protocol: SASL_PLAINTEXT sasl.mechanism: OAUTHBEARER oauth_cb: oauth:create_token # sink configsKafka Sink 示例MSK IAM 认证sink: type: datahub-kafka config: connection: bootstrap: b-1.msk.us-west-2.amazonaws.com:9098 schema_registry_url: http://datahub-gms:8080/schema-registry/api/ producer_config: security.protocol: SASL_SSL sasl.mechanism: OAUTHBEARER sasl.oauthbearer.method: default oauth_cb: datahub_actions.utils.kafka_msk_iam:oauth_cb从源码看配置了 OAuth 回调后连接器会在创建 Consumer / AdminClient 时显式调用一次consumer.poll(timeout30)触发回调执行kafka.py确保后续元数据请求携带有效令牌。Avro 元数据自动映射将 Schema 属性转为 Owner、Tag 与 TermAvro 规范允许 Schema 携带规范未定义的附加属性arbitrary metadata业界常用它承载业务元数据。Kafka 连接器可以将这些属性直接转换为 DataHub 的 Owner、Tag 与 Glossary Term。注意元数据映射目前仅支持 Avro 模式且要求这些 Avro 模式已推送到 Schema Registry。同时需要enable_meta_mapping默认true开启映射处理。简单标签schema_tags_field如果 Avro Schema 中嵌入了标签列表顶层或字段级可用schema_tags_field指定存放标签的字段名默认tags。示例 Avro Schema{ name: sampleRecord, type: record, tags: [tag1, tag2], fields: [ { name: field_1, type: string, tags: [tag3, tag4] } ] }config: schema_tags_field: tags对应源码在 kafka.py连接器读取 Avro 顶层other_props中schema_tags_field指定的列表为每个标签加上tag_prefix默认空字符串后生成 DataHub 标签。meta_mapping 与 field_meta_mapping也可以将特定 Avro 字段精确映射为 Owner、Term 与 Tag示例 Avro Schema{ name: sampleRecord, type: record, owning_team: Data-Science, data_tier: Bronze, fields: [ { name: field_1, type: string, gdpr: { pii: true } } ] }对应的映射配置config: meta_mapping: owning_team: match: ^(.*) operation: add_owner config: owner_type: group data_tier: match: Bronze|Silver|Gold operation: add_term config: term: {{ $match }} field_meta_mapping: gdpr.pii: match: true operation: add_tag config: tag: pii其中meta_mapping针对顶层 Schema 属性field_meta_mapping针对字段级属性支持嵌套路径如gdpr.pii{{ $match }}可引用正则捕获的匹配值。底层实现通过OperationProcessorkafka.py执行add_owner/add_term/add_tag三类操作Owner 来源类型标记为SERVICE相关指令语义与 dbt 元数据自动映射 的实现一致那里提供了更丰富的示例可供参考。此外还有strip_user_ids_from_email从邮箱中剥离用户 ID与tag_prefix两个辅助配置项。多阶段 Schema 解析为缺失 Schema 的 Topic 自动兜底对于未在 Schema Registry 注册、或使用了非默认命名策略的 TopicDataHub 提供了多阶段 Schema 解析schema resolution。该能力独立于数据画像二者可分别开关但共享同一个profiling.max_workers并发配置。当schema_resolution.enabled为true时默认关闭连接器按下述顺序尝试解析TopicNameStrategy直接按topic-key/topic-value查找最常见TopicSubjectMap使用用户通过topic_subject_map配置的 Topic 到 Subject 映射RecordNameStrategy从消息数据中提取记录名查找record_name-key/record_name-valueTopicRecordNameStrategy组合 Topic 与记录名查找topic-record_name-key/record_name-valueSchema Inference最后兜底从消息数据分析推断 Schema。这一设计保证了与 Confluent Schema Registry 各种命名策略的最大兼容性。上述解析方法在 schema_resolution.py 中均有对应实现每种策略的结果还附带ResolutionMethod诊断标签如topic_name_strategy、record_name_strategy、schema_inference见 kafka_constants.py摄入日志中会记录每个 Topic 实际命中的解析路径。配置示例source: type: kafka config: schema_resolution: enabled: true # disabled by default sample_timeout_seconds: 2.0 offset_reset_strategy: hybrid # earliest, latest, or hybrid max_messages_per_topic: 10 profiling: max_workers: 20 # controls parallelization for both profiling and schema resolution nested_field_max_depth: 5Schema 推断的采样策略用于 Schema 推断的消息采样通过offset_reset_strategy控制读取起点hybrid默认优先尝试latest以追求速度若未发现近期消息则回退到earliestlatest只读近期消息。最快但在低频quietTopic 上可能失败earliest从 Topic 历史起点扫描。最全面但在大 Topic 上较慢。性能要点推断出的 Schema 会缓存60 分钟避免反复采样Worker 数量会根据 CPU 核数与 Topic 数量自动伸缩默认5 × CPU 核数见 schema_resolution.py将schema_resolution.enabled设为false则缺失 Schema 时仅产生警告而不自动解析。配置类定义见 kafka_config.pysample_timeout_seconds默认2.0秒单 Topic 采样的时间上限、max_messages_per_topic默认10条。数据画像从消息内容生成字段级统计与样本值Kafka 连接器支持对消息内容做数据画像Data Profiling产出字段级统计与样本值。画像与 Schema 解析相互独立可只开其一但两者的并发度都取自profiling.max_workers。完整配置source: type: kafka config: profiling: enabled: true sample_size: 200 # messages to sample per topic max_sample_time_seconds: 60 sampling_strategy: latest # latest, random, stratified, or full max_workers: 4 batch_size: 100 # Field-level statistics include_field_null_count: true include_field_distinct_count: true include_field_min_value: true include_field_max_value: true include_field_mean_value: true include_field_median_value: true include_field_stddev_value: true include_field_quantiles: false # expensive, disabled by default include_field_distinct_value_frequencies: false # expensive include_field_histogram: false # expensive include_field_sample_values: true # Nested field handling profile_nested_fields: true nested_field_max_depth: 10 # Scheduled profiling (optional) operation_config: lower_freq_profile_enabled: false profile_day_of_week: 1 # Monday0, Sunday6 profile_date_of_month: 15采样策略sampling_strategylatest默认从每个分区末尾采样最近的消息random在分区范围内随机 offset 采样stratified在 Topic 时间线上均匀分布采样full处理整个 Topic仍受sample_size上限约束。从源码看四种策略共享同一套消息解码管线kafka.pylatest/random/full仅在选择起始 offset 的函数_latest_offset/_random_offset/_full_offset上不同stratified则通过步长stride在分区内等距 seek 采样。样本按batch_size默认 100批量消费并受max_sample_time_seconds默认 60 秒时间窗约束。字段统计与解码细节画像结果基于每个字段生成DatasetFieldProfile包括空值计数/比例、唯一值计数/比例、最小值/最大值/均值/中位数/标准差以及按需分位数、去重值频率与直方图include_field_sample_values会输出实际样本值。字段类型识别依据 Avro 类型归类为NUMERIC/BOOLEAN/DATETIME/STRING/UNKNOWN分类映射见 kafka_constants.py。消息解码时连接器会识别 Confluent 线格式magic byte0x00 4 字节大端 Schema ID 的 5 字节头见 kafka_constants.py用缓存的 Schema 反序列化 Avro 消息再以nested_field_max_depth限制深度做嵌套字段展平。解码失败的消息会被跳过并计入profiling_samples_skipped/profiling_avro_decode_failures统计而不会向画像注入伪造字段。nested_field_max_depth默认 10用于防止深度嵌套或循环 JSON 结构引发递归错误对嵌套复杂的 Topic 建议调低。若 Topic 完全没有 Schema 信息既不在 Schema RegistrySchema 推断也未开启画像会自动跳过该 Topic。画像任务通过ThreadPoolExecutor并行执行max_workers上限见 kafka.py每个 Worker 持有独立的 Consumer 连接互不阻塞。另外profiling.operation_config支持按周/月计划如每周一、每月 15 日执行低频画像可参考 SQL 类连接器的操作配置语义kafka_config.py。已知限制PROTOBUF 模式类型的限制当前 PROTOBUF 支持存在以下限制不支持递归类型。Protobuf 编译使用进程级全局 descriptor pool第一个编译某个全限定消息类型的 Topic 会成功后续嵌入同一类型名的 Topic 会命中duplicate symbol错误导致该 Topic 没有 Schema 字段DataHub 记录警告后继续处理。在大量 Topic 共享公共 proto 类型的环境中第一个之后的大多数 protobuf Topic 都可能没有 Schema 字段。此时应开启schema_resolution让这些 Topic 回退到基于消息数据的 Schema 推断。此外map 会被表示为消息数组。例如message MessageWithMap { mapint, string map_1 1; }会变成message Map1Entry { int key 1; string value 2/ } message MessageWithMap { repeated Map1Entry map_1 1; }其他注意点Topic 之间的血缘lineage不支持如需 Kafka Connect 场景的血缘请使用 kafka-connect 连接器。状态化摄入Stateful Ingestion仅在为 Source 配置了 Platform Instance 时可用官方能力注释见 kafka_post.md用于陈旧元数据的删除检测。故障排查与性能优化摄入失败时首先依次核对凭证、权限、连通性、范围过滤scope filters然后查看摄入日志中的 Source 专属错误并据此调整配置。Schema 解析错误DataHub 会自动优雅处理 Schema 解析错误并继续处理无需中断。Avro 二进制编码错误avro.errors.InvalidAvroBinaryEncoding: Read 0 bytes, expected 1 bytes开启 Schema 解析以自动推断source: type: kafka config: schema_resolution: enabled: true offset_reset_strategy: hybridProtobuf 重复符号错误Couldnt build proto file into descriptor pool: duplicate symbolDataHub 记录警告并继续处理。存在 Schema 冲突的 Topic 在开启 Schema 解析后会自动使用推断 Schema。Schema Registry 连接问题当 Schema Registry 不可达时可将 Schema 解析作为兜底source: type: kafka config: connection: schema_registry_url: http://localhost:8081 schema_resolution: enabled: true大型集群性能优化对拥有大量 Topic 的 Kafka 集群建议用topic_patterns提前裁剪范围并合理设置画像与解析并发source: type: kafka config: topic_patterns: allow: [prod_.*, analytics_.*] deny: [.*_temp, .*_test] profiling: enabled: true max_workers: 20 # controls both profiling and schema resolution parallelization sample_size: 100 nested_field_max_depth: 10 schema_resolution: enabled: true sample_timeout_seconds: 1.0topic_patterns的默认值在 kafka_config.py 中定义为allow[.*]、deny[^_.*]即默认排除下划线开头的内部 Topic。大 Topic 内存问题对消息体较大或吞吐量高的 Topic降低采样规模与递归深度profiling: enabled: true sample_size: 50 nested_field_max_depth: 2源码侧还为内存保护预设了多项硬性上限kafka_constants.py单个样本值字符串最长 1000 字符、嵌套字典最多展平 100 个键、列表最多 50 个元素、直方图默认 10 个桶、去重值频率最多 10 项避免单个消息撑爆画像结果或内存。验证与深入阅读仓库为 Kafka 连接器配备了完整的单元与集成测试可作为理解行为的补充资料单元测试metadata-ingestion/tests/unit/kafka/test_kafka_source.py、test_kafka_config.py、test_kafka_profiler.py、test_kafka_sampling_strategies.py、test_kafka_schema_inference.py、test_kafka_confluent_catalog.py、test_kafka_protobuf_schema_handling.py集成测试与示例 Recipemetadata-ingestion/tests/integration/kafka/test_kafka.py配套kafka_to_file.yml、kafka_catalog_to_file.yml、kafka_to_file_oauth.yml、kafka_without_schemas_to_file.yml等示例核心实现kafka.py、kafka_config.py、schema_resolution.py、kafka_profiler.py相关模型Dataset 实体、Platform Instance 说明【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表