OpenMetadata Kinesis 连接器接入指南:AWS 权限配置、连接参数详解与流元数据摄取原理
2026/9/15 11:42:55 网站建设 项目流程

OpenMetadata Kinesis 连接器接入指南:AWS 权限配置、连接参数详解与流元数据摄取原理

【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata

OpenMetadata 通过 Kinesis 连接器,将 AWS Kinesis Data Streams 中的流(Stream)以 Topic 实体的形式纳入统一元数据目录,同时支持流内样本数据(Sample Data)采集,为数据资产检索、血缘追踪与 AI 语义层提供实时流数据的基础上下文。本文以 Kinesis 连接器文档 为主体,结合仓库内连接器源码实现,完整讲解 AWS 最小权限策略、全部连接参数的语义与取值范围、摄取流程的底层调用链,并给出可复用的工作流配置示例。读完本文,你将能够独立完成 Kinesis 服务的连接配置、权限校验与元数据摄取。

Kinesis 连接器能做什么

Kinesis 连接器通过 AWS 官方 SDK(boto3)客户端访问 Kinesis Data Streams,完成两类核心工作:

  • 元数据摄取:列出 AWS 账户下全部流,将每个流建模为 OpenMetadata 的Topic实体,记录分区(Shard)数量、保留时长(Retention Period)、最大消息大小等属性,并生成可跳转到 AWS 控制台对应流详情页的sourceUrl
  • 样本数据采集:按需从流中读取少量消息(Sample Data),帮助用户在元数据目录中直观预览流内数据结构,为后续 schema 识别与数据质量评估打基础。

从源码结构看,Kinesis 连接器由四个文件构成完整的服务规格(Service Spec)注册,参见 service_spec.py:

  • KinesisSource:元数据摄取主逻辑(metadata.py);
  • KinesisConnection:连接建立与连通性测试(connection.py);
  • KinesisSampler:样本数据采样器(当前为未实现的占位实现,见 sampler.py);
  • 各 API 响应对应的 pydantic 数据模型(models.py)。

前提条件:AWS 最小权限策略

Kinesis 连接器在摄取过程中会依次调用五类 Kinesis API,因此运行该连接器的 AWS 身份(IAM 用户或角色)必须被授予对应权限。官方推荐的最小权限策略如下:

{ "Version": "2012-10-17", "Statement": [ { "Sid": "KinesisPolicy", "Effect": "Allow", "Action": [ "kinesis:ListStreams", "kinesis:DescribeStreamSummary", "kinesis:ListShards", "kinesis:GetShardIterator", "kinesis:GetRecords" ], "Resource": "*" } ] }

这五类权限与源码中的实际调用一一对应,可逐一核对:

权限对应 API在摄取流程中的用途
kinesis:ListStreamslist_streams分页枚举账户下所有流名称,见 metadata.py
kinesis:DescribeStreamSummarydescribe_stream_summary获取流的摘要信息(如RetentionPeriodHours),用于计算 Topic 的保留时长
kinesis:ListShardslist_shards分页获取流的 Shard(分区)列表,决定 Topic 的partitions数量
kinesis:GetShardIteratorget_shard_iterator获取 Shard 迭代器,用于定位样本数据读取起点
kinesis:GetRecordsget_records从 Shard 迭代器读取实际消息记录,用于样本数据采集

权限授予后,连接器即可读取流元数据与样本数据。Kinesis 权限的更多说明可参考 AWS 官方文档(见原文档所引的 AWS Kinesis 权限控制指引)。

连接配置参数详解

配置 Kinesis 连接时,除AWS Region 为唯一必填项外,其余参数均可根据你的认证方式按需填写。以下逐一说明各参数的语义与适用场景。

AWS Access Key ID

访问 AWS 时,你需要提供安全凭证用于身份认证与请求授权。访问密钥由两部分组成:

  1. 访问密钥 ID,例如AKIAIOSFODNN7EXAMPLE
  2. 秘密访问密钥,例如wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY

两者必须成对使用才能完成认证。你可以按需轮换或管理访问密钥,具体操作参见 AWS IAM 官方文档。

AWS Secret Access Key

即上述访问密钥中的秘密访问密钥部分,与 Access Key ID 配套填写。

AWS Region

