【问题标题】:Transform / specify table name in jdbc sync connector在 jdbc 同步连接器中转换/指定表名
【发布时间】:2021-03-22 11:25:06
【问题描述】:

我需要将 SQL Server 数据库从本地位置迁移到 GoogleCloud 使用 confluent/kafka 来做

我确实有源 debezium 连接器

{
  "name": "mssql_src",
  "config": {
    "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector",
    "tasks.max": "1",
    "key.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    ...
    ...
    ...
    "transforms": "unwrap",
    "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
    "transforms.unwrap.add.fields": "op,table,source.ts_ms",
    "transforms.unwrap.delete.handling.mode": "rewrite",

    "transforms": "Reroute",
    "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
    "transforms.Reroute.topic.regex": "source_dbname.dbo(.*)",
    "transforms.Reroute.topic.replacement": "target_dbname$1"
  }
}

重新路由转换不能与 unwrap 结合使用, 我仍然得到 source_dbname.dbo.* 主题而不是 target_dbname.*

我需要将数据插入 target_dbname 数据库 jdbc 同步连接器有以下配置

{
  "name": "mssql_trg",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "tasks.max": "1",
    "topics.regex": "source_dbname.dbo.*",
    "table.name.format": "${topic}",
    "connection.url": "jdbc:sqlserver://xxx.xxx.xxx.xxx:1433;DatabaseName=rocketlawyer3",
    "connection.user": "sqlserver",
    "connection.password": "sqlserver",
    "dialect.name": "SqlServerDatabaseDialect",
    "insert.mode": "upsert",
    "auto.create": true,
    "auto.evolve": true,
    "pk.mode": "record_value"
  }
}

显然它失败了,因为所有 SQL 操作都将表引用为 source_database_name.dbo.table_name

这里有两个问题:

  1. 如何使用 table.name.format 仅为 table_name 更改此字符串 source_database_name.dbo.table_name 或其他选项

  2. 如何使转换(重新路由和展开)在源连接器中工作

【问题讨论】:

    标签: sql-server apache-kafka apache-kafka-connect confluent-platform debezium


    【解决方案1】:

    您需要使用chained transformations。只需将unwrapRerout SMT 组合在一起即可:

    {
      "name": "mssql_src",
      "config": {
        "connector.class": "io.debezium.connector.sqlserver.SqlServerConnector",
        "tasks.max": "1",
        "key.converter": "io.confluent.connect.avro.AvroConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        ...
        ...
        ...
        "transforms": "unwrap, Reroute",
        "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState",
        "transforms.unwrap.add.fields": "op,table,source.ts_ms",
        "transforms.unwrap.delete.handling.mode": "rewrite",
    
        "transforms.Reroute.type": "io.debezium.transforms.ByLogicalTableRouter",
        "transforms.Reroute.topic.regex": "([^.]+)\\.([^.]+)\\.([^.]+)",
        "transforms.Reroute.topic.replacement": "$3"
      }
    }
    

    【讨论】:

    • 不客气!如果我的回答对您有帮助,请接受并投票
    猜你喜欢
    • 2021-06-28
    • 2017-09-08
    • 2021-01-14
    • 2014-08-30
    • 2013-08-13
    • 2019-11-10
    • 2012-05-16
    • 2021-03-04
    • 2018-10-21
    相关资源
    最近更新 更多