Amazon Kinesis 数据源深度指南:Firehose 跨平台血缘、Glue Schema Registry 与排障实战
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
本文基于 DataHub 官方
kinesis数据源文档(kinesis_post.md)并结合仓库源码、配置与测试用例编写,全面覆盖 Kinesis 连接器的 Capabilities、Firehose 跨平台血缘、Glue 表血缘、GSR 模式注册、过滤、标签所有权派生、Limitations 与 Troubleshooting。
引言:为什么需要一份面向 Kinesis 连接器的深度指南
AWS 实时数据管道中,Kinesis Data Streams(KDS)负责缓冲与分发,Amazon Data Firehose 负责将数据投递到 S3、Redshift、Snowflake、Iceberg、MongoDB、OpenSearch 等目标平台。DataHub 的kinesis连接器用一个 recipe、一份 IAM 策略、一个 ingestion job 同时接入这两类服务:
- KDS 流以Dataset(
Stream子类型)入库,携带区域 Container、StreamARN、分片数、保留期、加密与流模式等 custom properties; - Firehose 流以DataFlow(
Firehose Stream子类型)入库,内含单个DataJob(Delivery子类型),其dataJobInputOutput血缘边连接源 Kinesis 流与目标平台。
因此,本文深入讲解:
- Firehose 六种支持的目的地及其 URN 格式;
destination_platform_map跨平台 URN 覆盖机制;- Glue 表血缘(Parquet/ORC 格式转换);
- Glue Schema Registry(GSR)模式解析与
use_naming_convention的取舍; - 流/Firehose 过滤、标签所有权派生;
- 已知限制与常见故障排查。
读完本文,你将能正确配置 Kinesis 连接器、避免“血缘边存在但目标 URN 无效”的坑,并理解其内部实现。
Capabilities 总览
Kinesis 连接器在源码KinesisSource类上声明了以下能力(kinesis.py):
| Capability | 说明 | 状态 |
|---|---|---|
| Descriptions | 默认启用,生成实体描述 | Enabled by default |
| Containers | 区域级 Container | Region containers |
| Lineage (coarse) | Firehose → 目标平台血缘 | Firehose -> destination lineage |
| Tags | 从 AWS 资源标签生成 DataHub globalTags | From AWS resource tags |
| Schema Metadata | 通过 GSR 解析 schema(Avro / JSON / Protobuf) | Opt-in viaglue_schema_registry.enabled |
| Deletion Detection | 基于 stateful ingestion 的软删除 | Via stateful ingestion |
从源码结构看,连接器模块位于 metadata-ingestion/src/datahub/ingestion/source/kinesis/,包含:
kinesis.py—— 主入口与 region container、account_id 解析;kinesis_config.py—— 配置模型与destination_platform_map/ GSR 配置;kinesis_stream.py—— KDS 提取器;kinesis_firehose.py—— Firehose 提取器与 DataJob 血缘组装;kinesis_firehose_destinations.py—— 目的地 URN 处理器注册表;kinesis_schema_registry.py—— GSR 解析器;kinesis_tagging.py—— AWS 标签 → globalTags 工具;kinesis_report.py—— 报告字段。
Firehose 血缘:支持的目的地
六种支持的 AWS 目的地
Firehose 投递目标平台会产生dataJobInputOutput.outputDatasets血缘边。以下目的地类型受支持(文档表格):
| AWS destination | DataHub 平台 | URN 格式 |
|---|---|---|
| Amazon S3 / Extended S3 | s3 | urn:li:dataset:(urn:li:dataPlatform:s3,<bucket>[/<prefix>],...) |
| Amazon Redshift | redshift | urn:li:dataset:(urn:li:dataPlatform:redshift,<db>.<schema>.<table>,...) |
| Amazon OpenSearch / Elasticsearch | elasticsearch | urn:li:dataset:(urn:li:dataPlatform:elasticsearch,<index>,...) |
| Snowflake | snowflake | urn:li:dataset:(urn:li:dataPlatform:snowflake,<db>.<schema>.<table>,...) |
| Apache Iceberg | iceberg | urn:li:dataset:(urn:li:dataPlatform:iceberg,<namespace>.<table>,...) |
| MongoDB | mongodb | urn:li:dataset:(urn:li:dataPlatform:mongodb,<database>.<collection>,...) |
URN 名称默认统一小写(与 Snowflake 源默认convert_urns_to_lowercase=True保持一致)。如需按目的地覆盖,通过destination_platform_map.<platform>.convert_urns_to_lowercase: false关闭。
注意:不支持的 Firehose 目的地(HTTP、Datadog、Splunk、New Relic、Coralogix、LogicMonitor、Dynatrace、Honeycomb、Sumo Logic 等)不会产生血缘边 —— 连接器记录一条 “Unsupported Firehose destination” 警告,并将目的地配置作为 custom property 记在 DataJob 上,DataJob 本身仍会生成。
源码实现印证
在 kinesis_firehose_destinations.py 中,每个目的地都有对应的DestinationHandler子类(S3Destination/ExtendedS3Destination/RedshiftDestination/OpenSearchDestination/SnowflakeDestination/IcebergDestination/MongoDBDestination),它们通过matches()匹配 boto3DescribeDeliveryStream返回的 destination block,再用build_urns()构造 URN。DESTINATION_HANDLERS注册表按“先匹配 modern S3、再匹配 legacy”的顺序排列。
关键细节(源码注释明确):
- S3 处理器同时处理 legacy
S3DestinationDescription与 modernExtendedS3DestinationDescription,二者共享BucketARN+Prefix键;_build_s3_urn在 BucketARN 缺失时返回空列表,防御urn:li:dataset:(s3,,PROD)这类畸形 URN。 - Redshift 处理器从
CopyCommand.DataTableName取<schema>.<table>,从ClusterJDBCURL解析数据库名;数据库或表名缺失时拒绝生成 URN(避免urn:li:dataset:(redshift,my_table,PROD)这种语法合法但语义错误的 URN)。 - OpenSearch 处理器同时匹配 boto3 的
AmazonopensearchserviceDestinationDescription(老账户的 AWS 公共 API 名)与ElasticsearchDestinationDescription。 - Iceberg 是唯一支持“单条投递流 → 多张目标表”的目的地(
DestinationTableConfigurationList是列表),V1 假设 REST catalog 的点分隔 namespace 格式。
跨平台血缘:destination_platform_map
Firehose 目的地位于其他平台(S3、Redshift、Snowflake 等),因此 Kinesis 产生的血缘 URN 必须与这些平台自己的 DataHub 源的 URN 约定一致。destination_platform_map允许按目的地覆盖 URN 参数:
destination_platform_map: snowflake: platform_instance: "prod-snowflake-east" env: "PROD" # 仅当 Snowflake 源 recipe 也设置了 convert_urns_to_lowercase: false 时才需要 convert_urns_to_lowercase: false redshift: platform_instance: "analytics-cluster" env: "PROD" iceberg: # Iceberg catalog 区分大小写 —— 关闭小写化以保留原大小写 convert_urns_to_lowercase: false每个目的地平台有三个可调参数:
platform_instance—— 必须与目的地平台自己的源 recipe 中使用的字符串一致。否则 Firehose 血缘边会指向 DataHub UI 中无法解析的死 URN。env——PROD/DEV等。未设置时继承本源的env。convert_urns_to_lowercase—— 默认true。对于区分大小写的目的地(Iceberg、MongoDB)或使用自身convert_urns_to_lowercase=false导入的 Snowflake / Redshift 源,设为false。
源码实现印证
在 kinesis_config.py 中,DestinationPlatformDetail定义了这三个字段,其中env会被自动大写并校验为合法 FabricType(PROD/DEV/QA/STG 等),与顶层env的处理保持一致。
destination_platform_map的键被约束为DestinationPlatformLiteral(s3/redshift/elasticsearch/snowflake/iceberg/mongodb/glue),未知键在解析时直接报错,从源头杜绝拼写错误。
在 kinesis_firehose.py 的_destination_urn中实现了覆盖逻辑:
detail = self.config.destination_platform_map.get(platform) platform_instance = detail.platform_instance if detail else None resolved_env = env or (detail.env if detail and detail.env else self.config.env) if detail is None or detail.convert_urns_to_lowercase: name = name.lower() return make_dataset_urn_with_platform_instance(...)注意:convert_urns_to_lowercase的默认值同样是true(对齐 Snowflake 源),只有当 map 中显式设置false时才跳过小写化。
测试验证
test_kinesis_firehose.py 中test_destination_platform_map_overrides_snowflake_instance验证了platform_instance会被折叠进 URN 名称前缀(prod-sf.db.s.t);TestDestinationUrnCaseFolding类则专门验证默认小写化与按目的地关闭小写化的行为。配置层面的校验在 test_kinesis_config.py:未知平台(如datadog)会在启动时报错,glue作为合法的 map 键也被显式支持。
Glue 表血缘(Firehose 格式转换)
当 Firehose 流启用了Parquet/ORC 格式转换时,其SchemaConfiguration会引用一张 Glue 表来定义输出 schema。连接器会把它作为 Firehose delivery DataJob 的第二个上游输入暴露出来:
- 源 Kinesis 流仍是输入(既有行为);
- 目标 S3 路径仍是输出(既有行为);
- Glue 表被添加为第二个输入 —— 它的 schema 决定了写入 S3 路径的内容。
要让 Glue 表 URN 被发出,SchemaConfiguration必须包含DatabaseName和TableName。CatalogId不是必须的—— 当它与调用者账户相同时,AWS 会在DescribeDeliveryStream响应中省略它(按 AWS 文档,它是 input-side default)。存在SchemaConfiguration但缺少DatabaseName/TableName的情况,会被记录到 source report 的firehose_glue_schema_skipped字段中用于诊断。
如果 Glue catalog 是在非默认platform_instance下导入的,需要设置覆盖:
destination_platform_map: glue: platform_instance: "central-catalog" env: "PROD"整个行为由include_table_lineage标志控制 —— 关闭时不会发出任何 Glue 血缘。
源码实现印证
在 kinesis_firehose_destinations.py 中,ExtendedS3Destination.extract_schema_config_glue_urn从DataFormatConversionConfiguration.SchemaConfiguration读取DatabaseName/TableName,构造urn:li:dataset:(glue,<db>.<table>,<env>)格式的 URN;缺失字段时通过report_firehose_glue_schema_skipped记录原因。在 kinesis_firehose.py 中,_process_destination会对ExtendedS3Destination调用该方法,并把返回的 Glue URN追加到 inputs(上游)—— 因为该表的 schema 决定写入内容,S3 路径仍是数据目的地。
流与 Firehose 流过滤
stream_pattern和firehose_stream_pattern使用标准 DataHubAllowDenyPattern结构。一个常见的 deny 规则用于排除内部 / 审计 / 调试流:
stream_pattern: deny: - "^_.*" - ".*-debug$"stream_pattern过滤 Kinesis Data Streams(include_streams: true时生效);firehose_stream_pattern过滤 Firehose 流(include_firehose: true时生效)。
被过滤的流会记录在 report 的filtered_streams/filtered_firehose_streams字段中(kinesis_report.py)。
从标签派生所有权
连接器将 AWS 资源标签发布为 DataHubglobalTags(Key=Value变为urn:li:tag:Key:Value;仅 Key 的标签变为urn:li:tag:Key)。要把标签转成所有权,请应用内置的extract_ownership_from_tagstransformer —— 这使所有权处理与 DataHub 其他源保持一致,而不是在每个连接器里重新实现。
例如,把owner标签的值当作 corpuser 所有者:
transformers: - type: "extract_ownership_from_tags" config: tag_pattern: "owner:"该 transformer 还支持 corp group、owner types,以及追加 email domain —— 完整选项见 dataset_transformer.md(如owner_type、owner_type_urn、email_domain、extract_owner_type_from_tag_pattern等配置项)。
从源码看,所有权刻意不在连接器内派生 —— kinesis_tagging.py 的注释明确说明:它由
extract_ownership_from_tagstransformer 通用地处理,连接器只负责把 AWS 标签展平为GlobalTagsClass。
Glue Schema Registry(GSR)
GSR 是可选开启的,因为它需要额外的 IAM 权限(glue:Get*/glue:List*)。
开启与关闭的差别:
- 关闭(默认)—— 流在没有
schemaMetadataaspect 的情况下发出。所有其他元数据(properties、tags、ownership、lineage)不受影响,也不需要glue:*权限。 - 开启—— 对每个能解析出 schema 的流(见下面的解析顺序),连接器从 AWS Glue Schema Registry 获取 schema 并附带
schemaMetadataaspect(解析字段支持 Avro / JSON / Protobuf)。解析不到 schema 的流仍会发出,只是没有schemaMetadata—— 开启 GSR 永远不会丢流。需要GlueSchemaRegistryReadIAM statement。
开启配置:
glue_schema_registry: enabled: true registry_name: "default-registry" # 推荐:显式声明已知的 stream -> schema 关联 stream_schema_map: events: "events-v2" clicks: "click-events-schema" # 可选启发式 —— 见下文说明 use_naming_convention: false每个流的 schema 解析顺序:
- 如果流名是
stream_schema_map的键,使用映射的 schema 名; - 否则,如果
use_naming_convention: true,在配置的registry_name中查找与流名同名的 schema; - 否则,不带
schemaMetadata发出该流。
为什么 use_naming_convention 默认关闭?
与 Kafka + Confluent Schema Registry(定义了标准化的TopicNameStrategy,即<topic>-key/-valuesubject 命名)不同,AWS 没有定义 Kinesis Data Stream 与 Glue schema 之间的任何关系。schema 由生产者按记录选择(GlueSchemaRegistrySerializer会把 schema-id 嵌入每条记录),流本身没有 schema 绑定。多个生产者可以用不同 schema 写入同一个流,一个 schema 也可以被多个流复用。
有些组织把 “schema 名 == 流名” 作为内部约定,但这并非 AWS 最佳实践。如果你的组织遵循该约定,可以设置use_naming_convention: true;否则,在stream_schema_map中声明已知关联 —— 这是最可预测的模式。
(对于启用了 Parquet/ORC 格式转换的Firehose流,AWS 确实通过SchemaConfiguration定义了关系 —— 见上文 Glue 表血缘。该提取默认开启,不受此标志影响。)
源码实现印证
kinesis_config.py 中的KinesisGlueSchemaRegistryConfig有一个重要的 model validator:当enabled=false却设置了stream_schema_map或use_naming_convention时,直接抛配置错误—— 因为这两个激活开关在禁用状态下是静默无效的,用户几乎必然是“本想开启 GSR”。registry_name因有非空默认值而不作为激活信号。
kinesis_schema_registry.py 实现了解析逻辑:
_resolve_schema_name先查stream_schema_map,再按命名约定回退;get_schema_metadata调用glue:GetSchemaVersion(SchemaVersionNumber={"LatestVersion": True});- 命名约定探针未命中(
EntityNotFoundException且非显式映射)是预期结果,记录到gsr_naming_convention_misses字段而不产生 WARNING 噪音;而显式映射失败、AccessDenied、ValidationException 等真实错误会记入schema_resolution_failures并产生警告。
测试方面,test_kinesis_config.py 覆盖了:禁用时get_schema_metadata短路(不做任何 Glue 调用)、命名约定关闭时不解析、以及“配置了但未启用”在启动即被拒绝。
已知限制(Limitations)
对于任何使用非默认
platform_instance导入的目的地平台,destination_platform_map是必须的。否则 Firehose 血缘边会引用语法合法但在 DataHub 中解析不到任何东西的 URN —— 在发出的 JSON 里血缘看起来正确,但 UI 中目标是死链。始终为设置了platform_instance的目的地填充destination_platform_map:destination_platform_map: snowflake: platform_instance: "prod-snowflake-east" redshift: platform_instance: "analytics-cluster"Glue Schema Registry 的跨 schema 引用不会被解析。带有
$ref(JSON Schema)或命名导入(Avro / Protobuf)的 schema 只发出顶层 schema —— 嵌套引用不会被追踪。依赖跨 schema 导入的流会缺失部分字段。每个 recipe 一个区域。连接器每次运行只导入一个 AWS 区域。多区域账户需要为每个区域运行一个 recipe,并使用不同的
platform_instance。每个 recipe 一个 Glue Schema Registry。只查询
glue_schema_registry.registry_name指定的 registry。如果 schema 跨多个 registry,需要运行多个 recipe。没有
USAGE_STATS能力。每个流的读写吞吐量在 CloudWatch 中可用,但此连接器不读取。不支持的 Firehose 目的地。Splunk、HTTP、Datadog、New Relic、Coralogix、LogicMonitor、Dynatrace、Honeycomb、Sumo Logic 等目的地会发出没有输出血缘边的 DataJob,并在 report 中显示警告。六个受支持的目的地见上文的 Firehose 血缘表。
不从记录采样推断 schema。连接器依赖 AWS Glue Schema Registry 获取 schema —— 没有注册 schema 的流会在没有
schemaMetadata的情况下发出。没有 Lambda 消费者发现。连接器不会枚举消费某个流的 Lambda 函数,因此不产生
Stream → Lambda血缘。不支持 Kinesis Data Analytics(KDA / 托管 Flink)。该连接器不导入 KDA 应用。
Troubleshooting 排障指南
血缘边存在,但目标数据集在 DataHub 中找不到。你的destination_platform_map与目的地平台自身的导入设置不匹配。检查destination_platform_map.<platform>中的platform_instance和env是否与该平台源 recipe 使用的值完全一致。参见限制 #1。
Snowflake 或 Redshift 血缘无法解析,且目的地标识符为混合大小写。连接器默认将目的地 URN 名称小写化。如果你的 Snowflake / Redshift 源 recipe 设置了convert_urns_to_lowercase: false,请在 Kinesis 侧同样设置:
destination_platform_map: snowflake: convert_urns_to_lowercase: falseIceberg 或 MongoDB 血缘无法解析。这两个平台的标识符区分大小写。按上述方法为它们关闭小写化(convert_urns_to_lowercase: false)。
导入成功但搜索不显示实体。你的 DataHub 后端可能因搜索索引磁盘压力进入只读状态。在运行 DataHub 的主机上:
docker exec <opensearch-or-elasticsearch-container> \ curl -s localhost:9200/_cluster/allocation/explain如果集群报告超出disk.watermark.low,请释放磁盘空间并重新索引。
kinesis:ListStreams报AccessDeniedException。recipe 使用的 IAM 身份缺少 AWS IAM 权限 一节中的KinesisDataStreamsRead权限。首页拒绝会记录为警告并跳过 KDS 部分(用户可能故意只有 Firehose IAM);分页中途失败会升级为report.failure,以防止对未列出的流进行有状态软删除。
Firehose 部分静默为空,尽管 Firehose 流存在。IAM 身份缺少firehose:ListDeliveryStreams和/或firehose:DescribeDeliveryStream—— Kinesis Data Streams 权限不覆盖 Firehose。Firehose 权限缺失会被记录为警告(“Permission denied for Firehose”)而非失败;请检查导入运行报告。
流的 schema 找不到。glue:GetSchemaVersion对预期的 schema 名返回了EntityNotFoundException。要么在glue_schema_registry.stream_schema_map中添加显式条目,要么将 GSR schema 重命名为与流名一致(并开启use_naming_convention: true),要么接受该流将在没有schemaMetadata的情况下发出。
与 kinesis_pre.md / recipe 的衔接
本文聚焦于连接器的 Capabilities、配置细节与排障。前置的 IAM 权限、认证方式与最小 recipe 参见 kinesis_pre.md(IAM 策略包含kinesis:ListStreams/kinesis:DescribeStream/kinesis:ListTagsForStream,以及可选的 Firehose 与 Glue 权限;认证遵循 boto3 标准链:静态凭证 → 环境变量 → named profile → EC2/ECS/EKS IAM 角色 → SSO profile)。一个带注释、可运行的完整 recipe 示例在 kinesis_recipe.yml。
小结
kinesis连接器是 DataHub 将 AWS 实时数据管道纳入元数据治理的核心入口。理解 Firehose 目的地 URN 的构造规则、destination_platform_map的覆盖机制、GSR 的解析顺序与大小写折叠的默认行为,是避免“血缘边看起来正确但 UI 中全是死链”的关键。配合源码中的处理器注册表与单元测试,你可以自信地排查任何 lineage 或 schema 相关的问题。
【免费下载链接】datahubThe Context Platform for your Data and AI Stack项目地址: https://gitcode.com/GitHub_Trending/da/datahub
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考