每个 AWS Region 是一个独立的地理区域,AWS 在其中部署数据中心。由于同一服务可能部署在多个 Region,连接器必须知道目标服务所属的 Region 才能正确寻址。该参数是配置 Kinesis 连接时唯一必填的参数

Region 值同时参与两处逻辑:

  • 构造 AWS 客户端(通过AWSClient统一封装,见 connection.py);
  • 生成指向 AWS 控制台的流详情链接sourceUrl,形如https://{awsRegion}.console.aws.amazon.com/kinesis/home?region={awsRegion}#/streams/details/{streamName}/monitoring,见 metadata.py。

AWS Session Token

如果你使用临时凭证访问服务,除了 Access Key ID 与 Secret Access Key 外,还需要填写 AWS Session Token。临时凭证通常由 STS(Security Token Service)颁发,用于短期授权的场景。其使用规范参见 AWS 官方文档“Using temporary credentials with AWS resources”。

Endpoint URL

端点(Endpoint)是 AWS Web 服务的入口 URL。AWS SDK 与 AWS CLI 默认使用每个服务在对应 Region 的标准端点,但你可以为 API 请求指定替代端点。此参数在以下场景尤其有用:

  • 连接兼容 Kinesis API 的第三方实现(如本地模拟服务 LocalStack);
  • 通过代理或自定义网关访问 AWS。

各服务的端点列表参见 AWS 官方“AWS service endpoints”文档。

Profile Name

命名配置文件(Named Profile)是存储在 AWS 配置文件与凭证文件中的一组设置与凭证集合。指定 Profile Name 后,连接器将使用该配置文件中的凭证与设置执行命令,而非默认的default配置。若你的环境使用了多个 AWS 身份,可在此填写目标身份对应的 Profile 名称。更多说明参见 AWS CLI 官方“Named profiles”文档。

Assume Role ARN

AssumeRole常用于账户内角色切换或跨账户访问。在此字段填写目标角色(其他账户)的 ARN(Amazon Resource Name)。

需要注意:希望切换到其他账户角色的用户,必须获得账户管理员授予的权限——管理员需附加一条允许该用户以目标账户角色 ARN 调用AssumeRole的策略。如果你打算使用AssumeRole,该字段为必填。相关细节参见 AWS STSAssumeRoleAPI 文档。

Assume Role Session Name

被假定角色会话(Assumed Role Session)的标识符。当同一角色被不同主体或因不同目的假定(Assume)时,使用角色会话名称可唯一标识每次会话。默认值为OpenMetadataSession。更多说明参见 AWS STSAssumeRole文档中关于 Role Session Name 的部分。

Assume Role Source Identity

由调用AssumeRole操作的主体指定的源身份(Source Identity)。该信息会写入 AWS CloudTrail 日志,用于追溯“是谁通过该角色执行了操作”。若你需要在审计日志中区分实际操作者,可在此字段填写源身份标识。参见 AWS STSAssumeRole文档中 Source Identity 的说明。

认证优先级与 boto3 凭证解析

除显式填写的凭证外,boto3 客户端还会按照标准顺序自动解析环境变量、~/.aws/credentials~/.aws/config文件中的配置。因此实际部署中存在三种常见组合:

  • 仅填 Region:凭证全部来自环境变量或本机默认凭证文件,适合在已配置 AWS CLI 的主机上运行;
  • 显式 AK/SK:在连接配置中直接填写 Access Key ID 与 Secret Access Key,适合凭证集中管理或跨环境迁移;
  • AssumeRole 组合:填写 Assume Role ARN 与 Session Name,由连接器代为切换角色,适合跨账户元数据摄取。

源码中,所有 AWS 凭证的解析与 Kinesis 客户端创建统一收敛在AWSClient封装内(connection.py中通过AWSClient(connection.awsConfig).get_kinesis_client()获取客户端),各 Region 下的附加配置解析方式遵循 boto3 官方“Configuring credentials”文档。

摄取流程的源码级拆解

1. 流列表分页拉取

get_stream_names_listLimit=100为页大小循环调用list_streams(见 metadata.py):每次取出StreamNames后追加到结果集,若HasMoreStreams为真,则用最后一个流名称作为ExclusiveStartStreamName继续下一页,直至全部拉取完毕。

2. 逐流组装 Topic 元数据

