【问题标题】:Debezium source task fails to reconnect to postgresql DB when DB container is re-created重新创建数据库容器时,Debezium 源任务无法重新连接到 postgresql 数据库
【发布时间】:2018-12-24 22:33:00
【问题描述】:

我们有一个 kubernetes 集群,其中 Debezium 作为来自 Postgresql 的源任务运行并写入 kafka。 Debezium、postgres 和 kafka 都在不同的 pod 中运行。 当 postgres pod 被删除并且 kubernetes 重新创建 pod 时,debezium pod 无法重新连接。 来自 debezium pod 的日志:

    2018-07-17 08:31:38,311 ERROR  ||  WorkerSourceTask{id=inventory-connector-0} Task is being killed and will not recover until manually restarted   [org.apache.kafka.connect.runtime.WorkerTask]
    2018-07-17 08:31:38,311 INFO   ||  [Producer clientId=producer-4] Closing the Kafka producer with timeoutMillis = 30000 ms.   [org.apache.kafka.clients.producer.KafkaProducer]

Debezium 会继续尝试每隔一段时间刷新未完成的消息,但会出现以下异常:

    2018-07-17 08:32:38,167 ERROR  ||  WorkerSourceTask{id=inventory-connector-0} Exception thrown while calling task.commit()   [org.apache.kafka.connect.runtime.WorkerSourceTask]
    org.apache.kafka.connect.errors.ConnectException: org.postgresql.util.PSQLException: Database connection failed when writing to copy
    at io.debezium.connector.postgresql.RecordsStreamProducer.commit(RecordsStreamProducer.java:151)
    at io.debezium.connector.postgresql.PostgresConnectorTask.commit(PostgresConnectorTask.java:138)
    at org.apache.kafka.connect.runtime.WorkerSourceTask.commitSourceTask(WorkerSourceTask.java:437)
    at org.apache.kafka.connect.runtime.WorkerSourceTask.commitOffsets(WorkerSourceTask.java:378)
    at org.apache.kafka.connect.runtime.SourceTaskOffsetCommitter.commit(SourceTaskOffsetCommitter.java:108)
    at org.apache.kafka.connect.runtime.SourceTaskOffsetCommitter.access$000(SourceTaskOffsetCommitter.java:45)
    at org.apache.kafka.connect.runtime.SourceTaskOffsetCommitter$1.run(SourceTaskOffsetCommitter.java:82)
    at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
    at java.util.concurrent.FutureTask.runAndReset(FutureTask.java:308)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.access$301(ScheduledThreadPoolExecutor.java:180)
    at java.util.concurrent.ScheduledThreadPoolExecutor$ScheduledFutureTask.run(ScheduledThreadPoolExecutor.java:294)
    at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
    at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
    at java.lang.Thread.run(Thread.java:748)
    Caused by: org.postgresql.util.PSQLException: Database connection failed when writing to copy
    at org.postgresql.core.v3.QueryExecutorImpl.flushCopy(QueryExecutorImpl.java:942)
    at org.postgresql.core.v3.CopyDualImpl.flushCopy(CopyDualImpl.java:23)
    at org.postgresql.core.v3.replication.V3PGReplicationStream.updateStatusInternal(V3PGReplicationStream.java:176)
    at org.postgresql.core.v3.replication.V3PGReplicationStream.forceUpdateStatus(V3PGReplicationStream.java:99)
    at io.debezium.connector.postgresql.connection.PostgresReplicationConnection$1.doFlushLsn(PostgresReplicationConnection.java:246)
    at io.debezium.connector.postgresql.connection.PostgresReplicationConnection$1.flushLsn(PostgresReplicationConnection.java:239)
    at io.debezium.connector.postgresql.RecordsStreamProducer.commit(RecordsStreamProducer.java:146)
    ... 13 more
    Caused by: java.net.SocketException: Broken pipe (Write failed)
    at java.net.SocketOutputStream.socketWrite0(Native Method)
    at java.net.SocketOutputStream.socketWrite(SocketOutputStream.java:111)
    at java.net.SocketOutputStream.write(SocketOutputStream.java:155)
    at java.io.BufferedOutputStream.flushBuffer(BufferedOutputStream.java:82)
    at java.io.BufferedOutputStream.flush(BufferedOutputStream.java:140)
    at org.postgresql.core.PGStream.flush(PGStream.java:553)
    at org.postgresql.core.v3.QueryExecutorImpl.flushCopy(QueryExecutorImpl.java:939)
    ... 19 more

