如何用Confluent SFTP CSV Source Connector生成嵌套JSON写入Kafka
调整方案
你只需要修改连接器配置里的value.schema字段,定义嵌套的STRUCT结构即可实现要求的嵌套JSON输出,具体修改步骤如下:
1. 调整value.schema结构
修改后的格式化value.schema如下:
{ "name": "com.example.users.User", "type": "STRUCT", "isOptional": false, "fieldSchemas": { "name": { "isOptional": false, "type": "STRING" }, "address": { "name": "com.example.users.Address", "type": "STRUCT", "isOptional": false, "fieldSchemas": { "street": { "isOptional": false, "type": "STRING" }, "number": { "isOptional": false, "type": "STRING" } } } } }
2. 替换配置中的value.schema字符串
将上面的结构压缩为转义后的JSON字符串,替换原连接器配置里的value.schema值即可,转义后的字符串如下:
"{\"name\":\"com.example.users.User\",\"type\":\"STRUCT\",\"isOptional\":false,\"fieldSchemas\":{\"name\":{\"isOptional\":false,\"type\":\"STRING\"},\"address\":{\"name\":\"com.example.users.Address\",\"type\":\"STRUCT\",\"isOptional\":false,\"fieldSchemas\":{\"street\":{\"isOptional\":false,\"type\":\"STRING\"},\"number\":{\"isOptional\":false,\"type\":\"STRING\"}}}}}"
3. 配置校验注意事项
- 保持
csv.first.row.as.header为false、csv.skip.lines为1的配置不变,schema字段的顺序和CSV列顺序严格对应:第一列对应顶级name字段,第二列对应address.street,第三列对应address.number,和你现有映射逻辑一致 - 保持
value.converter.schemas.enable为false,输出的JSON就不会携带schema信息,完全符合你要求的结构
内容的提问来源于stack exchange,提问作者Alberto Pires
相关产品推荐
相关产品推荐

