如何通过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
相关产品推荐
相关产品推荐