有没有办法让 debezium 在可用时重新建立与 postgres 的连接? 还是我缺少一些配置?

  • Debezium 0.8 版
  • kubernetes 版本 1.10.3
  • postgres 9.6 版

【问题讨论】:

  • 你能分享更多关于 kubernetes 的细节,例如 yaml 文件吗?如果您提供有关环境的更多详细信息,将会很有帮助。

标签: postgresql apache-kafka kubernetes debezium


【解决方案1】:

看起来这是一个常见问题,并且在 debezium 和 kafka 中都有开放的功能请求

https://issues.jboss.org/browse/DBZ-248

https://issues.apache.org/jira/browse/KAFKA-5352

虽然这些是开放的,但看起来这是预期的行为

作为一种解决方法,我已将此活性探针添加到部署中

    livenessProbe:
        exec:
          command:
          - sh
          - -ec
          - ipaddress=$(ip addr | grep 'state UP' -A2 | tail -n1 | awk '{print $2}' | cut -f1  -d'/'); reply=$(curl -s $ipaddress:8083/connectors/inventory-connector/status | grep -o RUNNING | wc -l); if [ $reply -lt 2 ]; then exit 1; fi;
        initialDelaySeconds: 30
        periodSeconds: 5

第一个子句获取容器IP地址:

    ipaddress=$(ip addr | grep 'state UP' -A2 | tail -n1 | awk '{print $2}' | cut -f1 -d'/');

第二个子句发出请求并计算响应 json 中“RUNNING”的实例数:

    reply=$(curl -s $ipaddress:8083/connectors/inventory-connector/status | grep -o RUNNING | wc -l);

如果 'RUNNING' 出现的次数少于两次,则第三个子句返回退出代码 1

    if [ $reply -lt 2 ]; then exit 1; fi

它似乎正在进行初始测试 - 即重新启动 postgres DB 会触发 debezium 容器的重新启动。我猜想像这样的脚本(虽然可能是'robustified')可以包含在图像中以方便探测。

【讨论】:

  • 嘿,是的,目前我们没有在连接器本身内实现任何重新连接的概念。想法是,这可以通过编排层(在您的情况下为 Kubernetes)来完成,该层将通过 REST API 监控连接器状态,并在任务失败时重新启动。不过,这并不是一成不变的,我们最终也可能会提供内部解决方案,但我们还没有做到。
  • 好的,谢谢@Gunnar。很高兴知道我没有遗漏任何东西,这是意料之中的。我将尝试向 Debezium pod 添加 kubernetes liveness 和 readiness 探针来处理这种情况。
  • 是的,就是这样。您可以通过如下 URL 从 Kafka Connect REST API 获取连接器状态:connect-host:8083/connectors/your-connector/status。话虽如此,如果您设法实现上述探测,也许您有兴趣在 debezium.io 的博客上写一篇关于该主题的客座文章?它肯定也会对其他人有所帮助。
  • @Gunnar - 看起来 api 在任务正在运行和失败时都返回 200 OK。 kubelet 依靠状态码来确定 pod 是否处于活动状态/准备就绪状态。也许可以执行命令并解析 json 响应 - 我会试一试
  • 您可以使用jq 检索确切的任务状态: curl -s $ipaddress:8083/connectors/hiking-connector/status | jq '.tasks[0].state'
猜你喜欢
  • 2017-09-17
  • 1970-01-01
  • 2020-12-25
  • 1970-01-01
  • 2021-07-21
  • 2012-08-04
  • 2017-09-05
  • 2023-04-02
  • 1970-01-01
相关资源
最近更新 更多