OpenLineage
OpenLineage 集成
Debezium 内置了与 OpenLineage 的集成,可自动跟踪变更数据捕获(CDC)操作的数据血缘。通过 OpenLineage 集成,你可以全面了解数据流水线中的数据流转与转换过程。
关于数据血缘与 OpenLineage
数据血缘跟踪数据的来源、在系统之间的流转方式,以及所经历的转换过程。OpenLineage 是一个开放标准,Debezium 使用它来发出血缘元数据,并在数据流水线中共享这些元数据。
了解数据血缘对以下活动至关重要:
- 数据治理与合规
- 变更时的影响分析
- 调试数据质量问题
- 理解数据依赖关系
OpenLineage 是一个数据血缘开放标准,提供了一种统一的方式来跨多个数据系统收集和跟踪血缘元数据。该规范定义了一个用于描述数据集(Dataset)、作业(Job)和运行(Run)的通用模型,简化了在异构数据基础设施上构建完整血缘图的过程。
有关更多信息,请参阅 OpenLineage 的网站和文档。
运行事件
Debezium 会发出 OpenLineage 运行事件,以跟踪连接器从初始化到关闭的整个生命周期中的状态。这些事件让你能够了解连接器的健康状况,并支持对数据流水线操作进行实时监控。
连接器在发生以下状态变更后会发出运行事件:
START
报告连接器初始化。
RUNNING
在正常的流式处理操作期间以及处理单个表期间周期性发出。这些周期性事件确保了长时间运行的流式 CDC 操作能够持续被跟踪血缘。
COMPLETE
报告连接器已正常关闭。
FAIL
报告连接器遇到错误。
部署类型
Debezium 的 OpenLineage 集成适用于以下部署类型:
Kafka Connect
你将 Debezium 连接器作为 Kafka Connect 插件运行。
Debezium Server
你以独立服务器模式运行 Debezium。
各种部署类型均提供相同的 OpenLineage 事件模型和功能,但配置集成和安装依赖项所需的操作步骤各不相同。
Debezium 如何与 OpenLineage 集成
Debezium 将其内部对象和事件映射到 OpenLineage 数据模型,其中连接器表示为 Job,数据库表或 Kafka 主题表示为 Dataset。以下各节介绍这些映射在源连接器、接收器连接器以及连接器生命周期中的工作方式。
OpenLineage 作业映射
每个 Debezium 连接器实例都表示为一个 OpenLineage Job,它记录连接器的身份、版本和完整配置,以支持血缘追踪和流水线调试。
Debezium 连接器被映射为一个 OpenLineage Job,它包含以下元素:
名称
作业名称继承自 Debezium 的 topic.prefix,并结合任务 ID(例如 inventory.0)。
命名空间
若指定了 openlineage.integration.job.namespace,则继承自该属性;否则默认使用 topic.prefix 的值。
Debezium 连接器版本
正在运行并生成血缘事件的 Debezium 连接器的版本。
完整的连接器配置
所有连接器配置属性,可实现数据管道的完全可复现性和可调试性。
作业元数据
描述、标签和所有者。
源连接器的数据集映射
对于源连接器,Debezium 会将其监控的数据库表映射为 OpenLineage 输入数据集,将其产生的 Kafka 主题映射为输出数据集,并捕获两者的 schema 信息。
可能的数据集映射如下:
输入数据集
表示 Debezium 配置为捕获变更的数据库表。OpenLineage 集成会根据连接器配置自动创建输入数据集。集成在创建数据集映射时遵循以下原则:
- 连接器监控的每个表都会成为一个输入数据集。
- 每个数据集都会捕获对应源表的 schema 信息,包括每一列的名称和数据类型。
- 源表中的 DDL 变更会动态反映在数据集 schema 中。
输出数据集
表示写入 CDC 事件的 Kafka 主题。对于 Kafka Connect 部署,需要应用 OpenLineage 单消息转换(SMT)才会创建输出数据集。对于 Debezium Server 部署,输出数据集会被自动捕获。
输出数据集映射按照以下原则创建:
- 连接器产生的每个 Kafka 主题都会成为一个输出数据集。
- 输出数据集会捕获完整的 CDC 事件结构,包括元数据字段。
- 数据集的名称基于连接器的主题前缀配置。
接收连接器的数据集映射
对于接收连接器,数据流与源连接器相反。
输入数据集
表示接收连接器读取的 Kafka 主题。这些主题通常包含来自 Debezium 源连接器的 CDC 事件。定义输入数据集时遵循以下原则:
- 接收连接器消费的每个 Kafka 主题都代表一个输入数据集。
- 输入数据集会指定 Kafka 主题的 schema 和元数据。
- 命名空间格式遵循
kafka://bootstrap-server:port,其中 bootstrap 服务器通过openlineage.integration.dataset.kafka.bootstrap.servers属性指定。
输出数据集
表示接收连接器写入数据的目标数据库或集合。
定义输出数据集映射时遵循以下原则:
- 每个目标数据库表或集合代表一个输出数据集。
- 输出数据集指定目标数据库的架构信息。
- 命名空间的格式取决于目标数据库系统。有关更多信息,请参阅数据集命名空间格式。
Kafka Connect 与 Debezium Server 中的 Sink 连接器可用性
可用于 OpenLineage 集成的 sink 连接器因部署平台而异。
在 Kafka Connect 环境中,你可以配置以下 Debezium sink 连接器以使用 OpenLineage 集成:
MongoDB sink 连接器
将 CDC 事件写入 MongoDB 集合。
JDBC sink 连接器
将 CDC 事件写入关系数据库表。
在 Debezium Server 环境中,OpenLineage 集成仅可与 Kafka sink 一起使用。MongoDB sink 和 JDBC sink 连接器不适用于 Debezium Server。路径:integrations/openlineage.adoc
准备部署 Debezium OpenLineage 集成
在使用 Debezium OpenLineage 集成之前,请根据你的部署类型(Kafka Connect 或 Debezium Server)安装所需的依赖项。
所需依赖项
OpenLineage 集成需要若干 JAR 文件,这些文件已打包在 debezium-openlineage-core-libs 归档文件中。
Kafka Connect
要在 Kafka Connect 中使用 Debezium 与 OpenLineage,必须先获取所需的依赖项。
操作步骤
- 下载 OpenLineage 核心归档文件。
- 将归档文件的内容解压到 Kafka Connect 环境中的 Debezium 插件目录中。
Debezium Server
要使用 Debezium Server 与 OpenLineage,必须先获取所需的依赖项。
操作步骤
- 下载 OpenLineage 核心归档文件。
- 解压归档文件的内容。
- 将所有 JAR 文件复制到 Debezium Server 安装目录中的
/debezium/lib目录。
配置集成
要启用集成,必须配置 Debezium 连接器和 OpenLineage 客户端。Kafka Connect 和 Debezium Server 部署的配置方式有所不同。
Kafka Connect 配置
要在 Kafka Connect 中启用 Debezium 与 OpenLineage 的集成,请在连接器配置中添加属性,如下例所示:
# Enable OpenLineage integration
openlineage.integration.enabled=true
# Path to OpenLineage configuration file
openlineage.integration.config.file.path=/path/to/openlineage.yml
# Job metadata (optional but recommended)
openlineage.integration.job.namespace=myNamespace
openlineage.integration.job.description=CDC connector for products database
openlineage.integration.job.tags=env=prod,team=data-engineering
openlineage.integration.job.owners=Alice Smith=maintainer,Bob Johnson=Data EngineerDebezium Server 配置
要启用 Debezium Server 与 OpenLineage 的集成,请按照以下示例,将 OpenLineage 属性添加到 application.properties 文件中。OpenLineage 属性使用 debezium.source. 前缀。
# Enable OpenLineage integration
debezium.source.openlineage.integration.enabled=true
# Path to OpenLineage configuration file
debezium.source.openlineage.integration.config.file.path=config/openlineage.yml
# Job metadata (optional but recommended)
debezium.source.openlineage.integration.job.description=CDC connector for products database
debezium.source.openlineage.integration.job.tags=env=prod,team=data-engineering
debezium.source.openlineage.integration.job.owners=Alice Smith=maintainer,Bob Johnson=Data Engineer配置 OpenLineage 客户端
创建 openlineage.yml 文件来配置 OpenLineage 客户端。openlineage.yml 配置文件在 Kafka Connect 和 Debezium Server 部署中都会用到。
对于 Kafka Connect 部署,请将 openlineage.yml 文件放置在 Kafka Connect 工作节点能够读取的位置。同样,对于 Debezium Server 部署,请将 openlineage.yml 文件放置在 Debezium Server 进程能够读取的位置。
在 Kafka Connect 上存储该文件的常见位置包括 /kafka/openlineage.yml 和 /etc/debezium/openlineage.yml。对于 Debezium Server,通常将该文件存放在 Debezium Server 安装目录的子目录中,例如 config/openlineage.yml。.Procedure . 在相应位置创建 openlineage.yml 文件,并以下面的示例为参考,根据你的需要进行配置:
+
transport:
type: http
url: http://your-openlineage-server:5000
endpoint: /api/v1/lineage
auth:
type: api_key
api_key: your-api-key
# Alternative: Console transport for testing
# transport:
# type: console- 编辑你所用平台的部署配置 YAML 文件中的
config部分,指定openlineage.yml文件的路径。 - 重启 Debezium Server 或 Kafka Connect worker 以应用更改。有关 OpenLineage 客户端配置选项的详细信息,请参阅 OpenLineage 客户端文档。
Debezium OpenLineage 配置属性
下表列出了两种部署类型的 OpenLineage 配置属性。
对于 Debezium Server,所有属性名称都需加上 debezium.source. 前缀(例如 debezium.source.openlineage.integration.enabled)。 |
|---|
| 属性(Kafka Connect) | 说明 | 是否必需 | 默认值 |
|---|---|---|---|
openlineage.integration.enabled | 启用或禁用 OpenLineage 集成。 | 是 | false |
openlineage.integration.config.file.path | OpenLineage YAML 配置文件的路径。 | 是 | 无默认值 |
openlineage.integration.job.namespace | 作业所使用的命名空间。 | 否 | topic.prefix 的值 |
openlineage.integration.job.description | 便于人类阅读的作业描述。 | 否 | 无默认值 |
openlineage.integration.job.tags | 以逗号分隔的键值标签列表。 | 否 | 无默认值 |
openlineage.integration.job.owners | 以逗号分隔的"名称-角色"归属条目列表。 | 否 | 无默认值 |
openlineage.integration.dataset.kafka.bootstrap.servers(仅限源连接器) | 用于获取 Kafka 主题元数据的 Kafka 引导服务器。对于源连接器,如果你未指定该值,则会使用 schema.history.internal.kafka.bootstrap.servers 的值。对于接收器连接器,你必须为该属性指定一个值。 | 是(对于接收器连接器) | schema.history.internal.kafka.bootstrap.servers 的值(仅限源连接器) |
作业元数据增强
为了提升所生成血缘数据的实用性,你可以为每个作业添加元数据,以提供标签和所有者信息。标签是自定义的任意元数据,你可以将其附加到作业上,提供额外的上下文信息以帮助对作业进行分类。所有者则表明谁对该作业负责。请在 OpenLineage 配置属性中使用以下格式指定标签和所有者元数据:
标签列表格式
将标签指定为以逗号分隔的键值对列表,如下例所示:
openlineage.integration.job.tags=environment=production,team=data-platform,criticality=high所有者列表格式
以逗号分隔的「姓名-角色」配对列表来指定所有者,如下例所示:
openlineage.integration.job.owners=John Doe=maintainer,Jane Smith=Data Engineer,Team Lead=owner输出数据集血缘
Debezium 可以捕获输出数据集血缘(Kafka 主题),用于追踪 CDC 事件的目标位置。Kafka Connect 和 Debezium Server 的配置方式有所不同。
Kafka Connect 输出数据集血缘
要在 Kafka Connect 中捕获输出数据集血缘,请将 Debezium 配置为使用 OpenLineage 单消息转换(SMT):
# Add OpenLineage transform
transforms=openlineage
transforms.openlineage.type=io.debezium.transforms.openlineage.OpenLineage
# Required: Configure schema history with Kafka bootstrap servers
schema.history.internal.kafka.bootstrap.servers=your-kafka:9092该 SMT 会捕获 Debezium 写入 Kafka 主题的变更事件的详细 schema 信息。该转换会捕获包含以下内容的 schema 数据:
- 事件结构(before、after、source、事务元数据)
- 字段类型与嵌套结构
- 主题名称与命名空间
Debezium Server 输出数据集血缘
对于 Debezium Server 部署,当启用 OpenLineage 集成时,输出数据集血缘会被自动捕获。无需额外的配置或转换,因为 Debezium Server 可以完全控制输出记录。
完整配置示例
以下示例展示了在 Kafka Connect 和 Debezium Server 中启用 OpenLineage 集成的完整配置。
Kafka Connect 完整配置示例
以下示例展示了在 Kafka Connect 中启用 PostgreSQL 连接器与 OpenLineage 集成的完整配置:
{
"name": "inventory-connector-postgres",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"tasks.max": "1",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "postgres",
"database.password": "postgres",
"database.dbname": "postgres",
"topic.prefix": "inventory",
"snapshot.mode": "initial",
"slot.name": "inventory",
"schema.history.internal.kafka.bootstrap.servers": "kafka:9092",
"schema.history.internal.kafka.topic": "schema-changes.inventory",
"openlineage.integration.enabled": "true",
"openlineage.integration.config.file.path": "/kafka/openlineage.yml",
"openlineage.integration.job.description": "CDC connector for inventory database",
"openlineage.integration.job.tags": "env=production,team=data-platform,database=postgresql",
"openlineage.integration.job.owners": "Data Team=maintainer,Alice Johnson=Data Engineer",
"transforms": "openlineage",
"transforms.openlineage.type": "io.debezium.transforms.openlineage.OpenLineage"
}
}Debezium Server 完整配置示例
以下示例展示了在使用 Kafka sink 的 Debezium Server 中,为 PostgreSQL 连接器启用 OpenLineage 集成所需的完整 application.properties 配置:
# Sink configuration (Kafka)
debezium.sink.type=kafka
debezium.sink.kafka.producer.key.serializer=org.apache.kafka.common.serialization.StringSerializer
debezium.sink.kafka.producer.value.serializer=org.apache.kafka.common.serialization.StringSerializer
debezium.sink.kafka.producer.bootstrap.servers=kafka:9092
# Source connector configuration
debezium.source.connector.class=io.debezium.connector.postgresql.PostgresConnector
debezium.source.offset.storage=org.apache.kafka.connect.storage.MemoryOffsetBackingStore
debezium.source.offset.flush.interval.ms=0
debezium.source.database.hostname=postgres
debezium.source.database.port=5432
debezium.source.database.user=postgres
debezium.source.database.password=postgres
debezium.source.database.dbname=postgres
debezium.source.topic.prefix=tutorial
debezium.source.schema.include.list=inventory
# OpenLineage integration
debezium.source.openlineage.integration.enabled=true
debezium.source.openlineage.integration.config.file.path=config/openlineage.yml
debezium.source.openlineage.integration.job.description=CDC connector for products database
debezium.source.openlineage.integration.job.tags=env=prod,team=cdc
debezium.source.openlineage.integration.job.owners=Mario=maintainer,John Doe=Data scientist
# Logging configuration (optional)
quarkus.log.console.json=falseMongoDB sink 连接器配置示例
以下示例展示了在 Kafka Connect 中启用 MongoDB sink 连接器以与 OpenLineage 集成的完整配置:
{
"name": "mongodb-sink",
"config": {
"connector.class": "io.debezium.connector.mongodb.MongoDbSinkConnector",
"tasks.max": "1",
"mongodb.connection.string": "mongodb://admin:admin@mongodb:27017",
"topics": "inventory.inventory.products",
"sink.database": "inventory2",
"openlineage.integration.enabled": "true",
"openlineage.integration.config.file.path": "/kafka/openlineage.yml",
"openlineage.integration.job.description": "Sink connector for MongoDB",
"openlineage.integration.job.tags": "env=prod,team=cdc",
"openlineage.integration.job.owners": "Mario=maintainer,John Doe=Data scientist",
"openlineage.integration.dataset.kafka.bootstrap.servers": "kafka:9092"
}
}对于接收器连接器,需要 openlineage.integration.dataset.kafka.bootstrap.servers 属性,以便从 Kafka 主题检索输入数据集的元数据。与源连接器不同,接收器连接器无法通过 Kafka Connect 框架直接获取 Kafka 主题的元数据,必须显式建立连接以获取 schema 信息。 |
|---|
JDBC 接收器连接器配置示例
以下示例展示了在 Kafka Connect 中启用 JDBC 接收器连接器并集成 OpenLineage 的完整配置:
{
"name": "jdbc-sink",
"config": {
"connector.class": "io.debezium.connector.jdbc.JdbcSinkConnector",
"tasks.max": "1",
"connection.url": "jdbc:postgresql://postgres:5432/inventory",
"connection.username": "postgres",
"connection.password": "postgres",
"topics": "inventory.inventory.customers",
"insert.mode": "upsert",
"primary.key.mode": "record_key",
"openlineage.integration.enabled": "true",
"openlineage.integration.config.file.path": "/kafka/openlineage.yml",
"openlineage.integration.job.description": "Sink connector for JDBC",
"openlineage.integration.job.tags": "env=prod,team=data-engineering",
"openlineage.integration.job.owners": "Data Team=maintainer,Alice Johnson=Data Engineer",
"openlineage.integration.dataset.kafka.bootstrap.servers": "kafka:9092"
}
}数据集命名空间格式
数据集命名空间用于标识输入数据集和输出数据集的所在位置及来源系统。命名空间的格式因数据库系统和连接器类型(源连接器或接收连接器)而异,遵循 OpenLineage 数据集命名规范。
其他资源
输入数据集命名空间
输入数据集命名空间用于标识源数据库,其格式因数据库系统而异。以下示例说明了 Debezium 为不同数据库系统格式化输入数据集命名空间的方式。
PostgreSQL 输入数据集(用于源连接器)
- 命名空间:
postgres://hostname:port - 名称:
schema.table - Schema:源表中的列名和类型
Kafka 输入数据集(用于接收连接器)
- 命名空间:
kafka://kafka-broker:9092 - 名称:
inventory.inventory.products - Schema:来自源连接器的 CDC 事件结构
具体的命名空间格式取决于你的数据库系统,并遵循 OpenLineage 数据集命名规范。
源连接器的输出数据集命名空间
输出数据集命名空间用于标识写入 CDC 事件的 Kafka 主题。以下示例说明了 Debezium 为源连接器格式化输出数据集命名空间的方式。
Kafka 输出数据集(用于源连接器)
- 命名空间:
kafka://bootstrap-server:port - 名称:
topic-prefix.schema.table - Schema:包含元数据字段的完整 CDC 事件结构
接收连接器的输出数据集命名空间
输出数据集命名空间用于标识接收连接器写入数据的目标数据库。以下示例说明了 Debezium 接收连接器为不同数据库系统格式化输出数据集命名空间的方式。
MongoDB 输出数据集
- 命名空间:
mongodb://mongodb-host:27017 - 名称:
database.collection - Schema:目标集合的 Schema
JDBC 输出数据集(PostgreSQL)
- 命名空间:
postgres://postgres-host:5432 - 名称:
schema.table - Schema:目标表的 Schema
监控与故障排查
你可以通过检查连接器日志、查看 OpenLineage 事件以及解决常见的配置问题,来验证 Debezium OpenLineage 集成是否正常工作并诊断问题。
验证集成
要验证 OpenLineage 集成是否正常工作,请完成以下步骤:
操作步骤
检查连接器日志中与 OpenLineage 相关的消息。
如果你配置了 HTTP 传输,请验证事件是否出现在你的 OpenLineage 后端中。
出于测试目的,你可以配置控制台传输,以便直接在日志中查看事件,如下例所示:
transport: type: console
常见问题
以下参考内容介绍了可能导致 OpenLineage 集成无法正常工作的常见问题,以及相应的解决步骤。
集成无法正常工作
- 确认
openlineage.integration.enabled已设置为true。 - 检查连接器配置中指定的 OpenLineage 配置文件路径是否正确,且 Debezium 能够访问目标文件。
- 确保 OpenLineage 配置文件中的 YAML 语法有效。
- 确认 classpath 中已包含所有必需的 JAR 依赖。
输出数据集缺失
- 确认已将连接器配置为使用 OpenLineage 转换。
- 检查是否已在连接器配置中设置
schema.history.internal.kafka.bootstrap.servers属性。
连接问题
- 确认已在 OpenLineage 客户端配置中指定正确的服务器 URL 和身份验证信息。
- 检查 Debezium 与 OpenLineage 服务器之间的网络连通性。
依赖问题
- 确认所有必需的 JAR 文件均已存在,且其版本兼容。
- 检查是否存在与现有依赖的 classpath 冲突。
接收器连接器的输入数据集缺失
- 确认已配置
openlineage.integration.dataset.kafka.bootstrap.servers属性。 - 确认连接器能够访问 Kafka 引导服务器。
- 确认
topics配置中指定的 Kafka 主题存在,且连接器能够访问这些主题。
错误事件
当连接器发生故障时,请在 OpenLineage 的 FAIL 事件中检查以下内容:
- 错误消息
- 堆栈跟踪
- 用于调试的连接器配置
评论
登录后参与评论
KnowForge