You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何为JDBC Sink Connector指定Postgres目标表名?

Debezium同步Postgres表到JDBC Sink时目标表名无法指定的问题

使用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连接器自动匹配目标表名:

  1. 在transforms中添加routeTopic:
transforms = timestampConverter,replaceField,routeTopic
  1. 添加转换器配置,将原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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.02 22:07:02