Kafka Connect JDBC源连接器连接Microsoft SQL Server故障排查
嘿,看了你的情况,问题一目了然——你的分布式Docker化Kafka Connect集群里还没安装Kafka JDBC源连接器插件!从你贴的connector-plugins API返回结果就能看出来,目前加载的只有这些Sink类型的插件:
[ { "class": "io.confluent.connect.elasticsearch.ElasticsearchSinkConnector", "type": "sink", "version": "4.0.0" }, { "class": "io.confluent.connect.hdfs.HdfsSinkConnector", "type": "sink", "version": "4.0.0" }, { "class": "io.confluent.connect.hdfs.tools.SchemaSourceConnector", "type": "source", "version": "4.0.0" } ]
(注:你贴的内容截断了,我补全了常见的HDFS工具类插件)
接下来,针对Docker化的分布式集群,我给你捋清楚怎么一步步解决:
1. 先选对JDBC连接器包
你需要用Confluent的Kafka JDBC连接器(包含源和Sink实现),注意要和你现有插件的版本(4.0.0)兼容,最好直接选同版本的包,避免版本冲突。另外,别忘了还要准备Microsoft SQL Server的JDBC驱动包(比如mssql-jdbc-*.jar),这个是连接SQL Server必须的。
2. 在Docker环境里安装插件
因为是Docker化的集群,有两种靠谱的方式:
方式一:构建自定义Connect镜像
基于官方的Confluent Connect镜像,把JDBC连接器包和SQL Server驱动复制到镜像的插件目录里。写个简单的Dockerfile就行:FROM confluentinc/cp-kafka-connect:4.0.0 # 用confluent-hub安装JDBC连接器(自动下载对应版本) RUN confluent-hub install --no-prompt confluentinc/kafka-connect-jdbc:4.0.0 # 把本地的SQL Server驱动复制到插件目录 COPY mssql-jdbc-9.4.1.jre8.jar /usr/share/java/kafka-connect-jdbc/构建完这个新镜像后,替换掉原有Connect集群里所有节点的镜像,保证所有节点插件一致。
方式二:用挂载卷加载插件
如果不想重新构建镜像,也可以把JDBC连接器包和SQL Server驱动放到宿主机的某个目录,然后挂载到每个Connect容器的插件目录下。比如启动容器时加挂载参数:docker run -d \ --name connect-node-1 \ -v /your/local/path/kafka-connect-jdbc:/usr/share/java/kafka-connect-jdbc \ -v /your/local/path/mssql-jdbc.jar:/usr/share/java/kafka-connect-jdbc/mssql-jdbc.jar \ # 这里加上你原来的其他参数,比如Kafka地址、Connect的配置等 confluentinc/cp-kafka-connect:4.0.0注意:分布式集群的所有Connect节点都要这么配置,不然会出现节点插件不一致的问题。
3. 重启Connect集群
不管用哪种方式,都需要重启所有Connect节点,让新插件被加载。重启后,再调用/connector-plugins API,你就能看到io.confluent.connect.jdbc.JdbcSourceConnector的条目了。
4. 配置JDBC源连接器
插件加载好后,就可以创建连接器的配置了,给你个针对SQL Server的示例配置,你根据自己的实际情况改:
{ "name": "mssql-jdbc-source-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSourceConnector", "connection.url": "jdbc:sqlserver://<你的SQL Server地址>:1433;databaseName=<数据库名>;user=<用户名>;password=<密码>", "mode": "incrementing", // 增量同步模式,适合有自增ID的表 "incrementing.column.name": "id", // 替换成你表里的自增字段 "topic.prefix": "mssql-", // Kafka主题前缀,最终主题名是前缀+表名 "poll.interval.ms": 5000, // 每5秒轮询一次数据库 "table.whitelist": "你的表名", // 指定要同步的表,多个用逗号分隔 "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter.schemas.enable": "false", "value.converter.schemas.enable": "false" } }
最后提醒一句:SQL Server的JDBC驱动必须和JDBC连接器放在同一个目录下,不然Connect会找不到驱动类,同步就会失败哦!
内容的提问来源于stack exchange,提问作者gunj_desai

