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
相关产品推荐
相关产品推荐

