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

如何通过Java集成Kafka Connect读取CSV文件并转发数据给消费者

基于Java对接Kafka Connect实现CSV文件读取推送Kafka方案

前置依赖

  • 已部署可用的Kafka集群、Kafka Connect服务(生产环境推荐使用分布式模式)
  • 已将社区开源的kafka-connect-spooldir源连接器插件放到Kafka Connect服务的plugin.path配置路径下,重启Connect服务生效

实现步骤

1. 构造CSV源连接器配置

以下是标准的连接器配置JSON,可根据实际业务调整参数:

{
  "name": "csv-source-connector",
  "config": {
    "connector.class": "com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector",
    "tasks.max": "1",
    "input.path": "/opt/kafka_connect/csv_input",
    "finished.path": "/opt/kafka_connect/csv_finished",
    "error.path": "/opt/kafka_connect/csv_error",
    "input.file.pattern": ".*\\.csv",
    "csv.first.row.as.header": "true",
    "topic.creation.enable": "true",
    "topic.prefix": "csv.",
    "topic.creation.default.replication.factor": "3",
    "topic.creation.default.partitions": "3",
    "key.converter": "org.apache.kafka.connect.storage.StringConverter",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false"
  }
}

配置说明:

  • topic.creation.enable设为true即可开启连接器自动创建Kafka主题的能力,无需手动预先创建
  • 最终生成的Kafka主题名为topic.prefix + CSV文件名(不含后缀),例如CSV文件名为user_info.csv,自动创建的主题名为csv.user_info

2. Java代码提交连接器任务

通过调用Kafka Connect原生的REST API提交连接器即可,无需引入额外的Kafka Connect依赖,示例代码如下:

import okhttp3.*;
import java.io.IOException;

public class KafkaConnectCsvDemo {
    private static final String CONNECT_REST_URL = "http://<你的Kafka Connect服务地址>:8083/connectors";
    private static final MediaType JSON_MEDIA_TYPE = MediaType.get("application/json; charset=utf-8");

    public static void main(String[] args) throws IOException {
        // 替换为你自己的连接器配置JSON字符串
        String connectorConfig = "{\n" +
                "  \"name\": \"csv-source-connector\",\n" +
                "  \"config\": {\n" +
                "    \"connector.class\": \"com.github.jcustenborder.kafka.connect.spooldir.SpoolDirCsvSourceConnector\",\n" +
                "    \"tasks.max\": \"1\",\n" +
                "    \"input.path\": \"/opt/kafka_connect/csv_input\",\n" +
                "    \"finished.path\": \"/opt/kafka_connect/csv_finished\",\n" +
                "    \"error.path\": \"/opt/kafka_connect/csv_error\",\n" +
                "    \"input.file.pattern\": \".*\\\\.csv\",\n" +
                "    \"csv.first.row.as.header\": \"true\",\n" +
                "    \"topic.creation.enable\": \"true\",\n" +
                "    \"topic.prefix\": \"csv.\",\n" +
                "    \"topic.creation.default.replication.factor\": \"3\",\n" +
                "    \"topic.creation.default.partitions\": \"3\",\n" +
                "    \"key.converter\": \"org.apache.kafka.connect.storage.StringConverter\",\n" +
                "    \"value.converter\": \"org.apache.kafka.connect.json.JsonConverter\",\n" +
                "    \"value.converter.schemas.enable\": \"false\"\n" +
                "  }\n" +
                "}";
        
        OkHttpClient client = new OkHttpClient();
        RequestBody body = RequestBody.create(connectorConfig, JSON_MEDIA_TYPE);
        Request request = new Request.Builder()
                .url(CONNECT_REST_URL)
                .post(body)
                .build();
        
        try (Response response = client.newCall(request).execute()) {
            if (response.isSuccessful()) {
                System.out.println("CSV源连接器提交成功");
            } else {
                System.out.println("连接器提交失败,错误信息:" + response.body().string());
            }
        }
    }
}

3. 功能验证

  • 把需要处理的CSV文件放到你配置的input.path对应的目录下
  • 执行kafka-topics.sh --list --bootstrap-server <Kafka集群地址>命令,确认对应主题已自动创建
  • 执行kafka-console-consumer.sh --topic <对应主题名> --bootstrap-server <Kafka集群地址> --from-beginning命令,可查看到CSV文件每行转换后的JSON格式消息

常见调整项

  • 若CSV文件分隔符不是逗号,可新增csv.separator配置指定分隔符,比如csv.separator: "|"指定竖线分隔
  • 若需要自定义消息的key,可新增key.schema相关配置指定CSV某一列作为消息key
  • 若需要控制CSV文件的读取频率,可新增poll.interval.ms配置调整轮询间隔,默认是1000ms

内容的提问来源于stack exchange,提问作者ShadowClan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 20:18:00