You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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驱动

二、准备连接器与驱动

  1. Debezium SQL Server连接器:下载与KSQLDB兼容的2.4.x版本连接器,解压后放到本地./debezium-connector-sqlserver目录
  2. SQL Server JDBC驱动:下载mssql-jdbc-12.6.3.jre11.jar,放到本地./mssql-jdbc目录

三、启用Azure SQL Server CDC功能

  1. 启用数据库级CDC:
    -- 启用数据库CDC
    EXEC sys.sp_cdc_enable_db;
    
    -- 验证启用状态
    SELECT name, is_cdc_enabled FROM sys.databases WHERE name = '你的数据库名';
    
  2. 启用目标表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 = '你的表名';
    
  3. 确保SQL账号拥有db_owner角色,或具备CDC相关的CONTROL、ALTER、SELECT权限

四、部署Debezium CDC连接器

  1. 进入KSQL CLI:
    docker exec -it ksqldb-cli ksql http://ksqldb-server:8088
    
  2. 创建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'
    );
    
  3. 验证连接器状态:
    SHOW CONNECTORS;
    DESCRIBE CONNECTOR mssql_cdc_connector;
    
    成功后,CDC事件会发送到Kafka主题mssql-azure.dbo.你的表名

五、KSQLDB处理CDC流

  1. 创建流消费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'
    );
    
  2. 过滤转换事件(仅保留新增/更新操作):
    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

  1. 准备Event Hub连接器:下载Confluent Azure Event Hub连接器,放到./debezium-connector-sqlserver目录
  2. 创建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'
    );
    
  3. 验证Sink连接器:
    SHOW CONNECTORS;
    DESCRIBE CONNECTOR eventhub_sink;
    

七、整体验证

  1. 启动服务:docker-compose up -d
  2. 检查容器状态:docker-compose ps
  3. 在Azure SQL Server中插入/更新测试数据
  4. 在KSQL CLI中查询处理后的流:SELECT * FROM processed_cdc_stream EMIT CHANGES LIMIT 5;
  5. 查看Event Hub控制台,确认消息已送达

内容的提问来源于stack exchange,提问作者HPAmaris

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.21 06:29:54