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

如何解决ClickHouse Kafka Connect Sink无法处理DELETE操作的问题?

解决方案:Kafka Connect CDC同步MySQL到ClickHouse处理DELETE操作

针对ClickHouse Kafka Connect Sink不支持DELETE操作的问题,提供以下三种可行方案:

方案一:修改Debezium源连接器,将DELETE转为带标记的UPDATE

核心思路是让Debezium把MySQL的DELETE事件转换为ClickHouse Sink可处理的UPDATE事件,同时添加删除标记字段。

配置步骤:

  1. 调整Debezium MySQL源连接器的关键参数:
    • delete.handling.mode=rewrite:将DELETE事件转为op="u"的UPDATE事件,避免被Sink过滤
    • include.deleted.fields=*:让DELETE事件包含被删除记录的所有字段,防止更新时覆盖原有非主键字段
  2. 添加InsertField转换,插入删除标记字段is_deleted=1

完整配置示例:

{
    "name": "mysql-source-connector",
    "connector.class": "io.debezium.connector.mysql.MySqlConnector",
    "database.hostname": "mysql-host",
    "database.port": "3306",
    "database.user": "cdc-user",
    "database.password": "cdc-pass",
    "database.server.id": "1",
    "database.server.name": "mysql-cdc-server",
    "database.include.list": "your_db",
    "table.include.list": "your_db.your_table",
    "delete.handling.mode": "rewrite",
    "include.deleted.fields": "*",
    "transforms": "AddDeletedFlag",
    "transforms.AddDeletedFlag.type": "org.apache.kafka.connect.transforms.InsertField$Value",
    "transforms.AddDeletedFlag.static.field": "is_deleted",
    "transforms.AddDeletedFlag.static.value": 1
}
  1. ClickHouse目标表需预先定义is_deleted UInt8 DEFAULT 0字段,建议使用ReplacingMergeTree引擎确保更新逻辑生效:
CREATE TABLE your_target_table
(
    id Int64,
    name String,
    -- 其他业务字段
    is_deleted UInt8 DEFAULT 0
)
ENGINE = ReplacingMergeTree(is_deleted)
ORDER BY id;

方案二:用Kafka Streams预处理DELETE事件

通过Kafka Streams过滤并转换Topic中的DELETE事件,将其转为UPDATE事件后写入新Topic,再由ClickHouse Sink消费新Topic。

核心逻辑示例(Java):

import org.apache.kafka.streams.KafkaStreams;
import org.apache.kafka.streams.StreamsBuilder;
import org.apache.kafka.streams.kstream.KStream;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.common.serialization.Serdes;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.node.JsonNodeFactory;
import com.fasterxml.jackson.databind.node.ObjectNode;
import java.util.Properties;

public class CdcDeleteProcessor {
    public static void main(String[] args) {
        StreamsBuilder builder = new StreamsBuilder();
        // 消费原始CDC Topic
        KStream<String, JsonNode> cdcStream = builder.stream("mysql-cdc-topic");

        // 处理DELETE事件:转为op="u"的UPDATE,添加is_deleted=1
        KStream<String, JsonNode> deleteStream = cdcStream
                .filter((key, value) -> "d".equals(value.get("op").asText()))
                .mapValues(value -> {
                    ObjectNode updatedValue = JsonNodeFactory.instance.objectNode();
                    // 保留主键和原有字段
                    updatedValue.put("id", value.get("id").asLong());
                    updatedValue.put("name", value.get("name").asText());
                    // 添加删除标记
                    updatedValue.put("is_deleted", 1);
                    // 修改op为u,让Sink正常处理
                    updatedValue.put("op", "u");
                    return updatedValue;
                });

        // 合并原始的INSERT/UPDATE事件和处理后的DELETE事件
        KStream<String, JsonNode> mergedStream = cdcStream
                .filter((key, value) -> !"d".equals(value.get("op").asText()))
                .merge(deleteStream);

        // 写入新Topic供ClickHouse Sink消费
        mergedStream.to("mysql-cdc-processed-topic");

        KafkaStreams streams = new KafkaStreams(builder.build(), getStreamsConfig());
        streams.start();
    }

    // 初始化Kafka Streams配置
    private static Properties getStreamsConfig() {
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "cdc-delete-processor");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "kafka-broker:9092");
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, JsonSerde.class);
        return props;
    }
}

方案三:使用ClickHouse原生Kafka引擎表

绕过Kafka Connect Sink,直接让ClickHouse通过Kafka引擎表消费CDC Topic,再通过物化视图处理DELETE事件。

步骤:

  1. 创建Kafka引擎表消费原始CDC Topic:
CREATE TABLE cdc_kafka_source
(
    id Int64,
    name String,
    op String,
    is_deleted UInt8 DEFAULT 0
)
ENGINE = Kafka
SETTINGS
    kafka_broker_list = 'kafka-broker:9092',
    kafka_topic_list = 'mysql-cdc-topic',
    kafka_group_name = 'clickhouse-cdc-consumer',
    kafka_format = 'JSONEachRow',
    kafka_skip_broken_messages = 1;
  1. 创建ReplacingMergeTree类型的目标表:
CREATE TABLE your_target_table
(
    id Int64,
    name String,
    is_deleted UInt8 DEFAULT 0
)
ENGINE = ReplacingMergeTree(is_deleted)
ORDER BY id;
  1. 创建物化视图自动处理不同类型的CDC事件:
CREATE MATERIALIZED VIEW mv_cdc_sync TO your_target_table AS
SELECT
    id,
    name,
    CASE WHEN op = 'd' THEN 1 ELSE is_deleted END AS is_deleted
FROM cdc_kafka_source
WHERE op IN ('c', 'u', 'd');

该物化视图会自动将INSERT/UPDATE事件同步到目标表,将DELETE事件转为标记is_deleted=1的记录,ReplacingMergeTree在合并时会保留最新的标记状态。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.24 04:31:01