Kafka Connect
教程
本教程演示如何使用 Debezium 监控 MySQL 数据库。随着数据库中数据的变化,你将看到由此产生的事件流。
在本教程中,你将启动 Debezium 服务、运行一个带有简单示例数据库的 MySQL 数据库服务器,并使用 Debezium 监控该数据库的变化。
前提条件
已安装并正在运行 Docker。
本教程使用 Docker 和 Debezium 容器镜像来运行所需的服务。你应使用最新版本的 Docker。有关更多信息,请参阅 Docker Engine 安装文档。
| 本示例也可以使用 Podman 运行。有关更多信息,请参阅 Podman。 |
|---|
Debezium 简介
Debezium 是一个分布式平台,可将现有数据库中的信息转换为事件流,使应用程序能够检测数据库中的行级变化并立即做出响应。
Debezium 构建在 Apache Kafka 之上,提供了一组与 Kafka Connect 兼容的连接器。每个连接器适用于特定的数据库管理系统(DBMS)。连接器通过在数据变化发生时检测这些变化,并将每个变更事件的记录以流的形式发送到 Kafka 主题,从而记录 DBMS 中的数据变更历史。消费应用程序随后可以从 Kafka 主题中读取生成的事件记录。
借助 Kafka 可靠的流式平台,Debezium 使应用程序能够正确、完整地消费数据库中发生的变化。即使你的应用程序意外停止或断开连接,它也不会错过停机期间发生的事件。应用程序重新启动后,它会从上次中断的位置继续从主题中读取。
接下来的教程将向你展示如何使用简单的配置部署和使用 Debezium MySQL 连接器。有关部署和使用 Debezium 连接器的更多信息,请参阅连接器文档。
其他资源
- Debezium 连接器(Cassandra)
- Debezium 连接器(Db2)
- Debezium 连接器(MongoDB)
- Debezium 连接器(MySQL)
- Debezium 连接器(Oracle Database)
- Debezium 连接器(PostgreSQL)
- Debezium 连接器(SQL Server)
- Debezium 连接器(Vitess)
启动服务
使用 Debezium 需要两个独立的服务:Kafka 和 Debezium 连接器服务。在本教程中,你将使用 Docker 和 Debezium 容器镜像 各搭建一个服务实例。
要启动本教程所需的服务,你必须:
使用 Docker 运行 Debezium 的注意事项
本教程使用 Docker 和 Debezium 容器镜像 来运行 Kafka、Debezium 和 MySQL 服务。将每个服务运行在独立的容器中,可以简化配置过程,让你直观地看到 Debezium 的实际运行效果。
| 在生产环境中,你会运行每个服务的多个实例,以提供性能、可靠性、复制和容错能力。通常,你会将这些服务部署到类似 OpenShift 或 Kubernetes 这样的平台上,由平台管理运行在多台主机和机器上的多个 Docker 容器;或者,你会在专用硬件上进行安装。 |
|---|
使用 Docker 运行 Debezium 时,你应该注意以下几点:
- Kafka 的容器是临时性的(ephemeral)。
Kafka 通常会将数据存储在容器内部的本地目录中,因此你需要在宿主机上挂载目录作为卷。这样,当容器停止时,持久化的数据依然保留。不过,本教程跳过了这一步设置——当容器停止时,所有持久化数据都会丢失。这样一来,完成教程后的清理工作就非常简单。
| 关于存储持久化数据的更多信息,请参阅容器镜像的文档。 |
|---|
本教程要求你在不同的容器中运行每个服务。
为避免混淆,你将在一个单独的终端中以前台方式运行每个容器。这样,容器的所有输出都会显示在运行它的那个终端中。
Docker 还允许你以 detached(分离)模式运行容器(使用 -d选项),即容器启动后docker命令会立即返回。但是,分离模式下的容器不会在终端中显示其输出。要查看输出,你需要使用docker logs --follow --name <container-name>命令。更多信息请参阅 Docker 文档。
启动 Kafka
| Debezium 3.6.3.Final 已针对多个版本的 Kafka Connect 进行了测试。请参阅 Debezium 测试矩阵以确定 Debezium 与 Kafka Connect 之间的兼容性。 |
|---|
操作步骤
打开一个新的终端,并用它在容器中启动 Kafka。
此命令使用
quay.io/debezium/kafka镜像的 3.6 版本运行一个新容器:$ docker run -it --rm -p 9092:9092 \ --name kafka --hostname kafka \ -e CLUSTER_ID=<YOUR_UNIQUE_CLUSTER_IDENTIFIER> \ -e NODE_ID=1 \ -e NODE_ROLE=combined \ -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093 \ -e KAFKA_LISTENERS=PLAINTEXT://kafka:9092,CONTROLLER://kafka:9093 \ -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092 \ quay.io/debezium/kafka:3.6
以下列表说明了前述命令中各参数的用途:
-it
该容器为交互式容器,即终端的标准输入和输出都连接到该容器。
--rm
容器停止后将被删除。
--name kafka
容器的名称。
--hostname kafka
该标签确保正确的 Kafka 监听器与容器配对。
-p 9092:9092
将容器中的 9092 端口映射到 Docker 主机上的同一端口,以便容器外部的应用程序能够与 Kafka 通信。
-e CLUSTER_ID=<YOUR_UNIQUE_CLUSTER_IDENTIFIER>
集群中的唯一标识符。集群内的所有节点必须使用相同的集群 ID。
-e NODE_ID=1
Raft 协议中使用的节点标识符。在集群内必须唯一。
-e NODE_ROLE=combined
节点在集群中的角色。可以是控制器(controller)、代理(broker)或两者兼具(combined)。
-e KAFKA_CONTROLLER_QUORUM_VOTERS=1@kafka:9093
指定在 Kafka 控制器仲裁中充当投票者的节点。如果需要指定多个值,请使用以逗号分隔的列表,格式为 nodeId@host:port,例如:1@kafka1:9093,2@kafka2:9093,3@kafka3:9093。
-e KAFKA_LISTENERS=PLAINTEXT://kafka:9092,CONTROLLER://kafka:9093
定义 Kafka 监听器的内部端点和协议。
-e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://kafka:9092
指定代理向外部客户端公布自身的方式。
如果你使用 Podman,请运行以下命令:
$ podman run -it --rm --name kafka --pod dbz -e HOST_NAME=127.0.0.1 quay.io/debezium/kafka:3.6在本教程中,你将始终在 Docker 容器内部连接到 Kafka。这些容器中的任何一个都可以通过链接到 kafka 容器来与之通信。如果你需要从 Docker 容器外部连接到 Kafka,则必须设置 -e 选项,通过 Docker 主机公布 Kafka 地址(-e ADVERTISED_HOST_NAME=,后跟 Docker 主机的 IP 地址或可解析的主机名)。 |
|
|---|---|
| 2. 验证 Kafka 已启动。 |
你应该看到类似如下的输出:
...
2025-07-22T11:43:00,935 - INFO [main:AppInfoParser$AppInfo@125] - Kafka version: 4.0.0
2025-07-22T11:43:00,938 - INFO [main:AppInfoParser$AppInfo@126] - Kafka commitId: 985bc99521dd22bb
2025-07-22T11:43:00,939 - INFO [main:AppInfoParser$AppInfo@127] - Kafka startTimeMs: 1753184580923
2025-07-22T11:43:00,959 - INFO [main:Logging@66] - [KafkaRaftServer nodeId=1] Kafka Server started上述响应表明 Kafka 节点已成功启动并准备好接受客户端连接。随着 Kafka 持续产生输出,终端将不断显示新的内容。
启动 MySQL 数据库
启动一个预先配置了示例 inventory 数据库的 MySQL 服务器,以便 Debezium 有一个变更捕获的源。
前提条件
- Kafka 正在运行。
操作步骤
打开一个新终端,并使用该终端启动一个新容器,该容器运行的 MySQL 数据库服务器已预先配置了
inventory数据库。此命令使用
quay.io/debezium/example-mysql镜像的 3.6 版本启动一个新容器,该镜像基于 mysql:8.2 镜像(来源)。它还定义并填充了一个示例inventory数据库:$ docker run -it --rm --name mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=debezium -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw quay.io/debezium/example-mysql:3.6
-it
容器是交互式的,这意味着终端的标准输入和输出已连接到容器。
--rm
容器停止后将被自动移除。
--name mysql
容器的名称。
-p 3306:3306
将容器内的 3306 端口(MySQL 默认端口)映射到 Docker 主机上的同一端口,使容器外部的应用程序可以连接到数据库服务器。
-e MYSQL_ROOT_PASSWORD=debezium -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw
创建一个用户及其密码,该用户仅拥有 Debezium MySQL 连接器所需的最小权限。
如果你使用 Podman,请运行以下命令:
$ podman run -it --rm --name mysql --pod dbz -e MYSQL_ROOT_PASSWORD=debezium -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw quay.io/debezium/example-mysql:3.6验证 MySQL 服务器是否启动。
MySQL 服务器会在配置修改过程中启动并停止几次。你应该看到类似如下的输出:
... [System] [MY-010931] [Server] /usr/sbin/mysqld: ready for connections. Version: '8.0.27' socket: '/var/run/mysqld/mysqld.sock' port: 3306 MySQL Community Server - GPL. [System] [MY-011323] [Server] X Plugin ready for connections. Bind-address: '::' port: 33060, socket: /var/run/mysqld/mysqlx.sock
启动 MySQL 命令行客户端
启动 MySQL 后,启动一个 MySQL 命令行客户端,以便访问示例 inventory 数据库。
操作步骤
打开一个新终端,并使用它在容器中启动 MySQL 命令行客户端。
此命令使用 mysql:8.2 镜像运行一个新容器,并定义一个 shell 命令,以正确的选项运行 MySQL 命令行客户端:
$ docker run -it --rm --name mysqlterm --link mysql mysql:8.2 sh -c 'exec mysql -h"$MYSQL_PORT_3306_TCP_ADDR" -P"$MYSQL_PORT_3306_TCP_PORT" -umysqluser -p"mysqlpw"'
-it该容器是交互式的,即终端的标准输入和输出都连接到该容器。
--rm容器停止时将被删除。
--name mysqlterm容器的名称。
--link mysql将该容器链接到
mysql容器。
如果你使用 Podman,请运行以下命令:
$ podman run -it --rm --name mysqlterm --pod dbz mysql:8.2 sh -c 'exec mysql -h 0.0.0.0 -uroot -pdebezium'确认 MySQL 命令行客户端已经启动。
你应该能看到类似如下的输出:
mysql: [Warning] Using a password on the command line interface can be insecure. Welcome to the MySQL monitor. Commands end with ; or \g. Your MySQL connection id is 9 Server version: 8.0.27 MySQL Community Server - GPL Copyright (c) 2000, 2021, Oracle and/or its affiliates. Oracle is a registered trademark of Oracle Corporation and/or its affiliates. Other names may be trademarks of their respective owners. Type 'help;' or '\h' for help. Type '\c' to clear the current input statement. mysql>在
mysql>命令提示符下,切换到 inventory 数据库:mysql> use inventory;列出该数据库中的表:
mysql> show tables; +---------------------+ | Tables_in_inventory | +---------------------+ | addresses | | customers | | geom | | orders | | products | | products_on_hand | +---------------------+ 6 rows in set (0.00 sec)使用 MySQL 命令行客户端浏览数据库,查看数据库中预先加载的数据。
例如:
mysql> SELECT * FROM customers; +------+------------+-----------+-----------------------+ | id | first_name | last_name | email | +------+------------+-----------+-----------------------+ | 1001 | Sally | Thomas | [email protected] | | 1002 | George | Bailey | [email protected] | | 1003 | Edward | Walker | [email protected] | | 1004 | Anne | Kretchmar | [email protected] | +------+------------+-----------+-----------------------+ 4 rows in set (0.00 sec)
启动 Kafka Connect
启动 Kafka Connect 服务,该服务会暴露一个 REST API,你可通过它管理 Debezium MySQL 连接器并监控其状态。
前提条件
- Kafka 和 MySQL 正在运行。
- 你已使用 MySQL 命令行客户端连接到
inventory数据库。
操作步骤
打开一个新的终端,用它在容器中启动 Kafka Connect 服务。
以下命令会使用 3.6 版本的
quay.io/debezium/connect镜像运行一个新容器:$ docker run -it --rm --name connect -p 8083:8083 -e GROUP_ID=1 -e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses --link kafka:kafka --link mysql:mysql quay.io/debezium/connect:3.6
-it
该容器是交互式的,即终端的标准输入和输出已连接到该容器。
--rm
容器停止后将被删除。
--name connect
容器的名称。
-p 8083:8083
将容器内的 8083 端口映射到 Docker 主机上的相同端口。这样,容器外部的应用程序就可以使用 Kafka Connect 的 REST API 来设置和管理新的容器实例。
-e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses
设置 Debezium 镜像所需的环境变量。
--link kafka:kafka --link mysql:mysql
将该容器链接到运行 Kafka 和 MySQL 服务器的容器。
如果你使用 Podman,请运行以下命令:
$ podman run -it --rm --name connect --pod dbz -e GROUP_ID=1 -e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses quay.io/debezium/connect:3.6如果你提供了 --hostname 命令选项,那么 Kafka Connect REST API 将不会监听 localhost 接口。这可能会在暴露 REST 端口时造成问题。
如果遇到此问题,请设置环境变量 REST_HOST_NAME=0.0.0.0,以确保 REST API 可以从所有接口访问。
验证 Kafka Connect 已启动并准备好接受连接。
你应该看到类似如下的输出:
... 2020-02-06 15:48:33,939 INFO || Kafka version: 3.0.0 [org.apache.kafka.common.utils.AppInfoParser] ... 2020-02-06 15:48:34,485 INFO || [Worker clientId=connect-1, groupId=1] Starting connectors and tasks using config offset -1 [org.apache.kafka.connect.runtime.distributed.DistributedHerder] 2020-02-06 15:48:34,485 INFO || [Worker clientId=connect-1, groupId=1] Finished starting connectors and tasks [org.apache.kafka.connect.runtime.distributed.DistributedHerder]使用 Kafka Connect REST API 检查 Kafka Connect 服务的状态。
Kafka Connect 提供了一个 REST API 来管理 Debezium 连接器。要与 Kafka Connect 服务通信,你可以使用
curl命令向 Docker 主机的 8083 端口发送 API 请求(在启动 Kafka Connect 时,你已在connect容器中将该端口映射为 8083)。这些命令使用 localhost。如果你使用的是非原生 Docker 平台(如 Docker Toolbox),请将localhost替换为 Docker 主机的 IP 地址。打开一个新终端并检查 Kafka Connect 服务的状态:
$ curl -H "Accept:application/json" localhost:8083/ {"version":"4.3.0","commit":"cb8625948210849f"}
前面的响应表明正在运行 Kafka Connect 4.3.0 版本。
2. 检查已注册到 Kafka Connect 的连接器列表:
```none
$ curl -H "Accept:application/json" localhost:8083/connectors/
[]
```空方括号([])表示没有连接器注册到 Kafka Connect。
部署 MySQL 连接器
启动 Debezium 和 MySQL 服务之后,你就可以部署 Debezium MySQL 连接器,让它开始监控示例 MySQL 数据库(inventory)。
此时,你已经运行了 Debezium 服务、一个带有示例 inventory 数据库的 MySQL 数据库服务器,以及连接到该数据库的 MySQL 命令行客户端。要部署 MySQL 连接器,你必须:
-
连接器注册后,它将开始监控数据库服务器的
binlog,并为发生变更的每一行生成变更事件。 -
在连接器启动时查看 Kafka Connect 的日志输出,有助于你更好地理解它在开始监控
binlog之前必须完成的每项任务。
注册连接器以监控 inventory 数据库
注册 Debezium MySQL 连接器后,该连接器将开始监控 MySQL 数据库服务器的 binlog。binlog 记录数据库的所有事务(例如单行数据的变更和模式的变更)。当数据库中的某一行发生变化时,Debezium 会生成一个变更事件。
| 在生产环境中,你通常要么使用 Kafka 工具手动创建必要的主题(包括指定副本数量),要么使用 Kafka Connect 机制来自定义自动创建的主题的设置。但是,对于本教程,Kafka 被配置为自动创建仅有一个副本的主题。 |
|---|
前提条件
- 你已熟悉连接器可用的配置选项。
操作步骤
查看你将要注册的 Debezium MySQL 连接器的配置。
在下一步中,你将注册由以下配置定义的连接器:
{ "name": "inventory-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "184054", "topic.prefix": "dbserver1", "database.include.list": "inventory", "schema.history.internal.kafka.bootstrap.servers": "kafka:9092", "schema.history.internal.kafka.topic": "schema-changes.inventory" } }
下面的列表描述了上一个连接器配置示例中的部分字段:
name
指定连接器的名称。
config
连接器的配置。
tasks.max
指定连接器一次可同时运行的任务数量。任何时候只应运行一个任务。因为 MySQL 连接器读取 MySQL 服务器的 binlog,使用单个连接器任务可确保事件顺序和处理正确。Kafka Connect 服务通过连接器启动一个或多个任务来执行工作,并会自动将运行中的任务分布到 Kafka Connect 服务集群中。如果有任何服务停止或崩溃,任务会被重新分配给正在运行的服务。
database.hostname
指定数据库主机,即运行 MySQL 服务器的容器名称(mysql)。Docker 操作容器内的网络栈,使得每个链接的容器都可以在 /etc/hosts 中通过容器名称作为主机名进行解析。如果 MySQL 运行在常规网络上,则此值应指定 IP 地址或可解析的主机名。
database.server.id
唯一的主题前缀。
topic.prefix
指定数据库服务器或集群的唯一主题前缀。该名称将作为所有接收来自该连接器的变更事件记录的 Kafka 主题的前缀。
database.include.list
指定连接器捕获其变更的数据库名称。此连接器仅捕获 inventory 数据库的变更。
schema.history.internal.kafka.bootstrap.servers
指定连接器使用 kafka:9092 作为其 Kafka 代理,以连接数据库架构历史主题来写入或读取 DDL 语句。连接器重启后,它会恢复连接器停止时点 binlog 中存在的数据库架构。
schema.history.internal.kafka.topic
指定连接器存储其数据库架构历史的主题名称。此主题仅供内部使用,不供消费者直接使用。
有关更多信息,请参阅 MySQL 连接器配置属性。
| 出于安全考虑,不应将密码或其他机密信息以明文形式放入连接器配置中。相反,应通过 KIP-297("Externalizing Secrets for Connect Configurations")中定义的机制将所有机密信息外部化。 |
|---|
打开一个新终端,使用
curl命令注册 Debezium MySQL 连接器。此命令使用 Kafka Connect 服务的 API,向
/connectors资源提交一个POST请求,请求中包含一个 JSON 文档,用于描述这个新连接器(名为inventory-connector)。此命令使用
localhost来连接 Docker 主机。如果你使用的是非原生 Docker 平台,请将localhost替换为 Docker 主机的 IP 地址。$ curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{ "name": "inventory-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "mysql", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "184054", "topic.prefix": "dbserver1", "database.include.list": "inventory", "schema.history.internal.kafka.bootstrap.servers": "kafka:9092", "schema.history.internal.kafka.topic": "schemahistory.inventory" } }'
Windows 用户可能需要对双引号进行转义。例如:
$ curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{ \"name\": \"inventory-connector\", \"config\": { \"connector.class\": \"io.debezium.connector.mysql.MySqlConnector\", \"tasks.max\": \"1\", \"database.hostname\": \"mysql\", \"database.port\": \"3306\", \"database.user\": \"debezium\", \"database.password\": \"dbz\", \"database.server.id\": \"184054\", \"topic.prefix\": \"dbserver1\", \"database.include.list\": \"inventory\", \"schema.history.internal.kafka.bootstrap.servers\": \"kafka:9092\", \"schema.history.internal.kafka.topic\": \"schemahistory.inventory\" } }'否则,你可能会看到类似如下的错误:
{"error_code":500,"message":"Unexpected character ('n' (code 110)): was expecting double-quote to start field name\n at [Source: (org.glassfish.jersey.message.internal.ReaderInterceptorExecutor$UnCloseableInputStream); line: 1, column: 4]"}如果你使用 Podman,请运行以下命令:
$ curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{ "name": "inventory-connector", "config": { "connector.class": "io.debezium.connector.mysql.MySqlConnector", "tasks.max": "1", "database.hostname": "0.0.0.0", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "184054", "topic.prefix": "dbserver1", "database.include.list": "inventory", "schema.history.internal.kafka.bootstrap.servers": "0.0.0.0:9092", "schema.history.internal.kafka.topic": "schemahistory.inventory" } }'确认
inventory-connector已包含在连接器列表中:$ curl -H "Accept:application/json" localhost:8083/connectors/ ["inventory-connector"]查看该连接器的任务:
$ curl -i -X GET -H "Accept:application/json" localhost:8083/connectors/inventory-connector
你应该看到类似下面的响应(为了便于阅读,已做了格式化):
HTTP/1.1 200 OK
Date: Thu, 06 Feb 2020 22:12:03 GMT
Content-Type: application/json
Content-Length: 531
Server: Jetty(9.4.20.v20190813)
{
"name": "inventory-connector",
...
"tasks": [
{
"connector": "inventory-connector",
"task": 0
}
]
}连接器正在运行单个任务(任务 0)来完成其工作。该连接器仅支持单个任务,因为 MySQL 会将其所有活动按顺序记录在一个 binlog 中。因此,连接器只需要一个读取器,即可获得所有事件的一致、有序视图。
观察连接器启动过程
注册连接器时,Kafka Connect 容器中会生成大量日志输出。通过查看这些输出,你可以更好地了解连接器从创建到开始读取 MySQL 服务器 binlog 所经历的过程。
注册 inventory-connector 连接器后,你可以查看 Kafka Connect 容器(connect)中的日志输出,以跟踪连接器的状态。
前几行显示了连接器(inventory-connector)的创建和启动:
...
2021-11-30 01:38:44,223 INFO || [Worker clientId=connect-1, groupId=1] Tasks [inventory-connector-0] configs updated [org.apache.kafka.connect.runtime.distributed.DistributedHerder]
2021-11-30 01:38:44,224 INFO || [Worker clientId=connect-1, groupId=1] Handling task config update by restarting tasks [] [org.apache.kafka.connect.runtime.distributed.DistributedHerder]
2021-11-30 01:38:44,224 INFO || [Worker clientId=connect-1, groupId=1] Rebalance started [org.apache.kafka.connect.runtime.distributed.WorkerCoordinator]
2021-11-30 01:38:44,224 INFO || [Worker clientId=connect-1, groupId=1] (Re-)joining group [org.apache.kafka.connect.runtime.distributed.WorkerCoordinator]
2021-11-30 01:38:44,227 INFO || [Worker clientId=connect-1, groupId=1] Successfully joined group with generation Generation{generationId=3, memberId='connect-1-7b087c69-8ac5-4c56-9e6b-ec5adabf27e8', protocol='sessioned'} [org.apache.kafka.connect.runtime.distributed.WorkerCoordinator]
2021-11-30 01:38:44,230 INFO || [Worker clientId=connect-1, groupId=1] Successfully synced group in generation Generation{generationId=3, memberId='connect-1-7b087c69-8ac5-4c56-9e6b-ec5adabf27e8', protocol='sessioned'} [org.apache.kafka.connect.runtime.distributed.WorkerCoordinator]
2021-11-30 01:38:44,231 INFO || [Worker clientId=connect-1, groupId=1] Joined group at generation 3 with protocol version 2 and got assignment: Assignment{error=0, leader='connect-1-7b087c69-8ac5-4c56-9e6b-ec5adabf27e8', leaderUrl='http://172.17.0.7:8083/', offset=4, connectorIds=[inventory-connector], taskIds=[inventory-connector-0], revokedConnectorIds=[], revokedTaskIds=[], delay=0} with rebalance delay: 0 [org.apache.kafka.connect.runtime.distributed.DistributedHerder]
2021-11-30 01:38:44,232 INFO || [Worker clientId=connect-1, groupId=1] Starting connectors and tasks using config offset 4 [org.apache.kafka.connect.runtime.distributed.DistributedHerder]
2021-11-30 01:38:44,232 INFO || [Worker clientId=connect-1, groupId=1] Starting task inventory-connector-0 [org.apache.kafka.connect.runtime.distributed.DistributedHerder]
...再往下,你应该能看到连接器输出类似如下的内容:
...
2021-11-30 01:38:44,406 INFO || Kafka version: 3.0.0 [org.apache.kafka.common.utils.AppInfoParser]
2021-11-30 01:38:44,406 INFO || Kafka commitId: 8cb0a5e9d3441962 [org.apache.kafka.common.utils.AppInfoParser]
2021-11-30 01:38:44,407 INFO || Kafka startTimeMs: 1638236324406 [org.apache.kafka.common.utils.AppInfoParser]
2021-11-30 01:38:44,437 INFO || Database schema history topic '(name=schemahistory.inventory, numPartitions=1, replicationFactor=1, replicasAssignments=null, configs={cleanup.policy=delete, retention.ms=9223372036854775807, retention.bytes=-1})' created [io.debezium.storage.kafka.history.KafkaSchemaHistory]
2021-11-30 01:38:44,497 INFO || App info kafka.admin.client for dbserver1-schemahistory unregistered [org.apache.kafka.common.utils.AppInfoParser]
2021-11-30 01:38:44,499 INFO || Metrics scheduler closed [org.apache.kafka.common.metrics.Metrics]
2021-11-30 01:38:44,499 INFO || Closing reporter org.apache.kafka.common.metrics.JmxReporter [org.apache.kafka.common.metrics.Metrics]
2021-11-30 01:38:44,499 INFO || Metrics reporters closed [org.apache.kafka.common.metrics.Metrics]
2021-11-30 01:38:44,499 INFO || Reconnecting after finishing schema recovery [io.debezium.connector.mysql.MySqlConnectorTask]
2021-11-30 01:38:44,524 INFO || Requested thread factory for connector MySqlConnector, id = dbserver1 named = change-event-source-coordinator [io.debezium.util.Threads]
2021-11-30 01:38:44,525 INFO || Creating thread debezium-mysqlconnector-dbserver1-change-event-source-coordinator [io.debezium.util.Threads]
2021-11-30 01:38:44,526 INFO || WorkerSourceTask{id=inventory-connector-0} Source task finished initialization and start [org.apache.kafka.connect.runtime.WorkerSourceTask]
2021-11-30 01:38:44,529 INFO MySQL|dbserver1|snapshot Metrics registered [io.debezium.pipeline.ChangeEventSourceCoordinator]
2021-11-30 01:38:44,529 INFO MySQL|dbserver1|snapshot Context created [io.debezium.pipeline.ChangeEventSourceCoordinator]
2021-11-30 01:38:44,534 INFO MySQL|dbserver1|snapshot No previous offset has been found [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,534 INFO MySQL|dbserver1|snapshot According to the connector configuration both schema and data will be snapshotted [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,534 INFO MySQL|dbserver1|snapshot Snapshot step 1 - Preparing [io.debezium.relational.RelationalSnapshotChangeEventSource]
...Debezium 的日志输出使用 映射诊断上下文(MDC)在日志中提供与线程相关的信息,从而更容易理解在多线程的 Kafka Connect 服务中正在发生什么。这些信息包括连接器类型(上述日志中的 MySQL)、连接器的逻辑名称(上述日志中的 dbserver1),以及连接器的活动(task、snapshot 和 binlog)。
在上面的日志输出中,前几行涉及连接器的 task 活动,并报告了一些簿记信息(在本例中,报告连接器是在没有先前偏移量的情况下启动的)。接下来的三行涉及连接器的 snapshot 活动,报告正使用 debezium MySQL 用户启动快照,以及与该用户关联的 MySQL 权限。
如果连接器无法连接,或者看不到任何表或 binlog,请检查这些权限,确保上面列出的权限都已包含在内。 |
|---|
接下来,连接器报告构成快照操作的各个步骤:
...
2021-11-30 01:38:44,534 INFO MySQL|dbserver1|snapshot Snapshot step 1 - Preparing [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:44,535 INFO MySQL|dbserver1|snapshot Snapshot step 2 - Determining captured tables [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:44,535 INFO MySQL|dbserver1|snapshot Read list of available databases [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,537 INFO MySQL|dbserver1|snapshot list of available databases is: [information_schema, inventory, mysql, performance_schema, sys] [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,537 INFO MySQL|dbserver1|snapshot Read list of available tables in each database [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,548 INFO MySQL|dbserver1|snapshot snapshot continuing with database(s): [inventory] [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,551 INFO MySQL|dbserver1|snapshot Snapshot step 3 - Locking captured tables [inventory.addresses, inventory.customers, inventory.geom, inventory.orders, inventory.products, inventory.products_on_hand] [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:44,552 INFO MySQL|dbserver1|snapshot Flush and obtain global read lock to prevent writes to database [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,557 INFO MySQL|dbserver1|snapshot Snapshot step 4 - Determining snapshot offset [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:44,560 INFO MySQL|dbserver1|snapshot Read binlog position of MySQL primary server [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,562 INFO MySQL|dbserver1|snapshot using binlog 'mysql-bin.000003' at position '156' and gtid '' [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:44,562 INFO MySQL|dbserver1|snapshot Snapshot step 5 - Reading structure of captured tables [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:44,562 INFO MySQL|dbserver1|snapshot All eligible tables schema should be captured, capturing: [inventory.addresses, inventory.customers, inventory.geom, inventory.orders, inventory.products, inventory.products_on_hand] [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:45,058 INFO MySQL|dbserver1|snapshot Reading structure of database 'inventory' [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:45,187 INFO MySQL|dbserver1|snapshot Snapshot step 6 - Persisting schema history [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,273 INFO MySQL|dbserver1|snapshot Releasing global read lock to enable MySQL writes [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:45,274 INFO MySQL|dbserver1|snapshot Writes to MySQL tables prevented for a total of 00:00:00.717 [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:38:45,274 INFO MySQL|dbserver1|snapshot Snapshot step 7 - Snapshotting data [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,275 INFO MySQL|dbserver1|snapshot Snapshotting contents of 6 tables while still in transaction [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,275 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.addresses' (1 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,276 INFO MySQL|dbserver1|snapshot For table 'inventory.addresses' using select statement: 'SELECT `id`, `customer_id`, `street`, `city`, `state`, `zip`, `type` FROM `inventory`.`addresses`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,295 INFO MySQL|dbserver1|snapshot Finished exporting 7 records for table 'inventory.addresses'; total duration '00:00:00.02' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,296 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.customers' (2 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,296 INFO MySQL|dbserver1|snapshot For table 'inventory.customers' using select statement: 'SELECT `id`, `first_name`, `last_name`, `email` FROM `inventory`.`customers`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,304 INFO MySQL|dbserver1|snapshot Finished exporting 4 records for table 'inventory.customers'; total duration '00:00:00.008' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,304 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.geom' (3 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,305 INFO MySQL|dbserver1|snapshot For table 'inventory.geom' using select statement: 'SELECT `id`, `g`, `h` FROM `inventory`.`geom`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,316 INFO MySQL|dbserver1|snapshot Finished exporting 3 records for table 'inventory.geom'; total duration '00:00:00.011' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,316 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.orders' (4 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,316 INFO MySQL|dbserver1|snapshot For table 'inventory.orders' using select statement: 'SELECT `order_number`, `order_date`, `purchaser`, `quantity`, `product_id` FROM `inventory`.`orders`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,325 INFO MySQL|dbserver1|snapshot Finished exporting 4 records for table 'inventory.orders'; total duration '00:00:00.008' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,325 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.products' (5 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,325 INFO MySQL|dbserver1|snapshot For table 'inventory.products' using select statement: 'SELECT `id`, `name`, `description`, `weight` FROM `inventory`.`products`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,343 INFO MySQL|dbserver1|snapshot Finished exporting 9 records for table 'inventory.products'; total duration '00:00:00.017' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,344 INFO MySQL|dbserver1|snapshot Exporting data from table 'inventory.products_on_hand' (6 of 6 tables) [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,344 INFO MySQL|dbserver1|snapshot For table 'inventory.products_on_hand' using select statement: 'SELECT `product_id`, `quantity` FROM `inventory`.`products_on_hand`' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,353 INFO MySQL|dbserver1|snapshot Finished exporting 9 records for table 'inventory.products_on_hand'; total duration '00:00:00.009' [io.debezium.relational.RelationalSnapshotChangeEventSource]
2021-11-30 01:38:45,355 INFO MySQL|dbserver1|snapshot Snapshot - Final stage [io.debezium.pipeline.source.AbstractSnapshotChangeEventSource]
2021-11-30 01:38:45,356 INFO MySQL|dbserver1|snapshot Snapshot ended with SnapshotResult [status=COMPLETED, offset=MySqlOffsetContext [sourceInfoSchema=Schema{io.debezium.connector.mysql.Source:STRUCT}, sourceInfo=SourceInfo [currentGtid=null, currentBinlogFilename=mysql-bin.000003, currentBinlogPosition=156, currentRowNumber=0, serverId=0, sourceTime=2021-11-30T01:38:45.352Z, threadId=-1, currentQuery=null, tableIds=[inventory.products_on_hand], databaseName=inventory], snapshotCompleted=true, transactionContext=TransactionContext [currentTransactionId=null, perTableEventCount={}, totalEventCount=0], restartGtidSet=null, currentGtidSet=null, restartBinlogFilename=mysql-bin.000003, restartBinlogPosition=156, restartRowsToSkip=0, restartEventsToSkip=0, currentEventLengthInBytes=0, inTransaction=false, transactionId=null, incrementalSnapshotContext =IncrementalSnapshotContext [windowOpened=false, chunkEndPosition=null, dataCollectionsToSnapshot=[], lastEventKeySent=null, maximumKey=null]]] [io.debezium.pipeline.ChangeEventSourceCoordinator]
...其中每一步都会报告连接器为执行一致性快照所做的工作。例如,第 6 步会对被捕获的表逆向生成 DDL create 语句,并在获取全局写锁仅仅 1 秒后就释放它;第 7 步会读取每张表中的所有行,并报告所耗时间以及找到的行数。在这种情况下,连接器在不到 1 秒的时间内就完成了它的一致性快照。
| 在你的数据库上,快照过程会花费更长时间,但连接器会输出足够多的日志消息,让你能够跟踪它的处理进度,即使表中的数据量很大也是如此。虽然快照过程开始时会使用排他写锁,但即使对于大型数据库,该锁也不会持续很长时间。这是因为在复制任何数据之前锁就已经被释放了。更多信息请参阅 MySQL 连接器文档。 |
|---|
接下来,Kafka Connect 报告了一些“错误”。不过,你可以放心忽略这些警告:这些消息只是意味着创建了新的 Kafka 主题,并且 Kafka 必须为每个主题分配一个新的领导者:
...
2021-11-30 01:38:45,555 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 3 : {dbserver1=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:45,691 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 9 : {dbserver1.inventory.addresses=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:45,813 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 13 : {dbserver1.inventory.customers=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:45,927 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 18 : {dbserver1.inventory.geom=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:46,043 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 22 : {dbserver1.inventory.orders=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:46,153 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 26 : {dbserver1.inventory.products=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
2021-11-30 01:38:46,269 WARN || [Producer clientId=connector-producer-inventory-connector-0] Error while fetching metadata with correlation id 31 : {dbserver1.inventory.products_on_hand=LEADER_NOT_AVAILABLE} [org.apache.kafka.clients.NetworkClient]
...最后,日志输出显示连接器已从快照模式切换为持续读取 MySQL 服务器的 binlog:
...
2021-11-30 01:38:45,362 INFO MySQL|dbserver1|streaming Starting streaming [io.debezium.pipeline.ChangeEventSourceCoordinator]
...
Nov 30, 2021 1:38:45 AM com.github.shyiko.mysql.binlog.BinaryLogClient connect
INFO: Connected to mysql:3306 at mysql-bin.000003/156 (sid:184054, cid:13)
2021-11-30 01:38:45,392 INFO MySQL|dbserver1|binlog Connected to MySQL binlog at mysql:3306, starting at MySqlOffsetContext [sourceInfoSchema=Schema{io.debezium.connector.mysql.Source:STRUCT}, sourceInfo=SourceInfo [currentGtid=null, currentBinlogFilename=mysql-bin.000003, currentBinlogPosition=156, currentRowNumber=0, serverId=0, sourceTime=2021-11-30T01:38:45.352Z, threadId=-1, currentQuery=null, tableIds=[inventory.products_on_hand], databaseName=inventory], snapshotCompleted=true, transactionContext=TransactionContext [currentTransactionId=null, perTableEventCount={}, totalEventCount=0], restartGtidSet=null, currentGtidSet=null, restartBinlogFilename=mysql-bin.000003, restartBinlogPosition=156, restartRowsToSkip=0, restartEventsToSkip=0, currentEventLengthInBytes=0, inTransaction=false, transactionId=null, incrementalSnapshotContext =IncrementalSnapshotContext [windowOpened=false, chunkEndPosition=null, dataCollectionsToSnapshot=[], lastEventKeySent=null, maximumKey=null]] [io.debezium.connector.mysql.MySqlStreamingChangeEventSource]
2021-11-30 01:38:45,392 INFO MySQL|dbserver1|streaming Waiting for keepalive thread to start [io.debezium.connector.mysql.MySqlStreamingChangeEventSource]
2021-11-30 01:38:45,393 INFO MySQL|dbserver1|binlog Creating thread debezium-mysqlconnector-dbserver1-binlog-client [io.debezium.util.Threads]
...查看变更事件
部署 Debezium MySQL 连接器后,它会开始监视 inventory 数据库中的数据变更事件。
在观察连接器启动的过程中,你会看到事件被写入到以 dbserver1(连接器的名称)为前缀的以下主题中:
dbserver1
模式变更主题,所有 DDL 语句都会写入其中。
dbserver1.inventory.products
捕获 inventory 数据库中 products 表的变更事件。
dbserver1.inventory.products_on_hand
捕获 inventory 数据库中 products_on_hand 表的变更事件。
dbserver1.inventory.customers
捕获 inventory 数据库中 customers 表的变更事件。
dbserver1.inventory.orders
捕获 inventory 数据库中 orders 表的变更事件。
在本教程中,我们将探索 dbserver1.inventory.customers 主题。在这个主题中,你会看到不同类型的变更事件,以了解连接器是如何捕获它们的:
查看 create 事件
通过查看 dbserver1.inventory.customers 主题,你可以看到 MySQL 连接器是如何捕获 inventory 数据库中的 create 事件的。在本例中,create 事件捕获的是新添加到数据库中的客户。
操作步骤
打开一个新的终端,使用它启动
watch-topic工具,从该主题的开头开始监视dbserver1.inventory.customers主题。watch-topic工具非常简单,功能有限。它并非供应用程序用来消费事件。在这种场景下,你应改用 Kafka 消费者以及提供完整功能和灵活性的相应消费者类库。此命令使用 3.6 版本的
debezium/kafka镜像在新容器中运行watch-topic工具:$ docker run -it --rm --name watcher --link kafka:kafka quay.io/debezium/kafka:3.6 watch-topic -a -k dbserver1.inventory.customers
-a
监视自主题创建以来的所有事件。如果不使用此选项,watch-topic 只会显示你开始监视之后记录的事件。
-k
指定输出中应包含事件的键。在本例中,这包含该行的主键。
如果你使用 Podman,请运行以下命令:
$ podman run -it --rm --name watcher --pod dbz quay.io/debezium/kafka:3.6 watch-topic -a -k dbserver1.inventory.customerswatch-topic 工具返回了 customers 表的事件记录。共有四个事件,表中每行对应一个事件。每个事件均为 JSON 格式,因为 Kafka Connect 服务就是这样配置的。每个事件包含两个 JSON 文档:一个作为键(key),一个作为值(value)。
你应该看到类似如下的输出:
Using KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://172.17.0.7:9092
Using KAFKA_BROKER=172.17.0.3:9092
Contents of topic dbserver1.inventory.customers:
{"schema":{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"}],"optional":false,"name":"dbserver1.inventory.customers.Key"},"payload":{"id":1001}}
...| 该工具会持续监控主题,因此只要工具在运行,任何新事件都会自动出现。 |
|---|
- 对于最后一条事件,检查 key 的详细信息。
以下是最后一条事件 key 的详细信息(为便于阅读已格式化):
{
"schema":{
"type":"struct",
"fields":[
{
"type":"int32",
"optional":false,
"field":"id"
}
],
"optional":false,
"name":"dbserver1.inventory.customers.Key"
},
"payload":{
"id":1004
}
}事件包含两个部分:一个 schema 和一个 payload。schema 包含一个 Kafka Connect 模式,用于描述 payload 中的内容。在本例中,payload 是一个名为 dbserver1.inventory.customers.Key 的 struct,该结构不可为可选值,且包含一个必填字段(类型为 int32 的 id)。
payload 中只有一个 id 字段,其值为 1004。
通过查看事件的 key,可以看到该事件适用于 inventory.customers 表中 id 主键列值为 1004 的那一行记录。
- 查看同一事件的 value 详细信息。
事件的 value 显示该行记录被创建,并描述了其包含的内容(在本例中,为插入行的 id、first_name、last_name 和 email)。
以下是最后一个事件的 value 详细信息(为便于阅读已进行格式化):
{
"schema": {
"type": "struct",
"fields": [
{
"type": "struct",
"fields": [
{
"type": "int32",
"optional": false,
"field": "id"
},
{
"type": "string",
"optional": false,
"field": "first_name"
},
{
"type": "string",
"optional": false,
"field": "last_name"
},
{
"type": "string",
"optional": false,
"field": "email"
}
],
"optional": true,
"name": "dbserver1.inventory.customers.Value",
"field": "before"
},
{
"type": "struct",
"fields": [
{
"type": "int32",
"optional": false,
"field": "id"
},
{
"type": "string",
"optional": false,
"field": "first_name"
},
{
"type": "string",
"optional": false,
"field": "last_name"
},
{
"type": "string",
"optional": false,
"field": "email"
}
],
"optional": true,
"name": "dbserver1.inventory.customers.Value",
"field": "after"
},
{
"type": "struct",
"fields": [
{
"type": "string",
"optional": true,
"field": "version"
},
{
"type": "string",
"optional": false,
"field": "name"
},
{
"type": "int64",
"optional": false,
"field": "server_id"
},
{
"type": "int64",
"optional": false,
"field": "ts_sec"
},
{
"type": "string",
"optional": true,
"field": "gtid"
},
{
"type": "string",
"optional": false,
"field": "file"
},
{
"type": "int64",
"optional": false,
"field": "pos"
},
{
"type": "int32",
"optional": false,
"field": "row"
},
{
"type": "boolean",
"optional": true,
"field": "snapshot"
},
{
"type": "int64",
"optional": true,
"field": "thread"
},
{
"type": "string",
"optional": true,
"field": "db"
},
{
"type": "string",
"optional": true,
"field": "table"
}
],
"optional": false,
"name": "io.debezium.connector.mysql.Source",
"field": "source"
},
{
"type": "string",
"optional": false,
"field": "op"
},
{
"type": "int64",
"optional": true,
"field": "ts_ms"
},
{
"type": "int64",
"optional": true,
"field": "ts_us"
},
{
"type": "int64",
"optional": true,
"field": "ts_ns"
}
],
"optional": false,
"name": "dbserver1.inventory.customers.Envelope",
"version": 1
},
"payload": {
"before": null,
"after": {
"id": 1004,
"first_name": "Anne",
"last_name": "Kretchmar",
"email": "[email protected]"
},
"source": {
"version": "3.6.3.Final",
"name": "dbserver1",
"server_id": 0,
"ts_sec": 0,
"gtid": null,
"file": "mysql-bin.000003",
"pos": 154,
"row": 0,
"snapshot": true,
"thread": null,
"db": "inventory",
"table": "customers"
},
"op": "r",
"ts_ms": 1486500577691,
"ts_us": 1486500577691547,
"ts_ns": 1486500577691547930
}
}这一部分比事件的其余部分长得多,但和事件的 key 一样,它也包含 schema 和 payload。schema 中包含一个名为 dbserver1.inventory.customers.Envelope(版本 1)的 Kafka Connect schema,该 schema 可以包含五个字段:
op
一个必填字段,包含描述操作类型的字符串值。MySQL 连接器的取值为:c 表示创建(或插入)、u 表示更新、d 表示删除、r 表示读取(在快照的情况下)。
before
一个可选字段,如果存在,则包含事件发生之前行的状态。其结构由 dbserver1.inventory.customers.Value Kafka Connect schema 描述,dbserver1 连接器对 inventory.customers 表中的所有行都使用该 schema。
after
一个可选字段,如果存在,则包含事件发生之后行的状态。其结构由与 before 中相同的 dbserver1.inventory.customers.Value Kafka Connect schema 描述。
source
一个必填字段,包含一个描述事件源元数据的结构。对于 MySQL 而言,它包含若干字段:连接器名称、记录该事件的 binlog 文件名、该事件在 binlog 文件中的位置、事件中的行(如果包含多行)、受影响的数据库和表的名称、执行该更改的 MySQL 线程 ID、该事件是否属于快照的一部分,以及(如果可用)MySQL 服务器 ID 和秒级时间戳。
ts_ms
一个可选字段,如果存在,则包含连接器处理该事件的时间(使用运行 Kafka Connect 任务的 JVM 的系统时钟)。
事件的 JSON 表示形式比它们所描述的行要长得多。这是因为,对于每个事件的键和值,Kafka Connect 都会随附描述 payload 的 schema。随着时间推移,这种结构可能会发生变化。然而,由于键和值的 schema 都包含在事件本身中,消费应用程序就更容易理解这些消息,尤其是在消息随时间不断演进的情况下。
Debezium MySQL 连接器根据数据库表的结构来构造这些 schema。如果你使用 DDL 语句修改 MySQL 数据库中的表定义,连接器会读取这些 DDL 语句并更新其 Kafka Connect schema。这是确保每个事件的结构与事件发生时其源表结构完全一致的唯一方式。不过,包含单个表所有事件的 Kafka 主题中,可能会存在对应于表定义各个历史状态的事件。
JSON 转换器会在每条消息中都包含键和值的模式,因此生成的事件非常冗长。你也可以使用 Apache Avro 作为序列化格式,这样得到的事件消息要小得多。这是因为 Avro 会把每个 Kafka Connect 模式转换为 Avro 模式,并将这些 Avro 模式存储在一个独立的 Schema Registry 服务中。因此,当 Avro 转换器序列化事件消息时,它只会在消息中放入该模式的唯一标识符以及值的 Avro 编码二进制表示。这样一来,在网络上传输并存储到 Kafka 中的序列化消息,比你在前面看到的要小得多。实际上,Avro 转换器能够利用 Avro 模式演进技术,在 Schema Registry 中维护每个模式的历史记录。
将事件的 键 和 值 模式与
inventory数据库的当前状态进行比较。在运行 MySQL 命令行客户端的终端中,执行以下语句:mysql> SELECT * FROM customers; +------+------------+-----------+-----------------------+ | id | first_name | last_name | email | +------+------------+-----------+-----------------------+ | 1001 | Sally | Thomas | [email protected] | | 1002 | George | Bailey | [email protected] | | 1003 | Edward | Walker | [email protected] | | 1004 | Anne | Kretchmar | [email protected] | +------+------------+-----------+-----------------------+ 4 rows in set (0.00 sec)
这表明你查看的事件记录与数据库中的记录相匹配。
更新数据库并查看 update 事件
通过更新数据库中的某条记录,并观察 Debezium MySQL 连接器捕获的变更事件,你可以了解如何识别一次提交中的变更内容,以及如何比较事件以确定它们的先后顺序。
操作步骤
在运行 MySQL 命令行客户端的终端中,执行以下语句:
mysql> UPDATE customers SET first_name='Anne Marie' WHERE id=1004; Query OK, 1 row affected (0.05 sec) Rows matched: 1 Changed: 1 Warnings: 0查看更新后的
customers表:mysql> SELECT * FROM customers; +------+------------+-----------+-----------------------+ | id | first_name | last_name | email | +------+------------+-----------+-----------------------+ | 1001 | Sally | Thomas | [email protected] | | 1002 | George | Bailey | [email protected] | | 1003 | Edward | Walker | [email protected] | | 1004 | Anne Marie | Kretchmar | [email protected] | +------+------------+-----------+-----------------------+ 4 rows in set (0.00 sec)切换到运行
watch-topic的终端,你会看到第五个新事件。由于修改了
customers表中的一条记录,Debezium MySQL 连接器生成了一个新事件。你应该能看到两个新的 JSON 文档:一个用于事件的键,另一个用于新事件的值。以下是 update 事件键的详细信息(已为便于阅读进行了格式化):
{ "schema": { "type": "struct", "name": "dbserver1.inventory.customers.Key", "optional": false, "fields": [ { "field": "id", "type": "int32", "optional": false } ] }, "payload": { "id": 1004 } }
此键与前一个事件的键相同。
下面是该新事件的值。schema 部分没有变化,因此这里只展示 payload 部分(已格式化以便阅读):
{
"schema": {...},
"payload": {
"before": {
"id": 1004,
"first_name": "Anne",
"last_name": "Kretchmar",
"email": "[email protected]"
},
"after": {
"id": 1004,
"first_name": "Anne Marie",
"last_name": "Kretchmar",
"email": "[email protected]"
},
"source": {
"name": "3.6.3.Final",
"name": "dbserver1",
"server_id": 223344,
"ts_sec": 1486501486,
"gtid": null,
"file": "mysql-bin.000003",
"pos": 364,
"row": 0,
"snapshot": null,
"thread": 3,
"db": "inventory",
"table": "customers"
},
"op": "u",
"ts_ms": 1486501486308,
"ts_us": 1486501486308910,
"ts_ns": 1486501486308910814
}
}下列列表介绍了上述 update 事件值负载示例中的部分字段:
before
before 字段显示更新之前该行中存在的值。原始的 first_name 值为 Anne。
after
after 字段显示 update 事件发生后该行的状态。first_name 的值现在是 Anne Marie。
source
source 字段结构中的大多数值与之前相同,只是 ts_sec 和 pos 字段发生了变化。在某些情况下,file 的值也可能发生变化。
op
显示该事件所描述的操作类型。u 值表示该变更是由 update 事件引起的。
ts_ms、ts_us、ts_ns
以毫秒、微秒和纳秒为单位显示时间戳,指示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
通过查看 payload 部分,您可以了解到有关 update 事件的几项重要信息:
- 通过比较
before和after结构,您可以确定由于此次提交,受影响的行中究竟发生了哪些变化。 - 通过查看
source结构,您可以找到 MySQL 对该变更的记录信息(提供可追溯性)。 - 通过将某个事件的
payload部分与同一主题(或不同主题)中的其他事件进行比较,您可以确定该事件是发生在另一个事件之前、之后,还是与另一个事件属于同一次 MySQL 提交。
在数据库中删除记录并查看 delete 事件
通过在数据库中删除记录并观察 Debezium MySQL 连接器捕获的变更事件,您可以了解 Debezium 如何表示删除操作,以及 Kafka 日志压缩如何处理墓碑事件。
操作步骤
在运行 MySQL 命令行客户端的终端中,执行以下语句:
mysql> DELETE FROM customers WHERE id=1004; Query OK, 1 row affected (0.00 sec)
如果上述命令因外键约束冲突而失败,则必须使用以下语句删除 addresses 表中对该客户地址的引用:
mysql> DELETE FROM addresses WHERE customer_id=1004;切换到运行
watch-topic的终端,可以看到两条新事件。通过删除
customers表中的一行,Debezium MySQL 连接器生成了两条新事件。查看第一条新事件的键和值。
以下是第一条新事件的键的详细信息(已为便于阅读进行了格式化):
{ "schema": { "type": "struct", "name": "dbserver1.inventory.customers.Key", "optional": false, "fields": [ { "field": "id", "type": "int32", "optional": false } ] }, "payload": { "id": 1004 } }
这个键与你刚才查看的前两个事件中的键相同。
下面是第一个新事件的值(已格式化以便阅读):
{
"schema": {...},
"payload": {
"before": {
"id": 1004,
"first_name": "Anne Marie",
"last_name": "Kretchmar",
"email": "[email protected]"
},
"after": null,
"source": {
"name": "3.6.3.Final",
"name": "dbserver1",
"server_id": 223344,
"ts_sec": 1486501558,
"gtid": null,
"file": "mysql-bin.000003",
"pos": 725,
"row": 0,
"snapshot": null,
"thread": 3,
"db": "inventory",
"table": "customers"
},
"op": "d",
"ts_ms": 1486501558315,
"ts_us": 1486501558315901,
"ts_ns": 1486501558315901687
}
}以下列表描述了前述 delete 事件值中的部分字段:
before
before 字段显示行在被删除之前的状态。
after
显示删除后的状态。after 字段的值为 null,因为该行已不存在。
source
source 字段结构中的大部分值与删除发生之前的值相同,但 ts_sec 和 pos 字段的值已发生变化。在某些情况下,file 的值也可能发生改变。
op
显示该事件所描述的操作类型。op 字段的值为 d,表示该记录已被删除。
ts_ms、ts_us、ts_ns
显示以毫秒、微秒和纳秒为单位的时间戳,指示连接器处理该事件的时间。该时间基于运行 Kafka Connect 任务的 JVM 中的系统时钟。
因此,该事件为消费者提供了处理行删除所需的信息。事件记录还提供了旧值,因为某些消费者可能需要这些值才能正确处理删除操作。
- 检查第二个新事件的 key 和 value。
以下是第二个新事件的 key(为便于阅读已格式化):
{
"schema": {
"type": "struct",
"name": "dbserver1.inventory.customers.Key"
"optional": false,
"fields": [
{
"field": "id",
"type": "int32",
"optional": false
}
]
},
"payload": {
"id": 1004
}
}再次说明,这个 key 与你之前查看的三个事件中的 key 完全相同。
下面是该事件的 value(为了便于阅读已格式化):
{
"schema": null,
"payload": null
}如果 Kafka 设置为 日志压缩(log compacted),当主题中存在相同键的较晚消息时,它会移除较早的消息。最后这条事件被称为 墓碑(tombstone)事件,因为它包含一个键和一个空值。这意味着 Kafka 将移除所有相同键的先前消息。尽管先前的消息会被移除,但墓碑事件意味着消费者仍可以从头读取主题,而不会遗漏任何事件。
重启 Kafka Connect 服务
当 Kafka Connect 服务下线并重新启动时,Debezium 会从中断的位置继续捕获数据库变更,不会遗漏任何事件。
Kafka Connect 服务会自动管理其已注册连接器的任务。因此,如果它下线,重启后将启动所有未运行的任务。这意味着即使 Debezium 未运行,它仍然能够报告数据库中的变更。
操作步骤
打开一个新终端,使用该终端停止运行 Kafka Connect 服务的
connect容器:$ docker stop connect
connect 容器已停止,Kafka Connect 服务正常关闭。
由于你使用 --rm 选项运行容器,Docker 会在容器停止后将其删除。
2. 在服务停止期间,切换到 MySQL 命令行客户端所在的终端,并添加几条记录:
mysql> INSERT INTO customers VALUES (default, "Sarah", "Thompson", "[email protected]");
mysql> INSERT INTO customers VALUES (default, "Kenneth", "Anderson", "[email protected]");记录已写入数据库。但是,由于 Kafka Connect 未运行,watch-topic 不会记录任何更新。
| 在生产系统中,你会有足够的 broker 来处理生产者和消费者,并为每个主题维持最小数量的同步副本(ISR)。因此,如果失败的 broker 足够多,导致同步副本数量低于最小值,Kafka 将变得不可用。在这种情况下,生产者(如 Debezium 连接器)和消费者会等待 Kafka 集群或网络恢复。这意味着,在此期间,随着数据库中的数据不断变化,你的消费者可能看不到任何变更事件,因为此时没有产生任何变更事件。一旦 Kafka 集群重新启动或网络恢复,Debezium 将继续产生变更事件,你的消费者也会从中断的地方继续消费事件。 | |
|---|---|
| 3. 打开一个新的终端,用它来在容器中重启 Kafka Connect 服务。 |
此命令使用与你最初启动时相同的选项来启动 Kafka Connect:
$ docker run -it --rm --name connect -p 8083:8083 -e GROUP_ID=1 -e CONFIG_STORAGE_TOPIC=my_connect_configs -e OFFSET_STORAGE_TOPIC=my_connect_offsets -e STATUS_STORAGE_TOPIC=my_connect_statuses --link kafka:kafka --link mysql:mysql quay.io/debezium/connect:3.6Kafka Connect 服务启动后,会连接到 Kafka,读取上一个服务的配置,并启动已注册的连接器,这些连接器将从上次中断的位置继续执行。
以下是该重启服务输出的最后几行:
...
2021-11-30 01:49:07,938 INFO || Get all known binlogs from MySQL [io.debezium.connector.mysql.MySqlConnection]
2021-11-30 01:49:07,941 INFO || MySQL has the binlog file 'mysql-bin.000003' required by the connector [io.debezium.connector.mysql.MySqlConnectorTask]
2021-11-30 01:49:07,967 INFO || Requested thread factory for connector MySqlConnector, id = dbserver1 named = change-event-source-coordinator [io.debezium.util.Threads]
2021-11-30 01:49:07,968 INFO || Creating thread debezium-mysqlconnector-dbserver1-change-event-source-coordinator [io.debezium.util.Threads]
2021-11-30 01:49:07,968 INFO || WorkerSourceTask{id=inventory-connector-0} Source task finished initialization and start [org.apache.kafka.connect.runtime.WorkerSourceTask]
2021-11-30 01:49:07,971 INFO MySQL|dbserver1|snapshot Metrics registered [io.debezium.pipeline.ChangeEventSourceCoordinator]
2021-11-30 01:49:07,971 INFO MySQL|dbserver1|snapshot Context created [io.debezium.pipeline.ChangeEventSourceCoordinator]
2021-11-30 01:49:07,976 INFO MySQL|dbserver1|snapshot A previous offset indicating a completed snapshot has been found. Neither schema nor data will be snapshotted. [io.debezium.connector.mysql.MySqlSnapshotChangeEventSource]
2021-11-30 01:49:07,977 INFO MySQL|dbserver1|snapshot Snapshot ended with SnapshotResult [status=SKIPPED, offset=MySqlOffsetContext [sourceInfoSchema=Schema{io.debezium.connector.mysql.Source:STRUCT}, sourceInfo=SourceInfo [currentGtid=null, currentBinlogFilename=mysql-bin.000003, currentBinlogPosition=156, currentRowNumber=0, serverId=0, sourceTime=null, threadId=-1, currentQuery=null, tableIds=[], databaseName=null], snapshotCompleted=false, transactionContext=TransactionContext [currentTransactionId=null, perTableEventCount={}, totalEventCount=0], restartGtidSet=null, currentGtidSet=null, restartBinlogFilename=mysql-bin.000003, restartBinlogPosition=156, restartRowsToSkip=0, restartEventsToSkip=0, currentEventLengthInBytes=0, inTransaction=false, transactionId=null, incrementalSnapshotContext =IncrementalSnapshotContext [windowOpened=false, chunkEndPosition=null, dataCollectionsToSnapshot=[], lastEventKeySent=null, maximumKey=null]]] [io.debezium.pipeline.ChangeEventSourceCoordinator]
2021-11-30 01:49:07,981 INFO MySQL|dbserver1|streaming Requested thread factory for connector MySqlConnector, id = dbserver1 named = binlog-client [io.debezium.util.Threads]
2021-11-30 01:49:07,983 INFO MySQL|dbserver1|streaming Starting streaming [io.debezium.pipeline.ChangeEventSourceCoordinator]
...这些日志行说明,服务找到了上一个任务在关闭之前记录的偏移量,连接到 MySQL 数据库,并从该位置开始读取 binlog,进而根据自那一刻起 MySQL 数据库中的所有变更生成事件。
4. 切换到运行 watch-topic 的终端,查看 Kafka Connect 离线期间你所创建的那两条新记录对应的事件:
{"schema":{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"}],"optional":false,"name":"dbserver1.inventory.customers.Key"},"payload":{"id":1005}} {"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"first_name"},{"type":"string","optional":false,"field":"last_name"},{"type":"string","optional":false,"field":"email"}],"optional":true,"name":"dbserver1.inventory.customers.Value","field":"before"},{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"first_name"},{"type":"string","optional":false,"field":"last_name"},{"type":"string","optional":false,"field":"email"}],"optional":true,"name":"dbserver1.inventory.customers.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":true,"field":"version"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"server_id"},{"type":"int64","optional":false,"field":"ts_sec"},{"type":"string","optional":true,"field":"gtid"},{"type":"string","optional":false,"field":"file"},{"type":"int64","optional":false,"field":"pos"},{"type":"int32","optional":false,"field":"row"},{"type":"boolean","optional":true,"field":"snapshot"},{"type":"int64","optional":true,"field":"thread"},{"type":"string","optional":true,"field":"db"},{"type":"string","optional":true,"field":"table"}],"optional":false,"name":"io.debezium.connector.mysql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"},{"type":"int64","optional":true,"field":"ts_us"},{"type":"int64","optional":true,"field":"ts_ns"}],"optional":false,"name":"dbserver1.inventory.customers.Envelope","version":1},"payload":{"before":null,"after":{"id":1005,"first_name":"Sarah","last_name":"Thompson","email":"[email protected]"},"source":{"version":"3.6.3.Final","name":"dbserver1","server_id":223344,"ts_sec":1490635153,"gtid":null,"file":"mysql-bin.000003","pos":1046,"row":0,"snapshot":null,"thread":3,"db":"inventory","table":"customers"},"op":"c","ts_ms":1490635181455,"ts_us":1490635181455501,"ts_ns":1490635181455501571}}
{"schema":{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"}],"optional":false,"name":"dbserver1.inventory.customers.Key"},"payload":{"id":1006}} {"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"first_name"},{"type":"string","optional":false,"field":"last_name"},{"type":"string","optional":false,"field":"email"}],"optional":true,"name":"dbserver1.inventory.customers.Value","field":"before"},{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"first_name"},{"type":"string","optional":false,"field":"last_name"},{"type":"string","optional":false,"field":"email"}],"optional":true,"name":"dbserver1.inventory.customers.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":true,"field":"version"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"server_id"},{"type":"int64","optional":false,"field":"ts_sec"},{"type":"string","optional":true,"field":"gtid"},{"type":"string","optional":false,"field":"file"},{"type":"int64","optional":false,"field":"pos"},{"type":"int32","optional":false,"field":"row"},{"type":"boolean","optional":true,"field":"snapshot"},{"type":"int64","optional":true,"field":"thread"},{"type":"string","optional":true,"field":"db"},{"type":"string","optional":true,"field":"table"}],"optional":false,"name":"io.debezium.connector.mysql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"},{"type":"int64","optional":true,"field":"ts_us"},{"type":"int64","optional":true,"field":"ts_ns"}],"optional":false,"name":"dbserver1.inventory.customers.Envelope","version":1},"payload":{"before":null,"after":{"id":1006,"first_name":"Kenneth","last_name":"Anderson","email":"[email protected]"},"source":{"version":"3.6.3.Final","name":"dbserver1","server_id":223344,"ts_sec":1490635160,"gtid":null,"file":"mysql-bin.000003","pos":1356,"row":0,"snapshot":null,"thread":3,"db":"inventory","table":"customers"},"op":"c","ts_ms":1490635181456,"ts_us":1490635181456101,"ts_ns":1490635181456101571}}这些事件与之前看到的 create 事件类似。如你所见,即使 Debezium 未在运行,它仍然会报告数据库中的所有变更(只要在 MySQL 数据库从 binlog 中清除掉错过的提交之前将其重新启动即可)。
清理
教程结束后,你可以使用 Docker 停止所有正在运行的容器。
操作步骤
停止每个容器:
$ docker stop mysqlterm watcher connect mysql kafka
Docker 会停止每个容器。由于你启动容器时使用了 --rm 选项,Docker 也会将这些容器删除。
如果你使用的是 Podman,请运行以下命令:
$ podman pod kill dbz
$ podman pod rm dbz验证所有进程均已停止并被移除:
$ docker ps -a
如果任何进程仍在运行,请使用 docker stop <process-name> 或 docker stop <containerId> 停止它们。
后续步骤
完成本教程后,可以考虑以下后续步骤:
进一步探索本教程。
使用 MySQL 命令行客户端向数据库表中添加、修改和删除行,并观察其对主题的影响。你可能需要为每个主题分别运行一个
watch-topic命令。请注意,无法删除被外键引用的行。尝试使用 Postgres、MongoDB、SQL Server 和 Oracle 的 Debezium 连接器来运行本教程。
你可以使用位于 Debezium 示例仓库中的本教程的 Docker Compose 版本。其中提供了使用 MySQL、Postgres、MongoDB、SQL Server 和 Oracle 运行本教程的 Docker Compose 文件。
评论
登录后参与评论
KnowForge