使用GPKAFKA向GreenPlum插入数据失败的问题解决咨询
问题现象
使用GPKafka将Kafka中的Avro/JSON格式数据插入GreenPlum目标表时,持续报错提示仅支持单个JSON列,即使外部表结构与目标表一致,调整消息结构后仍无法解决。
Avro格式数据导入时错误:
debug,rollback the batch 0 due to Failed to execute batch: pq: avro_import: only support single json column
JSON格式数据导入时错误:
debug,rollback the batch 0 due to Failed to execute batch: pq: json_import: only support single json/jsonb/gp_jsonb column
使用的GreenPlum版本:6.19.3
相关信息
Kafka Avro格式消息示例
{ "key1": "try81", "col1": { "int": 1 }, "col2": { "string": "def" }, "col3": { "string": "ghi" }, "col4": { "string": "jkl" }, "instance_id": { "string": "009" } }
目标表DDL
CREATE TABLE testgpkafka ( key1 CHARACTER VARYING(5) NOT NULL, col1 INTEGER, col2 CHARACTER VARYING(55), col3 CHARACTER VARYING(19), col4 CHARACTER VARYING(13), instance_id CHARACTER VARYING(15), PRIMARY KEY (key1) );
GPKafka配置文件(gpkafka.yaml)
DATABASE: gpdb_dev USER: -- PASSWORD: -- HOST: -- PORT: -- KAFKA: INPUT: SOURCE: BROKERS: 192.168.151.201:9092, 192.168.151.202:9092, 192.168.151.203:9092 TOPIC: mcp_kafka_net_21.mcp.testgpkafka PARTITIONS: (0) COLUMNS: FORMAT: avro AVRO_OPTION: SCHEMA_REGISTRY_ADDR: http://192.168.151.201:8081 OUTPUT: SCHEMA: smi TABLE: testgpkafka MODE: insert
已尝试的无效操作
- 使用扁平JSON格式消息,仍触发相同错误:
{ "key1": "tr159", "col1": 1, "col2": "def", "col3": "ghi", "col4": "jkl", "instance_id": "009" }
- 将消息包裹为单JSON列结构,依旧无效:
{ "data": { "key1": "tr150", "col1": 1, "col2": "def", "col3": "ghi", "col4": "jkl", "instance_id": "009" } }
解决方法
1. 补全GPKafka配置的列映射
当前配置中INPUT.COLUMNS为空,GPKafka默认把整个消息当作单个JSON列处理,必须明确指定Kafka消息字段与目标表列的映射关系。
针对Avro格式的修正配置示例:
KAFKA: INPUT: SOURCE: BROKERS: 192.168.151.201:9092, 192.168.151.202:9092, 192.168.151.203:9092 TOPIC: mcp_kafka_net_21.mcp.testgpkafka PARTITIONS: (0) COLUMNS: - NAME: key1 TYPE: varchar(5) - NAME: col1 TYPE: int - NAME: col2 TYPE: varchar(55) - NAME: col3 TYPE: varchar(19) - NAME: col4 TYPE: varchar(13) - NAME: instance_id TYPE: varchar(15) FORMAT: avro AVRO_OPTION: SCHEMA_REGISTRY_ADDR: http://192.168.151.201:8081 OUTPUT: SCHEMA: smi TABLE: testgpkafka MODE: insert
2. 修正Kafka Avro消息结构
现有Avro消息的字段为嵌套结构(如col1: {"int":1}),与目标表的扁平列结构不兼容,需调整Avro Schema为扁平字段,使消息中每个字段直接对应值,修正后的消息示例:
{ "key1": "try81", "col1": 1, "col2": "def", "col3": "ghi", "col4": "jkl", "instance_id": "009" }
3. 验证外部表结构
运行GPKafka生成外部表后,执行以下SQL确认外部表列与目标表完全匹配:
SELECT column_name, data_type FROM information_schema.columns WHERE table_schema = 'smi' AND table_name = 'testgpkafka_external';
确认无误后再执行插入操作。
4. 手动COPY验证(排查用)
若GPKafka仍报错,可先用COPY命令手动测试数据导入,确认数据格式兼容性:
COPY smi.testgpkafka FROM '/tmp/test_data.json' WITH FORMAT 'json';
如果手动导入成功,说明问题出在GPKafka的配置或消息映射上。
内容的提问来源于stack exchange,提问作者Maaheen

