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

如何配置Apache Kafka File Sink Connector实现JSON转竖线分隔串及每小时文件轮转?

Apache Kafka File Sink Connector 配置实现方案

需求概述

需配置Kafka File Sink Connector实现两个核心功能:

  1. 将Kafka主题中的JSON格式消息转换为竖线分隔字符串输出到文件
    • 输入示例:
      {"key1":"value1","key2":"value2","key3":"value3","key4":"value4","key5":"value5","key6":"value6"}
      
    • 预期输出:
      value1|value2|value3|value4|value5|value6
      
  2. 按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:

开发步骤

  1. 创建Java类,实现org.apache.kafka.connect.transforms.Transformation<ConnectRecord>接口
  2. 在apply方法中完成JSON解析与格式转换:
    • 用Jackson/Gson等库解析消息为Map结构
    • 按配置的字段顺序提取值并拼接成竖线分隔字符串
    • 返回包含转换后内容的新ConnectRecord
  3. 通过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");
    }
}

集成步骤

  1. 将代码打包为JAR文件(需包含依赖库如Jackson)
  2. 将JAR放入Kafka Connect的plugin.path指定目录
  3. 在连接器配置中启用自定义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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 07:33:19