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

如何在CnosDB中实现变更数据捕获?并通过Kafka增量同步至ClickHouse

实现CnosDB的CDC及通过Kafka增量同步到ClickHouse的方案

一、CnosDB中实现变更数据捕获(CDC)

CnosDB通过内置的数据订阅功能实现CDC,核心是将表的变更操作(插入、覆盖更新等时序场景常见操作)推送至Kafka,具体步骤如下:

  1. 启用订阅功能
    修改CnosDB配置文件cnosdb.toml,开启订阅模块:

    [subscription]
    enabled = true
    worker_count = 4 # 可根据并发量调整订阅线程数
    

    重启CnosDB使配置生效。

  2. 创建数据订阅规则
    用SQL创建针对目标表的订阅,指定将变更数据发送到Kafka的指定Topic:

    CREATE SUBSCRIPTION cnosdb_cdc_sub 
    ON your_iot_db.yot_data_table 
    WITH (
        TYPE = 'kafka',
        BROKERS = 'kafka_node1:9092,kafka_node2:9092',
        TOPIC = 'cnosdb_iot_cdc_topic',
        FORMAT = 'json' # 指定JSON格式,方便ClickHouse解析
    );
    

    注:CnosDB作为时序数据库,目前主要支持插入和覆盖更新操作的CDC捕获,部分版本暂不支持删除操作的捕获,需结合自身使用版本确认。

  3. 验证订阅有效性
    可通过查看Kafka Topic的消息内容,或CnosDB的运行日志,确认变更数据是否正常推送。

二、通过Kafka实现CnosDB到ClickHouse的增量同步

基于CnosDB推送的CDC数据,利用ClickHouse的Kafka引擎实现自动增量同步,步骤如下:

  1. 创建ClickHouse目标分析表
    构建与CnosDB源表匹配的结构,同时保留CDC操作标识:

    CREATE TABLE iot_exception_analysis (
        timestamp DateTime64(3),
        device_id String,
        metric Float64,
        is_exception Bool,
        op_type String COMMENT '操作类型:INSERT/UPDATE',
        PRIMARY KEY (device_id, timestamp)
    ) ENGINE = MergeTree()
    ORDER BY (device_id, timestamp);
    
  2. 创建Kafka引擎表作为消费入口
    该表作为ClickHouse与Kafka的桥梁,负责消费CnosDB推送的CDC数据:

    CREATE TABLE cnosdb_kafka_consumer (
        timestamp DateTime64(3),
        device_id String,
        metric Float64,
        is_exception Bool,
        op_type String
    ) ENGINE = Kafka(
        'kafka_node1:9092,kafka_node2:9092',
        'cnosdb_iot_cdc_topic',
        'clickhouse_iot_consumer_group',
        'JSONEachRow'
    ) SETTINGS 
        kafka_max_block_size = 1048576,
        kafka_skip_broken_messages = 10; # 跳过损坏消息,避免消费中断
    
  3. 创建物化视图实现自动同步
    物化视图会持续消费Kafka引擎表的数据,同时可直接过滤出异常数据:

    CREATE MATERIALIZED VIEW iot_sync_view TO iot_exception_analysis AS
    SELECT
        timestamp,
        device_id,
        metric,
        is_exception,
        op_type
    FROM cnosdb_kafka_consumer
    WHERE is_exception = true; # 仅同步异常数据,按需调整过滤条件
    
  4. 全量初始化与增量衔接

    • 先执行一次全量同步,将CnosDB已有数据导入ClickHouse:
      INSERT INTO iot_exception_analysis
      SELECT timestamp, device_id, metric, is_exception, 'INSERT' as op_type
      FROM cnosdb.yot_data_table; # 通过ClickHouse的CnosDB数据源连接查询,或导出导入
      
    • 全量完成后,确认物化视图已开始消费Kafka增量数据,实现全量+增量的完整同步。

关键注意事项

  • 数据一致性:通过ClickHouse系统表system.kafka_consumer监控消费偏移量,避免重复消费或数据丢失。
  • 性能调优:根据数据量调整Kafka引擎表的kafka_max_block_size、物化视图刷新间隔等参数。
  • 版本兼容性:确认使用的CnosDB版本支持数据订阅功能,ClickHouse版本兼容Kafka引擎配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.01 00:31:02