Ubuntu本地搭建Kafka-Connect适配Debezium与WSO2流处理器求助
在Ubuntu 18.04上独立安装Kafka Connect并配置Debezium MongoDB连接器
别担心,我一步步帮你搞定独立版Kafka Connect的安装,以及后续的Debezium MongoDB连接器配置,完全不用依赖Confluent套件。
一、准备工作:确认依赖环境
首先确保你的系统满足基础要求:
- 已经部署好可用的Kafka和ZooKeeper集群(你已经搞定这部分了,很好)
- 安装Java 8或更高版本:Ubuntu 18.04可以直接用OpenJDK 8,执行命令:
安装完成后用sudo apt update && sudo apt install openjdk-8-jdk -yjava -version验证版本是否符合要求。
二、安装Apache Kafka Connect(独立模式)
Kafka Connect是Apache Kafka的核心组件之一,不需要单独下载,直接用官方Kafka二进制包即可:
- 下载与你现有Kafka版本匹配的Apache Kafka二进制包(注意版本兼容性:Debezium 2.x系列建议搭配Kafka 2.8.x及以上版本)
- 解压到本地目录,比如:
tar -xzf kafka_2.13-3.5.1.tgz cd kafka_2.13-3.5.1 - 配置独立模式的Connect参数:
复制默认配置文件并修改关键项:
编辑cp config/connect-standalone.properties config/my-connect-standalone.propertiesmy-connect-standalone.properties,重点修改以下内容:
提前创建插件目录:# 指向你的Kafka Broker地址,本地部署一般是localhost:9092 bootstrap.servers=localhost:9092 # 使用JSON转换器,新手友好,不需要额外schema管理 key.converter=org.apache.kafka.connect.json.JsonConverter value.converter=org.apache.kafka.connect.json.JsonConverter # 关闭schema自动生成,简化配置 key.converter.schemas.enable=false value.converter.schemas.enable=false # 指定偏移量存储文件路径 offset.storage.file.filename=/tmp/connect.offsets # 插件目录,后续用来放Debezium连接器 plugin.path=/opt/kafka-connect-pluginssudo mkdir -p /opt/kafka-connect-plugins sudo chown $USER:$USER /opt/kafka-connect-plugins - 启动Kafka Connect独立模式:
启动后查看控制台日志,如果没有报错,说明Connect服务正常运行(默认监听8083端口)。bin/connect-standalone.sh config/my-connect-standalone.properties
三、安装Debezium MongoDB连接器
- 下载对应版本的Debezium MongoDB连接器插件包(版本要和你的Kafka Connect版本匹配,比如Debezium 2.4.0对应Kafka 2.8.x-3.5.x)
- 解压插件到之前配置的
plugin.path目录:unzip debezium-connector-mongodb-2.4.0.Final-plugin.zip -d /opt/kafka-connect-plugins/debezium-mongodb - 重启Kafka Connect服务,让它加载新的插件:
先停止之前的Connect进程(Ctrl+C),再重新执行启动命令即可。 - 验证插件是否加载成功:
用curl调用Connect的REST接口:
如果返回的列表中包含curl http://localhost:8083/connector-pluginsio.debezium.connector.mongodb.MongoDbConnector,说明插件加载成功。
四、配置并启动MongoDB连接器
- 首先确保你的MongoDB已经开启副本集模式(Debezium需要通过oplog捕获数据变更,单节点MongoDB必须转为副本集才能工作)
- 创建MongoDB连接器配置文件(比如
mongodb-connector.json),示例内容如下:
替换{ "name": "mongodb-connector", "config": { "connector.class": "io.debezium.connector.mongodb.MongoDbConnector", "mongodb.hosts": "rs0/localhost:27017", "mongodb.name": "mongodb-source", "collection.include.list": "your_database.your_collection", "tasks.max": "1" } }your_database.your_collection为你需要捕获变更的数据库和集合名称。 - 提交连接器配置到Kafka Connect:
curl -X POST -H "Content-Type: application/json" --data @mongodb-connector.json http://localhost:8083/connectors - 检查连接器运行状态:
如果状态显示curl http://localhost:8083/connectors/mongodb-connector/statusRUNNING,说明连接器已经正常工作,开始捕获MongoDB的数据变更并发送到Kafka主题。
一些关键注意事项
- 版本兼容性:一定要确保Debezium版本、Kafka版本、MongoDB版本三者兼容,避免出现奇怪的报错
- 权限:Kafka Connect运行用户需要对插件目录、偏移量文件拥有读写权限
- MongoDB副本集:如果之前是单节点MongoDB,需要先初始化副本集才能让Debezium正常捕获变更
内容的提问来源于stack exchange,提问作者Rahul Anand
相关产品推荐
相关产品推荐

