Postgres Debezium CDC不会发布对Kafka的更改



我当前的测试配置如下所示:

version: '3.7'
services:
postgres:
image: debezium/postgres
restart: always
ports:
- "5432:5432"
zookeeper:
image: debezium/zookeeper
ports:
- "2181:2181"
- "2888:2888"
- "3888:3888"
kafka:
image: debezium/kafka
restart: always
ports:
- "9092:9092"
links:
- zookeeper
depends_on:
- zookeeper
environment:
- ZOOKEEPER_CONNECT=zookeeper:2181
- KAFKA_GROUP_MIN_SESSION_TIMEOUT_MS=250
connect:
image: debezium/connect
restart: always
ports:
- "8083:8083"
links:
- zookeeper
- postgres
- kafka
depends_on:
- zookeeper
- postgres
- kafka
environment:
- BOOTSTRAP_SERVERS=kafka:9092
- GROUP_ID=1
- CONFIG_STORAGE_TOPIC=my_connect_configs
- OFFSET_STORAGE_TOPIC=my_connect_offsets
- STATUS_STORAGE_TOPIC=my_source_connect_statuses

我像这样使用 docker 撰写运行它:

$ docker-compose up

而且我没有看到任何错误消息。似乎一切运行正常。如果我执行docker ps,我看到所有服务都在运行。

为了检查 Kafka 是否正在运行,我在 Python 中制作了 Kafka 生产者和 Kafka 消费者:

# producer. I run it in one console window
from kafka import KafkaProducer
from json import dumps
from time import sleep
producer = KafkaProducer(bootstrap_servers=['localhost:9092'], value_serializer=lambda x: dumps(x).encode('utf-8'))
for e in range(1000):
data = {'number' : e}
producer.send('numtest', value=data)
sleep(5)
# consumer. I run it in other colsole window
from kafka import KafkaConsumer
from json import loads
consumer = KafkaConsumer(
'numtest',
bootstrap_servers=['localhost:9092'],
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='my-group',
value_deserializer=lambda x: loads(x.decode('utf-8')))
for message in consumer:
print(message)

而且效果绝对很棒。我看到我的生产者如何发布消息,我看到它们如何在消费者窗口中被使用。

现在我想让CDC发挥作用。首先,在Postgres容器中,我将角色密码设置为postgrespostgres

$ su postgres
$ psql
psql> password postgres
Enter new password: postgres

然后我创建了一个新的数据库test

psql> CREATE DATABASE test;

我创建了一个表:

psql> c test;
test=# create table mytable (id serial, name varchar(128), primary key(id));

最后,对于我的Debezium CDC堆栈,我创建了一个连接器:

$ curl -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/ -d '{
"name": "test-connector",
"config": {
"connector.class": "io.debezium.connector.postgresql.PostgresConnector",
"tasks.max": "1",
"plugin.name": "pgoutput",
"database.hostname": "postgres",
"database.port": "5432",
"database.user": "postgres",
"database.password": "postgres",
"database.dbname" : "test",
"database.server.name": "postgres",
"database.whitelist": "public.mytable",
"database.history.kafka.bootstrap.servers": "localhost:9092",
"database.history.kafka.topic": "public.some_topic"
}
}'
{"name":"test-connector","config":{"connector.class":"io.debezium.connector.postgresql.PostgresConnector","tasks.max":"1","plugin.name":"pgoutput","database.hostname":"postgres","database.port":"5432","database.user":"postgres","database.password":"postgres","database.dbname":"test","database.server.name":"postgres","database.whitelist":"public.mytable","database.history.kafka.bootstrap.servers":"localhost:9092","database.history.kafka.topic":"public.some_topic","name":"test-connector"},"tasks":[],"type":"source"}

如您所见,我的连接器是在没有任何错误的情况下创建的。现在我希望Debezium CDC发布对Kafka主题public.some_topic的所有更改。为了检查这一点,我创建了一个新的 Kafka 消费者:

from kafka import KafkaConsumer
from json import loads
consumer = KafkaConsumer(
'public.some_topic',
bootstrap_servers=['localhost:9092'],
auto_offset_reset='earliest',
enable_auto_commit=True,
group_id='my-group',
value_deserializer=lambda x: loads(x.decode('utf-8')))
for message in consumer:
print(message)

与第一个示例的唯一区别是,我正在观看public.some_topic。然后,我转到数据库控制台并进行插入:

test=# insert into mytable (name) values ('Tom Cat');    
INSERT 0 1
test=#

因此,插入了一个新值,但我看到消费者窗口中没有发生任何事情。换句话说,Debezium 不会向 Kafkapublic.some_topic发布事件。这有什么问题,我该如何解决?

使用您的 Docker Compose,创建连接器时,我在 Kafka Connect 工作线程日志中看到此错误:

Caused by: org.postgresql.util.PSQLException: ERROR: could not access file "pgoutput": No such file or directory
at org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2505)
at org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2241)
at org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:310)
at org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:447)
at org.postgresql.jdbc.PgStatement.execute(PgStatement.java:368)
at org.postgresql.jdbc.PgStatement.executeWithFlags(PgStatement.java:309)
at org.postgresql.jdbc.PgStatement.executeCachedSql(PgStatement.java:295)
at org.postgresql.jdbc.PgStatement.executeWithFlags(PgStatement.java:272)
at org.postgresql.jdbc.PgStatement.execute(PgStatement.java:267)
at io.debezium.connector.postgresql.connection.PostgresReplicationConnection.createReplicationSlot(PostgresReplicationConnection.java:288)
at io.debezium.connector.postgresql.PostgresConnectorTask.start(PostgresConnectorTask.java:126)
... 9 more

