【问题标题】:How kafka fetch data only when new row insert or update old row from mysql database using kafka streaming仅当新行使用kafka流从mysql数据库中插入或更新旧行时,kafka如何获取数据
【发布时间】:2021-08-30 10:49:27
【问题描述】:

我正在使用融合平台 JDBC 连接器,它是将数据从 mysql 流式传输到 kafka 消费者。通过应用程序将数据插入另一个数据库以进行报告。 这里的问题只是在某个时间间隔内一次又一次地流式传输所有数据。实际上只需要那些新插入的数据或以前记录的任何更新。

根据时间戳不能做,因为表不包含任何时间列。并且根据增量 id 也是不可能的。请分享任何解决方案。

我有示例配置文件。

demo.json
{
    "name":"mysql-connector-demo",
    "config":{
                "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector",
                "key.converter":"org.apache.kafka.connect.json.JsonConverter",
                "value.converter":"org.apache.kafka.connect.json.JsonConverter",
                "connection.url": "jdbc:mysql://localhost:3306/test",
                "connection.user": "root",
                "connection.password": "1234",
                "topic.prefix": "test",
                "catalog.pattern":"test",
                "mode": "bulk",
                "validate.non.null": false,
                "query": "select * from test ",
                "table.types": "TABLE",
                "topic.prefix": "test-jdbc-",
                "poll.interval.ms": 10000
                "schema.ignore": true,
                "key.converter.schemas.enable": "false",
                "value.converter.schemas.enable": "false"
    }
} 

但是这里新插入的记录和新更新的记录不影响kafka消费者。

【问题讨论】:

    标签: jdbc apache-kafka apache-kafka-connect confluent-platform


    【解决方案1】:

    只想要那些新插入的数据或以前记录的任何更新

    那么您应该使用更改数据捕获 (CDC) 而不是 JDBC 轮询。

    Debezium 就是这样一种解决方案

    【讨论】:

    • 对于 Debezium 需要特别更改配置文件 binlog-format = ROW。但是我使用的是mysql的master和replica结构,它已经为所有机器设置了binlog-format = MIXED。 Debzium 不允许 MIXED 格式并且不能更改,因为 master 和所有副本都已经配置了 MIXED。
    • 好吧,JDBC 连接器抓取整行,而不是您感兴趣的子集,并且批量配置模式将连续读取整个表,如文档中所述。也许您应该添加一个 lastUpdated 列,您可以将其添加到默认为当前时间的每一行
    猜你喜欢
    • 1970-01-01
    • 2018-11-01
    • 2019-12-20
    • 2018-02-19
    • 1970-01-01
    • 2017-05-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多