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

Kafka Connect同步Redis到PostgreSQL报错:MAP类型无对应SQL列类型

Kafka Connect Redis源到PostgreSQL目标同步失败问题解决

环境与连接器配置

Redis源连接器配置

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/   -d  '{"name": "redissource","config": {"connector.class": "com.redis.kafka.connect.RedisSourceConnector", "tasks.max": "1","topics": "mystream","redis.uri": "redis://virginia:virginia@172.18.1.41:1235","redis.cluster.enabled": "false",  "redis.stream.name":"mystream",   "key.converter":"org.apache.kafka.connect.storage.StringConverter",  "value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://schema-registry:8081" }}'

PostgreSQL目标连接器配置

curl -i -X POST -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/   -d  '{"name": "redtopgsink","config": { "connector.class"  : "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url":"jdbc:postgresql://navi1085_postgres_1:5432/postgres", "connection.user" :"postgres",  "connection.password" :"postgres_nave1085",  "topics"  :"mystream", "key.converter":"org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://schema-registry:8081" ,   "insert.mode" :"upsert", "schema.pattern":"public","auto.create"  :"true" ,"pk.mode"   :"record_key","pk.fields":"sensor_id",   "delete.enabled"  :"true","auto.evolve":"true"}}'

错误信息

{"name":"redtopgsink","connector":{"state":"RUNNING","worker_id":"connect:8083"},"tasks":[{"id":0,"state":"FAILED","worker_id":"connect:8083","trace":"org.apache.kafka.connect.errors.ConnectException: Exiting WorkerSinkTask due to unrecoverable exception.\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:618)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.poll(WorkerSinkTask.java:334)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.iteration(WorkerSinkTask.java:235)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.execute(WorkerSinkTask.java:204)\n\tat org.apache.kafka.connect.runtime.WorkerTask.doRun(WorkerTask.java:200)\n\tat org.apache.kafka.connect.runtime.WorkerTask.run(WorkerTask.java:255)\n\tat java.base/java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:515)\n\tat java.base/java.util.concurrent.FutureTask.run(FutureTask.java:264)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1128)\n\tat java.base/java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:628)\n\tat java.base/java.lang.Thread.run(Thread.java:829)\nCaused by: org.apache.kafka.connect.errors.ConnectException: null (MAP) type doesn't have a mapping to the SQL database column type\n\tat io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.getSqlType(GenericDatabaseDialect.java:1945)\n\tat io.confluent.connect.jdbc.dialect.PostgreSqlDatabaseDialect.getSqlType(PostgreSqlDatabaseDialect.java:332)\n\tat io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.writeColumnSpec(GenericDatabaseDialect.java:1861)\n\tat io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.lambda$writeColumnsSpec$39(GenericDatabaseDialect.java:1850)\n\tat io.confluent.connect.jdbc.util.ExpressionBuilder.append(ExpressionBuilder.java:560)\n\tat io.confluent.connect.jdbc.util.ExpressionBuilder$BasicListBuilder.of(ExpressionBuilder.java:599)\n\tat io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.writeColumnsSpec(GenericDatabaseDialect.java:1852)\n\tat io.confluent.connect.jdbc.dialect.GenericDatabaseDialect.buildCreateTableStatement(GenericDatabaseDialect.java:1769)\n\tat io.confluent.connect.jdbc.sink.DbStructure.create(DbStructure.java:121)\n\tat io.confluent.connect.jdbc.sink.DbStructure.createOrAmendIfNecessary(DbStructure.java:67)\n\tat io.confluent.connect.jdbc.sink.BufferedRecords.add(BufferedRecords.java:122)\n\tat io.confluent.connect.jdbc.sink.JdbcDbWriter.write(JdbcDbWriter.java:74)\n\tat io.confluent.connect.jdbc.sink.JdbcSinkTask.put(JdbcSinkTask.java:84)\n\tat org.apache.kafka.connect.runtime.WorkerSinkTask.deliverMessages(WorkerSinkTask.java:584)\n\t... 10 more\n"}],"type":"sink"}

问题分析

核心错误是Avro中的MAP类型没有对应的PostgreSQL列类型映射。启用auto.create: true时,JDBC Sink连接器会尝试根据Avro Schema自动创建PostgreSQL表,但默认无法将MAP类型映射到合适的SQL类型,导致表创建失败,进而任务报错。

解决方案

方案1:手动创建目标表,指定MAP类型对应PostgreSQL类型

PostgreSQL的JSONB类型可直接对应Avro的MAP类型,先手动创建表:

CREATE TABLE public.mystream (
    sensor_id VARCHAR PRIMARY KEY,
    -- 替换为实际的MAP字段名,比如"data"
    data JSONB,
    -- 根据Avro Schema添加其他字段,如时间戳
    event_timestamp TIMESTAMP
);

修改Sink连接器配置,关闭自动建表:

curl -i -X PUT -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/redtopgsink/config -d  '{"connector.class"  : "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url":"jdbc:postgresql://navi1085_postgres_1:5432/postgres", "connection.user" :"postgres",  "connection.password" :"postgres_nave1085",  "topics"  :"mystream", "key.converter":"org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://schema-registry:8081" ,   "insert.mode" :"upsert", "schema.pattern":"public","auto.create"  :"false" ,"pk.mode"   :"record_key","pk.fields":"sensor_id",   "delete.enabled"  :"true","auto.evolve":"true"}'

方案2:使用SMT将MAP类型转换为JSON字符串

通过Kafka Connect的Single Message Transform(SMT)将MAP转为字符串,存储到PostgreSQL的TEXT列中。修改Sink连接器配置添加SMT:

curl -i -X PUT -H "Accept:application/json" -H "Content-Type:application/json" localhost:8083/connectors/redtopgsink/config -d  '{"connector.class"  : "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url":"jdbc:postgresql://navi1085_postgres_1:5432/postgres", "connection.user" :"postgres",  "connection.password" :"postgres_nave1085",  "topics"  :"mystream", "key.converter":"org.apache.kafka.connect.storage.StringConverter", "value.converter": "io.confluent.connect.avro.AvroConverter","value.converter.schema.registry.url": "http://schema-registry:8081" ,   "insert.mode" :"upsert", "schema.pattern":"public","auto.create"  :"true" ,"pk.mode"   :"record_key","pk.fields":"sensor_id",   "delete.enabled"  :"true","auto.evolve":"true", "transforms": "MapToString", "transforms.MapToString.type": "org.apache.kafka.connect.transforms.Cast$Value", "transforms.MapToString.spec": "data:string"}'

注意:将spec中的data替换为实际的MAP字段名。

方案3:调整Redis源连接器输出的Schema

检查Redis Stream数据结构,确认是否存在不必要的MAP字段。可通过Redis源连接器的配置参数(如redis.stream.entry.parser)调整输出字段类型,避免生成MAP类型的Avro字段。

内容的提问来源于stack exchange,提问作者Naveenkumar seerangan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 08:04:44