如果您使用 Kafka Connect REST API 查询任务,这也反映在任务的状态中:

curl -s "http://localhost:8083/connectors?expand=info&expand=status" | jq '."test-connector".status'
{
"name": "test-connector",
"connector": {
"state": "RUNNING",
"worker_id": "192.168.16.5:8083"
},
"tasks": [
{
"id": 0,
"state": "FAILED",
"worker_id": "192.168.16.5:8083",
"trace": "org.apache.kafka.connect.errors.ConnectException: org.postgresql.util.PSQLException: ERROR: could not access file "pgoutput": No such file or directoryntat io.debezium.connector.postgresql.PostgresConnectorTask.start(PostgresConnectorTask.java:129)ntat io.debezium.connector.common.BaseSourceTask.start(BaseSourceTask.java:49)ntat org.apache.kafka.connect.runtime.WorkerSourceTask.execute(WorkerSourceTask.java:208)ntat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:177)ntat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:227)ntat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)ntat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)ntat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)ntat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)ntat java.base/java.lang.Thread.run(Thread.java:834)nCaused by: org.postgresql.util.PSQLException: ERROR: could not access file "pgoutput": No such file or directoryntat org.postgresql.core.v3.QueryExecutorImpl.receiveErrorResponse(QueryExecutorImpl.java:2505)ntat org.postgresql.core.v3.QueryExecutorImpl.processResults(QueryExecutorImpl.java:2241)ntat org.postgresql.core.v3.QueryExecutorImpl.execute(QueryExecutorImpl.java:310)ntat org.postgresql.jdbc.PgStatement.executeInternal(PgStatement.java:447)ntat org.postgresql.jdbc.PgStatement.execute(PgStatement.java:368)ntat org.postgresql.jdbc.PgStatement.executeWithFlags(PgStatement.java:309)ntat org.postgresql.jdbc.PgStatement.executeCachedSql(PgStatement.java:295)ntat org.postgresql.jdbc.PgStatement.executeWithFlags(PgStatement.java:272)ntat org.postgresql.jdbc.PgStatement.execute(PgStatement.java:267)ntat io.debezium.connector.postgresql.connection.PostgresReplicationConnection.createReplicationSlot(PostgresReplicationConnection.java:288)ntat io.debezium.connector.postgresql.PostgresConnectorTask.start(PostgresConnectorTask.java:126)nt... 9 moren"
}
],
"type": "source"

您正在运行的 Postgres 版本是

postgres=# SHOW server_version;
server_version
----------------
9.6.16

pgoutput仅在版本 10> = 可用。

我将您的 Docker Compose 更改为使用版本 10:

image: debezium/postgres:10

在弹跳堆栈以干净启动并按照您的说明进行操作后,我得到了一个正在运行的连接器:

curl -s "http://localhost:8083/connectors?expand=info&expand=status" | 
jq '. | to_entries[] | [ .value.info.type, .key, .value.status.connector.state,.value.status.tasks[].state,.value.info.config."connector.class"]|join(":|:")' | 
column -s : -t| sed 's/"//g'| sort
source  |  test-connector  |  RUNNING  |  RUNNING  |  io.debezium.connector.postgresql.PostgresConnector

和 Kafka 主题中的数据:

$ docker exec kafkacat kafkacat -b kafka:9092 -t postgres.public.mytable -C
{"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":true,"field":"name"}],"optional":true,"name":"postgres.public.mytable.Value","field":"before"},{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":true,"field":"name"}],"optional":true,"name":"postgres.public.mytable.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"version"},{"type":"string","optional":false,"field":"connector"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"ts_ms"},{"type":"string","optional":true,"name":"io.debezium.data.Enum","version":1,"parameters":{"allowed":"true,last,false"},"default":"false","field":"snapshot"},{"type":"string","optional":false,"field":"db"},{"type":"string","optional":false,"field":"schema"},{"type":"string","optional":false,"field":"table"},{"type":"int64","optional":true,"field":"txId"},{"type":"int64","optional":true,"field":"lsn"},{"type":"int64","optional":true,"field":"xmin"}],"optional":false,"name":"io.debezium.connector.postgresql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"}],"optional":false,"name":"postgres.public.mytable.Envelope"},"payload":{"before":null,"after":{"id":1,"name":"Tom Cat"},"source":{"version":"1.0.0.Final","connector":"postgresql","name":"postgres","ts_ms":1579172192292,"snapshot":"false","db":"test","schema":"public","table":"mytable","txId":561,"lsn":24485520,"xmin":null},"op":"c","ts_ms":1579172192347}}% Reached end of topic postgres.public.mytable [0] at offset 1

我将 kafkacat 添加到您的 Docker Compose 中:

kafkacat:
image: edenhill/kafkacat:1.5.0
container_name: kafkacat
entrypoint: 
- /bin/sh 
- -c 
- |
while [ 1 -eq 1 ];do sleep 60;done

编辑:保留以前的答案,因为它仍然有用且相关:

Debezium 将根据表的名称向主题写入消息。在您的示例中,这将是postgres.test.mytable.

这就是为什么kafkacat很有用,因为您可以运行

kafkacat -b broker:9092 -L 

以查看所有主题和分区的列表。一旦你有了主题

kafkacat -b broker:9092 -t postgres.test.mytable -C

从中阅读。

查看有关 kafkacat 的详细信息,包括如何使用 Docker 运行它

这里还有一个关于Docker Compose的演示。

相关内容

  • 没有找到相关文章

最新更新