如何配置Apache Kafka File Sink Connector实现JSON转竖线分隔串及每小时文件轮转?
Apache Kafka File Sink Connector 配置实现方案
需求概述
需配置Kafka File Sink Connector实现两个核心功能:
- 将Kafka主题中的JSON格式消息转换为竖线分隔字符串输出到文件
- 输入示例:
{"key1":"value1","key2":"value2","key3":"value3","key4":"value4","key5":"value5","key6":"value6"} - 预期输出:
value1|value2|value3|value4|value5|value6
- 输入示例:
- 按
output.txt_yyyyMMddHH格式每小时轮转输出文件
问题1:JSON转竖线分隔串的配置规则
可通过Kafka Connect内置的TemplateTransform实现无代码格式转换,核心配置如下:
# 配置JSON转换器解析消息内容 value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false # 启用Template转换并指定拼接规则 transforms=formatOutput transforms.formatOutput.type=org.apache.kafka.connect.transforms.TemplateTransform$Value # 按字段顺序定义竖线分隔模板,严格匹配输出顺序 transforms.formatOutput.template={value.key1}|{value.key2}|{value.key3}|{value.key4}|{value.key5}|{value.key6}
若存在字段缺失场景,可在模板中添加默认值避免异常,例如
{value.key1?:""}表示字段为空时输出空字符串。
问题2:每小时文件轮转的配置项
官方FileStreamSinkConnector不支持文件轮转,需替换为Confluent提供的FileRotateSinkConnector,核心配置如下:
# 替换为支持轮转的连接器类 connector.class=io.confluent.connect.file.FileRotateSinkConnector # 配置每小时轮转(单位:毫秒) rotate.interval.ms=3600000 # 指定轮转后的文件名格式(yyyyMMddHH对应年、月、日、24小时制小时) file.name.format=output.txt_yyyyMMddHH # 基础输出文件路径(轮转时自动添加时间后缀) file=/path/output.txt
若需基于固定时间点轮转(如整点),可使用
rotate.schedule.interval.ms配合rotate.schedule.timezone指定时区,确保轮转时间与业务时区一致。
问题3:自定义转换器的开发与集成
若内置Transform无法满足动态字段解析、复杂格式处理等需求,可开发自定义Transform:
开发步骤
- 创建Java类,实现
org.apache.kafka.connect.transforms.Transformation<ConnectRecord>接口 - 在
apply方法中完成JSON解析与格式转换:- 用Jackson/Gson等库解析消息为Map结构
- 按配置的字段顺序提取值并拼接成竖线分隔字符串
- 返回包含转换后内容的新ConnectRecord
- 通过
configure方法加载自定义配置(如字段顺序、分隔符)
示例核心代码:
import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.config.ConfigDef; import com.fasterxml.jackson.databind.ObjectMapper; import java.util.Map; public class JsonToPipeDelimitedTransform<R extends ConnectRecord<R>> implements Transformation<R> { private ObjectMapper objectMapper = new ObjectMapper(); private String[] fieldOrder; @Override public void configure(Map<String, ?> configs) { fieldOrder = ((String) configs.get("field.order")).split(","); } @Override public R apply(R record) { try { Map<String, String> jsonData = objectMapper.readValue(record.value().toString(), Map.class); StringBuilder sb = new StringBuilder(); for (int i = 0; i < fieldOrder.length; i++) { sb.append(jsonData.getOrDefault(fieldOrder[i], "")); if (i < fieldOrder.length - 1) sb.append("|"); } return record.newRecord(record.topic(), record.kafkaPartition(), record.keySchema(), record.key(), null, sb.toString(), record.timestamp()); } catch (Exception e) { throw new RuntimeException("Failed to convert JSON to pipe-delimited string", e); } } @Override public void close() {} @Override public ConfigDef config() { return new ConfigDef() .define("field.order", ConfigDef.Type.STRING, ConfigDef.Importance.HIGH, "Comma-separated list of field names in output order"); } }
集成步骤
- 将代码打包为JAR文件(需包含依赖库如Jackson)
- 将JAR放入Kafka Connect的
plugin.path指定目录 - 在连接器配置中启用自定义Transform:
transforms=customFormat transforms.customFormat.type=com.your.package.JsonToPipeDelimitedTransform transforms.customFormat.field.order=key1,key2,key3,key4,key5,key6
完整配置示例
结合以上配置,最终的连接器配置如下:
name=file-rotate-sink-connector connector.class=io.confluent.connect.file.FileRotateSinkConnector tasks.max=1 topics=my-topic # 文件轮转配置 file=/path/output.txt rotate.interval.ms=3600000 file.name.format=output.txt_yyyyMMddHH # JSON解析与格式转换配置 value.converter=org.apache.kafka.connect.json.JsonConverter value.converter.schemas.enable=false transforms=formatOutput transforms.formatOutput.type=org.apache.kafka.connect.transforms.TemplateTransform$Value transforms.formatOutput.template={value.key1}|{value.key2}|{value.key3}|{value.key4}|{value.key5}|{value.key6}
内容的提问来源于stack exchange,提问作者Vishnu
相关产品推荐
相关产品推荐

