如何解决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事件,同时添加删除标记字段。
配置步骤:
- 调整Debezium MySQL源连接器的关键参数:
delete.handling.mode=rewrite:将DELETE事件转为op="u"的UPDATE事件,避免被Sink过滤include.deleted.fields=*:让DELETE事件包含被删除记录的所有字段,防止更新时覆盖原有非主键字段
- 添加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 }
- 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事件。
步骤:
- 创建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;
- 创建ReplacingMergeTree类型的目标表:
CREATE TABLE your_target_table ( id Int64, name String, is_deleted UInt8 DEFAULT 0 ) ENGINE = ReplacingMergeTree(is_deleted) ORDER BY id;
- 创建物化视图自动处理不同类型的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
相关产品推荐
相关产品推荐

