知识卡片
数据缓存重用层:利用Kafka可重复消费特性,避免同一份数据被实时和离线各自重复清洗
内容
数据经过实时清洗(ETL:解压、解密、转义、补全、异常处理)之后,接下来实时计算和离线计算这两条完全不同的处理路径,都需要用到这份清洗过的数据——如果各自独立地从原始数据重新做一遍清洗,不仅浪费计算资源,还可能因为清洗逻辑的实现差异导致两条路径处理出来的数据口径不一致。某音乐公司的解法是引入一层”数据缓存重用层”:把经过实时清洗之后的数据重新写回Kafka集群、保留一定周期,这样同一份已经清洗好的数据可以被下游两条完全独立的路径各自消费——离线计算(批处理)通过自研的KG-Camus组件定时批量拉取到HDFS,实时计算则基于Storm/JStorm直接从Kafka消费(借助现成的storm-kafka组件)。这个设计能够成立的关键前提是Kafka原生支持消息的重复消费——同一份数据留在Kafka里,不会因为被实时计算路径消费过一次就”消失”,离线计算路径依然能够独立地把它完整读取出来,两条路径互不干扰、也不需要各自重复执行一遍清洗逻辑。这个案例给出了一条处理”同一份数据需要被多条独立下游路径消费”这类场景的通用思路:与其在每条下游路径各自重复处理原始数据(既浪费计算资源,又容易导致口径不一致),不如把”处理一次、共享给多方复用”作为架构设计的显式目标,专门在处理完成和多方消费之间插入一层可以被重复读取的缓冲(这里是支持重复消费的消息队列),让”清洗”这个动作真正只发生一次,后续所有下游路径都从这个统一、干净的结果出发,而不是各自从零开始。
参考来源
- 位置:《高可用架构(第1卷)》第6章《大数据与数据库》"6.1 某音乐公司的大数据实践"节,"6.1.2 某音乐公司大数据技术架构"(源文件:_epub-src/OEBPS/Text/Chapter6_1_3.xhtml)
- 结论依据:原文说明"决定把经过数据实时清洗后的数据重新写入Kafka并保留一定周期,离线计算(批处理)通过KG-Camus拉到HDFS……实时计算基于Storm/JStorm直接从Kafka消费,有很完美的解决方案——storm-kafka组件",直接支撑本卡片结论。
- 原始内容:决定把经过数据实时清洗后的数据重新写入Kafka并保留一定周期,离线计算(批处理)通过KG-Camus拉到HDFS(通过作业调度系统配置相应的作业计划),实时计算基于Storm/JStorm直接从Kafka消费,有很完美的解决方案——storm-kafka组件。