如何用Apache Kafka开源JDBC连接器实现MSSQL表间Kafka Connect同步
替换Confluent JDBC连接器为Apache Kafka开源版本的配置方法
要替换为Apache Kafka官方开源的JDBC连接器,只需修改配置中的connector.class字段,以下是调整后的完整配置:
源连接器(Source Connector)配置
将原Confluent的源连接器类io.confluent.connect.jdbc.JdbcSourceConnector替换为Apache Kafka开源版本的org.apache.kafka.connect.jdbc.JdbcSourceConnector,修改后的curl命令如下:
curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{ "name": "jdbc_source_mssql_01", "config": { "connector.class": "org.apache.kafka.connect.jdbc.JdbcSourceConnector", "connection.url": "jdbc:sqlserver://fulfillmentdbhost:1433;databaseName=fulfillmentdb", "connection.user": "fullfilment_user", "connection.password": "<password>", "topic.prefix": "order-status-update-", "mode":"timestamp", "table.whitelist" : "fulfullmentdb.status", "timestamp.column.name": "LAST_UPDATED", "validate.non.null": false } }'
注:这里同步修正了连接URL为MSSQL的标准格式(原配置写的是MySQL的URL,与你描述的MSSQL数据库匹配),如果你的实际环境有特殊配置可按需调整。
Sink连接器(Sink Connector)配置
同理,将Confluent的Sink连接器类io.confluent.connect.jdbc.JdbcSinkConnector替换为org.apache.kafka.connect.jdbc.JdbcSinkConnector,修改后的curl命令如下:
curl -X POST http://localhost:8083/connectors -H "Content-Type: application/json" -d '{ "name": "jdbc_sink_mssql_01", "config": { "connector.class": "org.apache.kafka.connect.jdbc.JdbcSinkConnector", "connection.url": "jdbc:sqlserver://crmdbhost:1433;databaseName=crmdb", "connection.user": "crm_user", "connection.password": "<password>", "topics": "order-status-update-status", "table.name.format" : "crmdb.order_status" } }'
额外说明
- 除
connector.class字段外,原Confluent JDBC连接器的大部分配置参数在Apache开源版本中均可兼容,无需额外修改。 - 需确保你的Kafka Connect环境中已部署Apache Kafka JDBC连接器的jar包,同时包含MSSQL对应的JDBC驱动(如
mssql-jdbc.jar)。
内容的提问来源于stack exchange,提问作者Shoaib Khan
相关产品推荐
相关产品推荐

