如何在Docker Compose中配置KSQLDB连接Azure SQL Server并对接Event Hub
连接Azure SQL Server与KSQLDB实现CDC流式处理并发送至Event Hub
我找不到SQL Server接入KSQLDB的最新文档,想了解最简方案:连接Azure SQL Server与KSQLDB,实现CDC表的流式处理并发送至Event Hub。以下是我的测试用docker-compose.yml,需要分步配置指导,包括在compose中添加SQL Server连接器,以及实现从KSQLDB到Event Hub的消息发送。
version: '3.9' # Updated version to ensure compatibility services: zookeeper: image: confluentinc/cp-zookeeper:7.4.0 hostname: zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:7.4.0 hostname: kafka container_name: kafka depends_on: - zookeeper ports: - "29092:29092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 ksql-server: image: confluentinc/ksqldb-server:0.29.0 hostname: ksqldb-server container_name: ksqldb-server ports: - "8088:8088" environment: KSQL_LISTENERS: http://0.0.0.0:8088 KSQL_BOOTSTRAP_SERVERS: kafka:9092 KSQL_CONNECT_BOOTSTRAP_SERVERS: kafka:9092 KSQL_CONNECT_PLUGIN_PATH: "/usr/share/kafka/plugins" KSQL_CONNECT_URL: "jdbc:sqlserver://xxxxx.database.windows.net:4093;database=xxxxx;encrypt=true;trustServerCertificate=false;hostNameInCertificate=*.database.windows.net;ApplicationIntent=readonly" KSQL_CONNECT_USER: "xxxx" KSQL_CONNECT_PASSWORD: "xxxxxx" KSQL_KSQL_INTERNAL_TOPIC_REPLICATION_FACTOR: 2 volumes: - "C:/Program\ Files/sqljdbc_12.6.3.0_enu/sqljdbc_12.6/enu/jars:/usr/share/kafka/plugins/mssql" depends_on: - kafka # KSQL server depends on the Kafka broker - zookeeper # KSQL server also depends on Zookeeper - maven # Add a Maven service dependency maven: image: maven:3.9.8-eclipse-temurin entrypoint: ["sh", "-c", "mvn dependency:copy -DremoteRepositories=https://repo1.maven.org/ -DgroupId=com.microsoft.sqlserver -DartifactId=mssql-jdbc -Dversion=12.6.3.0 -Ddest=/tmp/mssql-jdbc/"] volumes: - "C:/Program\ Files/sqljdbc_12.6.3.0_enu/sqljdbc_12.6/enu/jars:/usr/share/kafka/plugins/mssql" ksqldb-cli: image: confluentinc/ksqldb-cli:0.29.0 container_name: ksqldb-cli depends_on: - kafka - ksql-server entrypoint: /bin/sh tty: true
分步配置指导
一、修正docker-compose.yml配置
当前配置存在参数错误、冗余服务等问题,以下是修正后的完整配置:
version: '3.9' services: zookeeper: image: confluentinc/cp-zookeeper:7.4.0 hostname: zookeeper container_name: zookeeper ports: - "2181:2181" environment: ZOOKEEPER_CLIENT_PORT: 2181 ZOOKEEPER_TICK_TIME: 2000 kafka: image: confluentinc/cp-kafka:7.4.0 hostname: kafka container_name: kafka depends_on: - zookeeper ports: - "29092:29092" environment: KAFKA_BROKER_ID: 1 KAFKA_ZOOKEEPER_CONNECT: 'zookeeper:2181' KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: PLAINTEXT:PLAINTEXT,PLAINTEXT_HOST:PLAINTEXT KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092,PLAINTEXT_HOST://localhost:29092 KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 ksqldb-server: image: confluentinc/ksqldb-server:0.29.0 hostname: ksqldb-server container_name: ksqldb-server ports: - "8088:8088" environment: KSQL_LISTENERS: http://0.0.0.0:8088 KSQL_BOOTSTRAP_SERVERS: kafka:9092 KSQL_KSQL_INTERNAL_TOPIC_REPLICATION_FACTOR: 1 # 单节点Kafka必须设为1,避免副本数不足报错 KSQL_CONNECT_BOOTSTRAP_SERVERS: kafka:9092 KSQL_CONNECT_GROUP_ID: "ksqldb-connect-group" KSQL_CONNECT_CONFIG_STORAGE_TOPIC: "_ksqldb-connect-configs" KSQL_CONNECT_OFFSET_STORAGE_TOPIC: "_ksqldb-connect-offsets" KSQL_CONNECT_STATUS_STORAGE_TOPIC: "_ksqldb-connect-statuses" KSQL_CONNECT_CONFIG_STORAGE_REPLICATION_FACTOR: 1 KSQL_CONNECT_OFFSET_STORAGE_REPLICATION_FACTOR: 1 KSQL_CONNECT_STATUS_STORAGE_REPLICATION_FACTOR: 1 KSQL_CONNECT_PLUGIN_PATH: "/usr/share/kafka/plugins" volumes: # 挂载Debezium SQL Server连接器和JDBC驱动 - "./debezium-connector-sqlserver:/usr/share/kafka/plugins/debezium" - "./mssql-jdbc:/usr/share/kafka/plugins/mssql" depends_on: - kafka - zookeeper ksqldb-cli: image: confluentinc/ksqldb-cli:0.29.0 container_name: ksqldb-cli depends_on: - ksql-server entrypoint: /bin/sh tty: true
配置说明
- 移除冗余的
maven服务,改为手动下载插件更可靠 - 修正KSQL内部主题副本数为1(单节点Kafka环境强制要求)
- 添加Kafka Connect必要的集群配置(组ID、存储主题等)
- 调整插件挂载路径,用于放置Debezium连接器和SQL Server JDBC驱动
二、准备连接器与驱动
- Debezium SQL Server连接器:下载与KSQLDB兼容的2.4.x版本连接器,解压后放到本地
./debezium-connector-sqlserver目录 - SQL Server JDBC驱动:下载
mssql-jdbc-12.6.3.jre11.jar,放到本地./mssql-jdbc目录
三、启用Azure SQL Server CDC功能
- 启用数据库级CDC:
-- 启用数据库CDC EXEC sys.sp_cdc_enable_db; -- 验证启用状态 SELECT name, is_cdc_enabled FROM sys.databases WHERE name = '你的数据库名'; - 启用目标表CDC:
-- 对指定表启用CDC EXEC sys.sp_cdc_enable_table @source_schema = N'dbo', @source_name = N'你的表名', @role_name = NULL, @supports_net_changes = 1; -- 验证表CDC状态 SELECT name, is_tracked_by_cdc FROM sys.tables WHERE name = '你的表名'; - 确保SQL账号拥有
db_owner角色,或具备CDC相关的CONTROL、ALTER、SELECT权限
四、部署Debezium CDC连接器
- 进入KSQL CLI:
docker exec -it ksqldb-cli ksql http://ksqldb-server:8088 - 创建CDC连接器:
CREATE SOURCE CONNECTOR mssql_cdc_connector WITH ( 'connector.class' = 'io.debezium.connector.sqlserver.SqlServerConnector', 'tasks.max' = '1', 'database.hostname' = 'xxxxx.database.windows.net', 'database.port' = '1433', -- Azure SQL默认端口为1433,需确认实际端口 'database.user' = 'xxxx', 'database.password' = 'xxxxxx', 'database.dbname' = 'xxxxx', 'database.server.name' = 'mssql-azure', -- 自定义服务器标识,作为Kafka主题前缀 'database.history.kafka.bootstrap.servers' = 'kafka:9092', 'database.history.kafka.topic' = 'dbhistory.mssql-azure', 'include.schema.changes' = 'false', 'table.include.list' = 'dbo.你的表名', -- 指定要捕获CDC的表 'database.encrypt' = 'true' ); - 验证连接器状态:
成功后,CDC事件会发送到Kafka主题SHOW CONNECTORS; DESCRIBE CONNECTOR mssql_cdc_connector;mssql-azure.dbo.你的表名
五、KSQLDB处理CDC流
- 创建流消费CDC主题:
CREATE STREAM customer_cdc_stream ( before STRUCT<id INT, name STRING, email STRING>, after STRUCT<id INT, name STRING, email STRING>, op STRING ) WITH ( KAFKA_TOPIC = 'mssql-azure.dbo.你的表名', VALUE_FORMAT = 'JSON', KEY_FORMAT = 'KAFKA' ); - 过滤转换事件(仅保留新增/更新操作):
CREATE STREAM processed_cdc_stream AS SELECT after->id AS customer_id, after->name AS customer_name, after->email AS customer_email, op AS operation_type FROM customer_cdc_stream WHERE op IN ('I', 'U'); -- I=插入,U=更新,D=删除
六、发送至Azure Event Hub
- 准备Event Hub连接器:下载Confluent Azure Event Hub连接器,放到
./debezium-connector-sqlserver目录 - 创建Sink连接器:
CREATE SINK CONNECTOR eventhub_sink WITH ( 'connector.class' = 'io.confluent.connect.azure.eventhubs.EventHubSinkConnector', 'tasks.max' = '1', 'topics' = 'processed_cdc_stream', -- KSQL处理后的流对应的Kafka主题 'bootstrap.servers' = '<你的Event Hub命名空间>.servicebus.windows.net:9093', 'sasl.mechanism' = 'PLAIN', 'security.protocol' = 'SASL_SSL', 'sasl.jaas.config' = 'org.apache.kafka.common.security.plain.PlainLoginModule required username="$ConnectionString" password="你的Event Hub连接字符串";', 'key.converter' = 'org.apache.kafka.connect.storage.StringConverter', 'value.converter' = 'org.apache.kafka.connect.json.JsonConverter', 'value.converter.schemas.enable' = 'false' ); - 验证Sink连接器:
SHOW CONNECTORS; DESCRIBE CONNECTOR eventhub_sink;
七、整体验证
- 启动服务:
docker-compose up -d - 检查容器状态:
docker-compose ps - 在Azure SQL Server中插入/更新测试数据
- 在KSQL CLI中查询处理后的流:
SELECT * FROM processed_cdc_stream EMIT CHANGES LIMIT 5; - 查看Event Hub控制台,确认消息已送达
内容的提问来源于stack exchange,提问作者HPAmaris
相关产品推荐
相关产品推荐

