如何为JDBC Sink Connector指定Postgres目标表名?
使用Debezium Postgres源连接器(版本debezium/debezium-connector-postgresql:2.2.1)和JDBC Sink连接器(版本confluentinc/kafka-connect-jdbc:10.7.4),需将源库public schema下的parametros表数据同步到目标库的parametros_sistema表中。已通过org.apache.kafka.connect.transforms.ReplaceField$Value转换器完成字段名映射,但启动Sink连接器时触发以下错误:
ERROR [collaborator-sink-postgres-connector|task-0] WorkerSinkTask{id=collaborator-sink-postgres-connector-0} Task threw an uncaught and unrecoverable exception. Task is being killed and will not recover until manually restarted. Error: Table "collab-service"."public"."parametros" is missing and auto-creation is disabled (org.apache.kafka.connect.runtime.WorkerSinkTask:586)
核心问题:连接器默认尝试匹配源表名称创建目标表,但自动创建功能已关闭,且未正确指定目标表名。
附:现有配置与Topic消息
源连接器配置
name = collaborator-postgres-connector connector.class = io.debezium.connector.postgresql.PostgresConnector tasks.max = 1 topic.prefix = collab-service database.hostname = host database.server.name = collaborator-postgres-server database.port = 5432 database.user = user database.password = password database.dbname = sflm_dev plugin.name = pgoutput slot.name = collab_test decimal.handling.mode = double snapshot.mode = always schema.name.adjustment.mode = none table.include.list = public.parametros database.history.kafka.topic = postgres_history database.history.kafka.bootstrap.servers = kafka:9092 message.prefix.include.list = after
Sink连接器配置
name = collaborator-sink-postgres-connector connector.class = io.confluent.connect.jdbc.JdbcSinkConnector tasks.max = 1 topics = collab-service.public.parametros connection.url = host connection.user = user connection.password = password database = collaborator_suite auto.create = false insert.mode = upsert pk.mode = record_key transforms = timestampConverter,replaceField transforms.timestampConverter.type = org.apache.kafka.connect.transforms.TimestampConverter$Value transforms.timestampConverter.field = para_dt_cadastro, para_dt_ult_alt transforms.timestampConverter.target.type = Timestamp transforms.replaceField.type = org.apache.kafka.connect.transforms.ReplaceField$Value transforms.replaceField.renames = para_cd_id:pasi_cd_id, para_tx_dominio:pasi_tx_dominio, para_tx_descricao:pasi_tx_descricao, para_tx_valor:pasi_tx_valor, para_tx_tipo:pasi_tx_tipo, para_dt_ult_alt:pasi_dt_ult_alt, para_dt_cadastro:pasi_dt_cadastro, para_dt_cadastro:pasi_dt_cadastro, para_dt_ult_alt:pasi_dt_ult_alt, usua_cd_id_cadastro:usua_cd_id_cadastro, usua_cd_id_ult_alt:usua_cd_id_ult_alt, para_tx_sistema:pasi_tx_sistema db.timezone = UTC pk.fields = para_cd_id
Topic消息示例
[ { "topic": "collab-service.public.parametros", "partition": 0, "offset": 1, "timestamp": 1704296444971, "timestampType": "CREATE_TIME", "headers": [], "key": "Struct{para_cd_id=2}", "value": { "schema": { "type": "struct", "fields": [ { "type": "struct", "fields": [ { "type": "int32", "optional": false, "default": 0, "field": "para_cd_id" }, { "type": "string", "optional": false, "field": "para_tx_dominio" }, { "type": "string", "optional": false, "field": "para_tx_descricao" }, { "type": "string", "optional": false, "field": "para_tx_valor" }, { "type": "string", "optional": false, "field": "para_tx_tipo" }, { "type": "int64", "optional": true, "name": "io.debezium.time.MicroTimestamp", "version": 1, "field": "para_dt_cadastro" }, { "type": "int64", "optional": true, "name": "io.debezium.time.MicroTimestamp", "version": 1, "field": "para_dt_ult_alt" }, { "type": "int64", "optional": true, "field": "usua_cd_id_cadastro" }, { "type": "double", "optional": true, "field": "usua_cd_id_ult_alt" }, { "type": "string", "optional": true, "field": "para_tx_sistema" } ], "optional": true, "name": "collab-service.public.parametros.Value", "field": "before" }, { "type": "struct", "fields": [ { "type": "int32", "optional": false, "default": 0, "field": "para_cd_id" }, { "type": "string", "optional": false, "field": "para_tx_dominio" }, { "type": "string", "optional": false, "field": "para_tx_descricao" }, { "type": "string", "optional": false, "field": "para_tx_valor" }, { "type": "string", "optional": false, "field": "para_tx_tipo" }, { "type": "int64", "optional": true, "name": "io.debezium.time.MicroTimestamp", "version": 1, "field": "para_dt_cadastro" }, { "type": "int64", "optional": true, "name": "io.debezium.time.MicroTimestamp", "version": 1, "field": "para_dt_ult_alt" }, { "type": "int64", "optional": true, "field": "usua_cd_id_cadastro" }, { "type": "double", "optional": true, "field": "usua_cd_id_ult_alt" }, { "type": "string", "optional": true, "field": "para_tx_sistema" } ], "optional": true, "name": "collab-service.public.parametros.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,incremental" }, "default": "false", "field": "snapshot" }, { "type": "string", "optional": false, "field": "db" }, { "type": "string", "optional": true, "field": "sequence" }, { "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" }, { "type": "struct", "fields": [ { "type": "string", "optional": false, "field": "id" }, { "type": "int64", "optional": false, "field": "total_order" }, { "type": "int64", "optional": false, "field": "data_collection_order" } ], "optional": true, "name": "event.block", "version": 1, "field": "transaction" } ], "optional": false, "name": "collab-service.public.parametros.Envelope", "version": 1 }, "payload": { "before": null, "after": { "para_cd_id": 2, "para_tx_dominio": "SISTEMA_CLOUD_PROJECTID", "para_tx_descricao": "Variável para referenciar o id do projeto na nuvem do Google", "para_tx_valor": "collaborator-364516", "para_tx_tipo": "SISTEMA", "para_dt_cadastro": 1669047531542921, "para_dt_ult_alt": 1669047531542921, "usua_cd_id_cadastro": null, "usua_cd_id_ult_alt": null, "para_tx_sistema": "R2D2" }, "source": { "version": "2.2.1.Final", "connector": "postgresql", "name": "collab-service", "ts_ms": 1704296442190, "snapshot": "last", "db": "sflm_dev", "sequence": "[null,\"17335024878048\"]", "schema": "public", "table": "parametros", "txId": 989004, "lsn": 17335024878048, "xmin": null }, "op": "r", "ts_ms": 1704296444402, "transaction": null } } } ]
解决方法
方式1:直接通过table.name.format指定目标表名(推荐)
这是最直接的方案,无需额外转换器。在Sink连接器配置中添加以下参数,明确指定目标表的schema和名称:
table.name.format = public.parametros_sistema
同时需要更新主键配置,匹配目标表的主键字段(因已通过ReplaceField将para_cd_id映射为pasi_cd_id):
pk.fields = pasi_cd_id
修改后的完整Sink配置如下:
name = collaborator-sink-postgres-connector connector.class = io.confluent.connect.jdbc.JdbcSinkConnector tasks.max = 1 topics = collab-service.public.parametros connection.url = host connection.user = user connection.password = password database = collaborator_suite auto.create = false insert.mode = upsert pk.mode = record_key # 指定目标表名 table.name.format = public.parametros_sistema transforms = timestampConverter,replaceField transforms.timestampConverter.type = org.apache.kafka.connect.transforms.TimestampConverter$Value transforms.timestampConverter.field = para_dt_cadastro, para_dt_ult_alt transforms.timestampConverter.target.type = Timestamp transforms.replaceField.type = org.apache.kafka.connect.transforms.ReplaceField$Value transforms.replaceField.renames = para_cd_id:pasi_cd_id, para_tx_dominio:pasi_tx_dominio, para_tx_descricao:pasi_tx_descricao, para_tx_valor:pasi_tx_valor, para_tx_tipo:pasi_tx_tipo, para_dt_ult_alt:pasi_dt_ult_alt, para_dt_cadastro:pasi_dt_cadastro, usua_cd_id_cadastro:usua_cd_id_cadastro, usua_cd_id_ult_alt:usua_cd_id_ult_alt, para_tx_sistema:pasi_tx_sistema db.timezone = UTC # 匹配目标表主键 pk.fields = pasi_cd_id
方式2:使用RegexRouter转换器修改Topic名称(多表场景适用)
如果需要批量处理多表映射,可通过转换器修改Topic名称,让Sink连接器自动匹配目标表名:
- 在
transforms中添加routeTopic:
transforms = timestampConverter,replaceField,routeTopic
- 添加转换器配置,将原Topic中的表名替换为目标表名:
transforms.routeTopic.type = org.apache.kafka.connect.transforms.RegexRouter transforms.routeTopic.regex = collab-service.public.(.*) transforms.routeTopic.replacement = $1_sistema
此配置会将collab-service.public.parametros转换为parametros_sistema,连接器会自动匹配目标库中的同名表。
内容的提问来源于stack exchange,提问作者leandrofita

