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

如何将多Kafka Topic日志合并至ClickHouse表并实现日志互富化?

在高负载场景下用ClickHouse实现日志数据互富化(基于Kafka Engine)

针对你的需求,可以通过Kafka引擎表 + 事务状态表 + 物化视图的组合方案,实现日志数据的实时互富化,同时保留所有原始日志行,满足30k RPS的高负载场景要求。

核心思路

  1. 用Kafka引擎表直接接收两个Topic的原始日志数据;
  2. 维护一个ReplacingMergeTree类型的事务状态表,存储每个TransactionalId对应的最新完整字段值(CustomerId、SenderIp);
  3. 通过物化视图将Kafka的原始数据同步到状态表,自动更新事务的最新状态;
  4. 再通过物化视图将原始日志与状态表关联,补全缺失字段后写入最终存储表,保留所有原始日志行。

具体实现步骤

1. 创建Kafka引擎接收表

分别为两个Topic创建Kafka引擎表,用于实时消费日志数据:

-- 接收inbound topic的Kafka表
CREATE TABLE log_inbound_kafka
(
    Event String,
    TransactionalId UInt64,
    CustomerId UInt32,
    SenderIp String
)
ENGINE = Kafka
SETTINGS kafka_broker_list = 'kafka-broker:9092',
       kafka_topic_list = 'inbound-topic',
       kafka_group_name = 'clickhouse-log-group',
       kafka_format = 'JSONEachRow';

-- 接收IP topic的Kafka表
CREATE TABLE log_ip_kafka
(
    Event String,
    TransactionalId UInt64,
    CustomerId UInt32,
    SenderIp String
)
ENGINE = Kafka
SETTINGS kafka_broker_list = 'kafka-broker:9092',
       kafka_topic_list = 'ip-topic',
       kafka_group_name = 'clickhouse-log-group',
       kafka_format = 'JSONEachRow';

2. 创建事务状态表

用ReplacingMergeTree存储每个事务的最新完整状态,按TransactionalId作为主键,通过版本字段确保最新数据被保留:

CREATE TABLE transaction_state
(
    TransactionalId UInt64,
    CustomerId UInt32,
    SenderIp String,
    _version UInt64 DEFAULT now() -- 用当前时间戳作为版本标识
)
ENGINE = ReplacingMergeTree(_version)
ORDER BY TransactionalId
TTL TransactionalId + INTERVAL 7 DAY; -- 可选:自动清理7天前的旧事务数据

3. 创建物化视图同步状态表

将两个Kafka表的有效数据同步到状态表,更新对应字段:

-- 同步inbound数据(仅更新非空的CustomerId)
CREATE MATERIALIZED VIEW mv_inbound_to_state TO transaction_state AS
SELECT
    TransactionalId,
    CustomerId,
    SenderIp,
    now() AS _version
FROM log_inbound_kafka
WHERE CustomerId IS NOT NULL;

-- 同步IP数据(仅更新非空的SenderIp)
CREATE MATERIALIZED VIEW mv_ip_to_state TO transaction_state AS
SELECT
    TransactionalId,
    CustomerId,
    SenderIp,
    now() AS _version
FROM log_ip_kafka
WHERE SenderIp IS NOT NULL;

4. 创建富化日志表及物化视图

将原始日志与状态表关联,补全缺失字段后写入最终存储表,保留所有原始日志行:

-- 最终富化日志存储表
CREATE TABLE enriched_logs
(
    Event String,
    TransactionalId UInt64,
    CustomerId UInt32,
    SenderIp String
)
ENGINE = MergeTree
ORDER BY (TransactionalId, Event);

-- 富化inbound日志的物化视图
CREATE MATERIALIZED VIEW mv_enrich_inbound TO enriched_logs AS
SELECT
    l.Event,
    l.TransactionalId,
    COALESCE(l.CustomerId, s.CustomerId) AS CustomerId, -- 优先取原始值,无则用状态表值
    COALESCE(l.SenderIp, s.SenderIp) AS SenderIp
FROM log_inbound_kafka l
LEFT JOIN transaction_state s ON l.TransactionalId = s.TransactionalId;

-- 富化IP日志的物化视图
CREATE MATERIALIZED VIEW mv_enrich_ip TO enriched_logs AS
SELECT
    l.Event,
    l.TransactionalId,
    COALESCE(l.CustomerId, s.CustomerId) AS CustomerId,
    COALESCE(l.SenderIp, s.SenderIp) AS SenderIp
FROM log_ip_kafka l
LEFT JOIN transaction_state s ON l.TransactionalId = s.TransactionalId;

高负载场景优化建议

  • Kafka消费配置:调整kafka_max_block_size=10485760(10MB)、kafka_flush_interval_ms=100,匹配30k RPS的吞吐量,减少频繁提交带来的开销;
  • 状态表性能:ReplacingMergeTree的合并操作后台异步执行,不会阻塞数据写入;可根据业务需求调整TTL,避免状态表过度膨胀;
  • 关联性能:transaction_state表的主键是TransactionalId,等值关联时会利用主键索引加速查询,确保高负载下的关联效率;
  • 资源隔离:为Kafka消费、状态表更新、日志富化分配独立的ClickHouse节点或资源池,避免互相影响。

方案验证

当你的示例数据进入Kafka后:

  1. log_inbound_kafka收到第一条数据,通过mv_inbound_to_state更新transaction_state中对应TransactionalId的CustomerId;
  2. log_ip_kafka收到第二条数据,通过mv_ip_to_state更新transaction_state中对应TransactionalId的SenderIp;
  3. 两个物化视图分别将原始日志与状态表关联,补全缺失字段后写入enriched_logs,最终得到你需要的富化结果,且保留了所有原始日志行。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 22:47:13