【发布时间】: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
这里有两个问题:
-
如何使用 table.name.format 仅为 table_name 更改此字符串 source_database_name.dbo.table_name 或其他选项
-
如何使转换(重新路由和展开)在源连接器中工作
【问题讨论】:
标签: sql-server apache-kafka apache-kafka-connect confluent-platform debezium