get_topic_list对每个流并行获取两类信息(见 metadata.py):

  • 流摘要:调用describe_stream_summary,重点读取StreamDescriptionSummary.RetentionPeriodHours
  • 分区列表:调用list_shards,通过NextToken处理分页(当响应中无NextToken时显式置为None以结束循环)。

随后yield_topic将这些信息映射为CreateTopicRequest(见 metadata.py):

  • partitions:取 Shard 列表长度,即分区数;
  • retentionTime:将RetentionPeriodHours(小时)换算为毫秒,hours * 3_600_000,见_compute_retention_time
  • maximumMessageSize:固定为1_000_000字节(源码常量MAX_MESSAGE_SIZE),对应 Kinesis 单条记录上限;
  • sourceUrl:按 Region 与流名拼装 AWS 控制台监控页链接。

3. 样本数据采集

当工作流开启generateSampleData且未在全局层面被禁用时(见 metadata.py),yield_topic_sample_data会对每个分区执行:

  1. TRIM_HORIZON迭代器类型调用get_shard_iterator获取迭代器(对应KinesisEnum.TRIM_HORIZON);
  2. 调用get_records读取记录;
  3. 对每条记录的Data优先尝试 Base64 解码后按 UTF-8 转字符串,解码失败(binascii.ErrorUnicodeDecodeError)时回退为直接 UTF-8 解码,见_get_sample_records(metadata.py)。

任意分区取到样本数据后即停止,避免对流的过度读取。

4. 连通性测试

KinesisConnection.test_connectionlist_streams作为唯一的测试步骤(见 connection.py):只要凭证有效且具备kinesis:ListStreams权限,即判定连接成功。该测试既可在元数据工作流中执行,也可通过自动化工作流(Automation Workflow)触发,默认超时时间为 3 分钟。

5. 数据模型与注册

所有 Kinesis API 响应均通过 pydantic 模型一次性转换(见 models.py),模型采用ConfigDict(extra="allow")以容忍 AWS 返回的多余字段。连接器整体通过BaseSpec注册为服务规格,供 OpenMetadata 服务目录动态加载。

完整工作流配置示例

以下 YAML 给出一个可直接参考的 Kinesis 元数据摄取工作流配置(字段与上文参数一一对应):

source: type: kinesis serviceName: kinesis_production serviceConnection: config: type: Kinesis awsConfig: awsAccessKeyId: AKIAIOSFODNN7EXAMPLE awsSecretAccessKey: wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY awsRegion: us-east-1 awsSessionToken: "" # 使用临时凭证时必填 endPointURL: "" # 自定义端点,如 LocalStack profileName: "" # 使用非 default 命名配置时填写 assumeRoleArn: "" # 使用 AssumeRole 时必填 assumeRoleSessionName: OpenMetadataSession # 默认值 assumeRoleSourceIdentity: "" # 可选,用于 CloudTrail 审计溯源 sourceConfig: config: type: MessagingMetadata generateSampleData: true # 是否采集流样本数据 sink: type: metadata-rest config: {} workflowConfig: openMetadataServerConfig: hostPort: http://localhost:8585/api authProvider: openmetadata

配置完成后,可在 OpenMetadata UI 的 Services 页面新建 Kinesis 服务并执行“Test Connection”,通过后再创建摄取管道运行元数据同步。摄取完成后,每条 Kinesis 流将以 Topic 实体呈现于数据资产目录,包含分区数、保留时长、最大消息大小等属性与样本数据预览。

小结

Kinesis 连接器的接入要点可归纳为三件事:给足权限(五类 Kinesis API 的最小策略)、选对认证方式(Region 必填,其余按凭证来源组合填写)、按需开启样本数据generateSampleData)。源码层面,连接器围绕AWSClient统一封装凭证、以 pydantic 模型承接 API 响应、通过BaseSpec注册服务规格,其摄取调用链清晰且易于扩展——若需要为 Kinesis 流补充 schema 自动分类等能力,KinesisSampler的占位实现(sampler.py)即是一个明确的切入点。

【免费下载链接】OpenMetadataThe Open Context Layer for Data and AI , OpenMetadata is the open platform for building trusted data context and business semantics for humans, AI assistants, and agents.项目地址: https://gitcode.com/GitHub_Trending/op/OpenMetadata

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询