ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Kappa架构与Kafka实战:从核心原理到集群搭建与调优

Kappa架构与Kafka实战:从核心原理到集群搭建与调优 Kappa架构这几年又重新被频繁提起很多团队在从Lambda架构往这边迁移。我自己的感受是Kappa架构之所以能成为流处理领域的“屠龙刀”本质上是因为Kafka这个组件把“存储”和“计算”彻底解耦了。这篇文章我想从Kappa架构的核心思想出发结合Kafka的底层原理把集群搭建、参数调优、消息顺序性和延迟排查这些实操内容一次性讲透适合正在做流处理架构选型、或者想用Kafka重构数据处理链路的读者参考。1. 为什么Kappa架构值得被称为“屠龙刀”1.1 Lambda架构的痛点两套代码的维护噩梦在聊Kappa之前绕不开Lambda架构。Lambda架构把数据处理分为批处理层、速度层和服务层批处理层用离线计算引擎处理全量数据速度层用流处理引擎处理实时增量数据最后在服务层把两边的结果合并。听起来很完美实际上落地后问题非常现实——同一套业务逻辑需要写两遍一遍给离线引擎一遍给流处理引擎。我有段时间维护过这种架构最痛苦的不是写代码而是改需求。今天产品说指标口径要调批处理改完等调度任务跑完流处理还要单独发布一次两边结果稍微对不上就得查半天是哪个环节出了偏差。Lambda的问题还在于存储和计算被硬生生拆成了两条链路。批处理层存储的是历史快照速度层存储的是增量窗口服务层要同时读两个存储系统做合并如果两个系统的一致性模型不同数据对不齐就是家常便饭。更麻烦的是离线任务和实时任务的计算逻辑很难做到真正的同一个口径数据延迟和准确性的争论会持续整个项目周期。1.2 Kappa架构的本质一切皆流重放即批处理Kappa架构的核心思想其实非常朴素——既然实时处理链路天然就能消费数据为什么还需要一套独立的批处理链路只要底层存储能把数据完整保留下来任何时刻想“重新计算一遍历史数据”不需要去跑离线任务只需要启动一个新的流作业从Kafka的指定offset开始重新消费一遍就可以了。这就是它区别于Lambda的关键不区分批和流全部当作流来处理离线计算只是流的“重放”。这个思想的成立有一个硬性前提消息中间件必须支持长时间的数据保留而且重放的成本要足够低。Kafka的设计恰好满足了这一点。Kafka的存储是分段的顺序日志消费进度由offset记录消费者只要seek到任意offset就能从那个位置把数据重新读一遍。这意味着历史计算不再依赖外部存储的快照而是直接把Kafka当作一个可回放的分布式日志系统。1.3 Kafka在Kappa中的角色不只是消息队列在大多数人的认知里Kafka是一个消息队列负责把数据从生产端搬到消费端。但在Kappa架构里Kafka的角色要重要得多——它同时承担了存储层和事件总线两层职责。Kappa架构的全部可靠性都构建在Kafka的消息持久化之上如果Kafka丢了数据或性能出现瓶颈整个架构的数据基础就会塌方。这种设计带来的好处是Kappa架构的数据链路只有一个核心组件。数据进Kafka之后所有计算任务都从同一个地方消费数据同一个topic既可以被实时作业消费用于在线计算也可以被一个刚启动的重放作业消费用于回溯计算。结果是存储只用一套、接入只用一套、计算的代码逻辑也只需要维护一份。相比Lambda架构运维复杂度和代码维护量都有很大幅度的下降。2. 锻造刀胚Kafka的核心机制决定了Kappa的上限2.1 日志即存储Kafka的数据结构设计Kafka之所以能做到长时间保留海量数据底层依赖的是分段日志结构。每个topic的一个分区在磁盘上对应一个目录目录下是一组segment文件每个segment的大小默认1GB包含一个.log文件和一个.index文件。消息写入时不断追加到当前活跃segment的尾部当segment写满后Kafka就创建一个新的segment开始继续写。这种设计有几个好处首先追加写入的顺序IO性能非常高每秒几十万的写入吞吐就是这么来的其次segment文件写满后是不可变的Kafka可以基于segment做定时清理比如按保留时间删除过期的segment或者按总大小删除最早的segment最后消费者读数据时只需要做顺序扫描和少量二分查找效率远高于随机读写。Kappa架构要求Kafka能长时间保留数据这个保留能力本质上就是segment文件的滚动和清理机制配合的结果。2.2 offset与消费者组重放能力的技术底座Kafka的消费者组是支撑Kappa架构“重放”能力的关键机制。每个消费者组对于同一个分区维护着一个当前的消费位置offset通过提交offset来记录“我已经消费到哪一条了”。当你需要重新计算历史数据只需要用一个新的group.id启动消费者它就能从最早可用的offset开始消费或者用现有group的消费者实例通过assign和seek方法手动指定要消费的分区和offset位置。举一个实际场景。假设业务方在一个月后发现计算逻辑有bug需要把所有历史数据重新计算一遍。在Lambda架构下你要跑一个月前的全量离线数据耗时很长。在Kappa架构下只要topic的数据保留期覆盖了需要回溯的范围内直接新起一个流处理作业从这个topic的最早offset开始消费重算结果写入一个新的输出topic修正后再切换流量即可。整个过程就像录好的视频重新播放一遍而Kafka的日志就是那盘录像带。2.3 分区与顺序性只能在分区内保证有序消息顺序性是流处理里一个绕不开的话题。Kafka的顺序保证说的是“同一个分区内的消息是有序的”而不是全局有序。生产者在发送消息时可以指定keyKafka会对key进行哈希相同key的消息进入同一个分区因此同一个key的消息天然保持顺序。在Kappa架构中顺序性通常意味着你需要选对消息的key。比如处理用户行为日志时应该把userId作为消息key这样同一个用户的所有事件就会进入同一个分区流处理作业在消费时就不会出现同一个用户的事件被不同分区并行消费导致事件乱序的问题。如果一个topic不指定key而是随机轮询写入分区那消费端的顺序就无法保证了。消费端同样有顺序性的要求。Kafka允许单个消费者实例启用多个线程消费但是多个线程同时处理同一分区的消息时完成顺序是不可控的。如果业务对顺序敏感有两种方案最简单的方案是把消费线程数设为1让一个分区由单线程串行处理复杂的方案是分区内按key做缓冲区排序当然这会增加很大的代码复杂度。我建议普通场景直接单线程消费分区Kafka单线程消费一个分区的吞吐量其实已经很高远没有到成为瓶颈的程度。2.4 高吞吐与扩展性分区就是并发的天花板Kafka的高吞吐来自分区并行。一个topic被拆成多个分区分区分摊到集群的多台broker上生产者的写入和消费者的读取都可以在不同分区上并行。消费者组里可以有多个消费者实例每个实例负责一部分分区当一个实例挂掉时其他实例会抢占它负责的分区完成故障转移。在Kappa架构里分区的数量决定了整个流处理管道的并行度上限。选择分区数量时需要同时考虑目标吞吐量、单分区吞吐能力和consumer实例数量。一个粗略的估算方式是目标吞吐除以单消费者实例的吞吐能力就是所需的最小分区数。比如每秒需要处理10万条消息单实例每秒能处理1万条那么至少要10个分区。另外还要考虑未来流量增长建议按峰值流量的2倍预留分区数因为分区一旦创建扩容涉及数据再平衡成本不小。3. 实操从零搭建Kafka集群并落地Kappa架构3.1 环境准备与集群规划搭建Kafka集群之前先规划集群规模。Kafka本身对硬件的要求不算极端但存储是重点——因为Kappa架构下Kafka要长时间保留全量数据磁盘空间需要按“峰值流入速率乘以数据保留时长”来估算。假设每天流入100GB数据需要保留7天那至少准备1TB以上的可用存储空间还要考虑副本因子占用的额外空间。我这里给出一个通用的3节点集群配置参考适合中等规模业务使用节点角色配置建议broker-1Kafka Broker Controller4核8G1TB SSDbroker-2Kafka Broker4核8G1TB SSDbroker-3Kafka Broker4核8G1TB SSD仲裁节点KRaft Controller2核4G100G SSD注意从Kafka 2.8开始社区推荐用KRaft模式替代ZooKeeper3.5之后ZooKeeper被标记为弃用4.0已经移除了。新项目建议直接上KRaft模式省掉一套ZooKeeper集群的运维负担。3.2 Docker Compose快速搭建Kafka集群如果你不想在物理机上手工配置用Docker Compose可以在几分钟内拉起一套KRaft模式的Kafka集群。下面这个配置我实测过可以直接保存为docker-compose.yml使用version: 3.8 services: kafka1: image: bitnami/kafka:3.7 container_name: kafka1 ports: - 19092:9092 environment: - KAFKA_ENABLE_KRAFTyes - KAFKA_CFG_PROCESS_ROLESbroker,controller - KAFKA_CFG_CONTROLLER_LISTENER_NAMESCONTROLLER - KAFKA_CFG_INTER_BROKER_LISTENER_NAMEINTERNAL - KAFKA_CFG_LISTENERSINTERNAL://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERSINTERNAL://kafka1:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS1kafka1:9093,2kafka2:9093,3kafka3:9093 - KAFKA_KRAFT_CLUSTER_IDabcdefghijklmnopqrstuv - KAFKA_CFG_NODE_ID1 volumes: - kafka1_data:/bitnami/kafka kafka2: image: bitnami/kafka:3.7 container_name: kafka2 ports: - 19093:9092 environment: - KAFKA_ENABLE_KRAFTyes - KAFKA_CFG_PROCESS_ROLESbroker,controller - KAFKA_CFG_CONTROLLER_LISTENER_NAMESCONTROLLER - KAFKA_CFG_INTER_BROKER_LISTENER_NAMEINTERNAL - KAFKA_CFG_LISTENERSINTERNAL://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERSINTERNAL://kafka2:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS1kafka1:9093,2kafka2:9093,3kafka3:9093 - KAFKA_KRAFT_CLUSTER_IDabcdefghijklmnopqrstuv - KAFKA_CFG_NODE_ID2 volumes: - kafka2_data:/bitnami/kafka kafka3: image: bitnami/kafka:3.7 container_name: kafka3 ports: - 19094:9092 environment: - KAFKA_ENABLE_KRAFTyes - KAFKA_CFG_PROCESS_ROLESbroker,controller - KAFKA_CFG_CONTROLLER_LISTENER_NAMESCONTROLLER - KAFKA_CFG_INTER_BROKER_LISTENER_NAMEINTERNAL - KAFKA_CFG_LISTENERSINTERNAL://:9092,CONTROLLER://:9093 - KAFKA_CFG_ADVERTISED_LISTENERSINTERNAL://kafka3:9092 - KAFKA_CFG_LISTENER_SECURITY_PROTOCOL_MAPCONTROLLER:PLAINTEXT,INTERNAL:PLAINTEXT - KAFKA_CFG_CONTROLLER_QUORUM_VOTERS1kafka1:9093,2kafka2:9093,3kafka3:9093 - KAFKA_KRAFT_CLUSTER_IDabcdefghijklmnopqrstuv - KAFKA_CFG_NODE_ID3 volumes: - kafka3_data:/bitnami/kafka volumes: kafka1_data: kafka2_data: kafka3_data:启动集群后进入任意一个容器验证一下状态docker exec -it kafka1 bash kafka-topics.sh --bootstrap-server localhost:9092 --create --topic test-events --partitions 12 --replication-factor 2 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic test-events如果能看到每个分区有2个副本且状态正常集群就算搭好了。3.3 关键参数配置保留时间、副本因子与ack模型Kafka集群的运行效果很大程度取决于几个关键参数Kappa架构下更要重点关注。首先是数据保留策略这是Kappa架构能否成型的核心参数。默认的log.retention.hours是168小时7天如果你的业务需要重放一个月的数据就要调大这个值。注意调整保留时间会增加磁盘占用一定要先算清楚磁盘容量再改。kafka-configs.sh --bootstrap-server localhost:9092 --entity-type topics --entity-name test-events --alter --add-config retention.ms2592000000这条命令把test-events的保留时间改成了30天。300GB的分区数据就算30天不清理磁盘吃紧时会自动清除最早的segment不需要人工干预。其次是副本因子。Kafka允许生产者从配置为acks0、acks1或acksall生产端还需要配合min.insync.replicas参数来保证数据的持久性。如果设置acksall但min.insync.replicas1其实在副本都挂到只剩一个时仍然会写入成功并没有达到预期的可靠性。一般建议min.insync.replicas设为2配合副本因子3这样最多容忍一个broker宕机而不丢消息。生产端的幂等性也建议开启。enable.idempotencetrue可以让生产者对消息做去重配合acksall解决网络重试导致的重复消息问题。幂等开启后生产者的初始化ProducerId和序列号会带来少量开销但对数据链路来说这是值得的成本。3.4 用Kafka Streams实现Kappa架构的核心链路集群搭好了下面用Kafka Streams写一个Kappa架构的流处理作业演示数据进入Kafka后如何完成实时计算和重放计算。Kafka Streams是Kafka自带的流处理库最大的好处是它和Kafka深度集成分区机制、消费者组机制、状态存储都是现成的不需要额外部署流处理框架。假设业务场景是统计每个用户的点击量写入results topic。代码结构如下Properties props new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, click-analytics-app); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, localhost:19092,localhost:19093,localhost:19094); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.STATE_DIR_CONFIG, /tmp/kafka-streams); props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, earliest); StreamsBuilder builder new StreamsBuilder(); KStreamString, String source builder.stream(raw-events); KTableString, Long counts source .flatMapValues(value - Arrays.asList(value.split(\\s))) .groupBy((key, word) - word) .count(Materialized.as(CountsStore)); counts.toStream().to(results, Produced.with(Serdes.String(), Serdes.Long())); KafkaStreams streams new KafkaStreams(builder.build(), props); streams.start();重点说几个设计细节第一application.id决定了这个作业的消费者组。当你要做重放计算时换一个新的application.id启动同一个代码它就会从topic维护的最早offset开始消费从而触发全量重算。这个步骤就是Kappa架构实现“批处理”效果的关键。第二状态存储的含义要理解清楚。上面的groupBy操作会创建一个内部状态存储CountsStoreKafka Streams把它物化到本地磁盘并同步一份到内部的changelog topic。由于有了这个状态存储即使某个实例挂了重启它也能从changelog恢复状态不会丢失计数数据。第三同一个topic被多个不同application.id消费时各个作业的offset互相独立互不影响。这样可以把pipeline拆解为“原始数据层、实时计算层、重放计算层”多个计算作业并行消费原始topic在物理上做到互不干扰。3.5 数据流向与Topic设计在Kappa架构里Topic的规划直接决定了整个数据管道的清晰度。建议采用“分层topic”的方式组织数据原始数据层raw-events保存所有最原始的事件数据保留时间最长禁止业务方直接消费做自定义计算。明细计算层processed-events由流处理作业从原始层计算后写入保留时间可以缩短到3天左右供下游在线查询。结果聚合层aggregated-results统计汇总后的结果可以持久化到数据库供前端展示。这样做的优势是分工明确。原始层保证数据的完整可追溯体重放计算只需基于原始层重跑明细层降低了下游消费大数据量的压力结果层提供低延迟查询能力。每个层级的保留策略可以独立配置Kafka的磁盘空间也能更合理地利用。4. 实战中的坑消息延迟、顺序性与运维排查4.1 消息延迟高的排查路径消息延迟高是Kafka使用中最高频的问题。首先要区分是生产端延迟还是消费端延迟。生产端延迟的常见原因包括broker写入瓶颈、acks等待时间过长、批量参数设置不合理。消费端延迟则需要看消费者组当前的lag值。查看消费组滞后程度的方法kafka-consumer-groups.sh --bootstrap-server localhost:19092 --describe --group click-analytics-app输出的LAG列如果持续累积且不回落说明消费速度低于生产速度。排查顺序我一般这样走第一检查消费者的单条消息处理耗时。如果处理逻辑里包含外部IO数据库写入、RPC调用大概率瓶颈在这里。优化方案有批量写入、异步化处理、增加消费者实例数量。第二检查分区数量是否充足。消费者组的最大并行度等于订阅topic的分区总数如果分区数少于消费者实例数多出来的消费者其实处于空闲状态。这种场景下增加分区数就是最直接的扩容方式。第三检查网络和磁盘指标。broker的磁盘IO利用率超过80%或者网络带宽打满都会明显影响吞吐。用JMX监控或者Prometheus采集Broker的指标可以做到提前预警。4.2 消息顺序性被破坏的几种情况消息乱序是Kappa架构里最隐蔽的数据正确性问题。我总结了一下实际中容易踩的坑场景一生产者重试机制导致乱序。默认情况下生产者发送消息到同一个分区如果第一条失败后重试第二条已经发送成功那么失败的那条重试成功后会排在后面打破了原本的顺序。解决办法是设置max.in.flight.requests.per.connection1或者启用幂等性。幂等性开启后Kafka在broker端会对同一个ProducerId的消息做序号校验从而在允许in-flight并发的情况下依然保证顺序。场景二消费端多线程处理导致乱序。消费端如果开多个线程来处理同一个分区的数据线程A和线程B各自执行完成时间不可控后续聚合时就可能出现顺序颠倒。解决办法是单线程处理分区消息或者采用分区内按key排序的方式。场景三重放作业和实时作业的合并顺序。当Kappa架构中重放计算和实时计算同时运行时两个作业会分别写结果数据到达输出topic的顺序会乱七八糟下游如果直接读输出topic可能看到旧数据比新数据晚到达。我的处理方式是重放作业写一个临时结果topic等重放完成并验证数据一致性后再通过一个切换操作把消费流量切换到新topic。4.3 消费者再均衡导致的问题消费者再均衡是Kafka消费端另一个常见问题的来源。当消费者组里的实例数量变化、或者订阅的topic分区数变化时Kafka会触发一次rebalance把分区重新分配。在rebalance期间整个消费组会停止消费如果这个过程频繁发生就会看到消费停顿和延迟飙升。再均衡有两种触发因素一种是主动的比如有新的消费者实例加入另一种是被动的比如某个消费者实例在session.timeout.ms时间内没有发送心跳被判定为下线。对于偶发的rebalance影响不大但如果频繁触发要重点检查消费者实例的心跳线程是否被长任务阻塞了。如果消息处理时间过长可以考虑调大max.poll.interval.ms让Kafka不要在长时间处理时误判消费者已死。另一种更隐蔽的情况是消费者处理消息过程中发生异常导致offset一直没有提交消费者组不断重新消费同一批消息输出topic中就会出现大量重复数据。对策是处理好“处理”和“提交offset”之间的时机关系同时在下游做幂等写入。4.4 Kafka可视化运维工具怎么选底层的命令行工具适合排查问题但日常运维、查看topic消息内容、监控消费组lag有个可视化工具会高效得多。我用过几款简单对比下工具特点适用场景Kafdrop轻量级Web UI能查看topic、分区、消息内容快速查看消息内容开发调试Kafka UI功能全面支持消费者组管理、消息查询、CMAK迁移日常运维和监控Kafka Manager/CMAK老牌工具偏集群管理老集群监控、分区分配调整Offset Explorer桌面客户端跨平台个人电脑连测试环境快速排查如果只是想在本地快速看一下某个topic有没有消息、内容是什么我推荐KafdropDocker部署几分钟就能用。如果是生产环境长期运维用Kafka UI更合适它的消费者组页面能直接看到每个分区的lag和历史趋势。5. 大消息、不同语言客户端与性能边界5.1 大消息如何处理1MB限制与调优策略Kafka默认单条消息的最大大小是1MB由broker端的message.max.bytes和topic端的max.message.bytes共同控制。实际业务中如果出现接近或超过1MB的消息最常见的是埋点批量上报、日志聚合包这种场景。如果你确实需要传输更大的消息可以调大三个地方的参数topic的max.message.bytes、broker的message.max.bytes、消费者的fetch.max.bytes。但我的建议是尽量先在客户端做拆分。一个10MB的事件报文切成多段Kafka端到端延迟和数据可靠性都会好很多。大消息在Kafka里有两个副作用一是会占用较大的网络和内存缓冲影响整体吞吐二是消费者拉取时单个大消息会成为分区消费的阻塞点其他消息都要等它处理完。除非业务确实绕不开否则不建议走大消息路线。5.2 多语言客户端接入要点Kafka的生态里Java客户端是官方维护最完善、功能最丰富的。但如果你的技术栈是非Java比如在Qt应用里通过C接入Kafka或者用Python写流处理脚本客户端选择需要注意以下几点。C/Qt场景我实际用过librdkafka它是Kafka官方推荐的C库性能接近Java客户端。Qt接入librdkafka的方法是先编译出动态库然后用C封装producer和consumer。需要注意在MinGW编译环境下librdkafka需要依赖OpenSSL和zlib的MinGW版本否则编译时容易出现链接错误。另一个容易踩坑的地方是librdkafka的配置项名字和Java客户端不完全一致比如Java里的bootstrap.servers在librdkafka中写作bootstrap.servers但部分配置需要通过rd_kafka_conf_set接口设置字符串类型的值。Python端更简单confluent-kafka-python或者kafka-python都能快速上手。生产环境我推荐confluent-kafka-python它底层也是librdkafka性能和功能都更完整但需要注意它不支持Windows环境的pip安装编译Windows下建议用WSL或者Docker容器来跑Python消费者。5.3 Kafka面试题背后的核心知识网络Kafka面试题里出现频率最高的几个问题其实都是Kafka原理的核心节点我把它们串起来讲一下。第一个问题Kafka为什么快几个关键词回答顺序写磁盘、零拷贝、页缓存、批量消息压缩。顺序写利用磁盘的顺序IO性能零拷贝通过sendfile让数据从磁盘到网卡不经过用户态页缓存让热数据的读都在内存中完成。这些机制叠加Kafka的单分区吞吐可以达到每秒几十万条。第二个问题Kafka消息不丢失怎么做分三个端来看生产者端开启幂等acksall重试broker端设置足够的副本因子且min.insync.replicas2消费者端关掉自动提交改为业务处理成功后再手动提交offset。第三个问题如何保证消息顺序生产者端key相同进同分区单分区内由Kafka保证顺序消费者端单线程消费在重试场景下开启幂等保证in-flight重试不破坏顺序。这些问题本质上考的是对Kafka分布式日志模型、副本协议、消费者组机制的综合理解。看参考答案只能应付面试真正把这些知识串起来用在自己的架构里才是Kappa架构落地时最有价值的能力。6. 我的落地体会与扩展建议踩过不少坑之后我个人的落地体会是Kappa架构真正拉开差距的地方不在流处理引擎本身而是在Kafka的数据治理功底上。Kafka的topic规划、保留策略、key设计、消费组管理体系决定了Kappa架构能支撑多大的业务复杂度。很多团队一开始用Kafka只当消息管道进了Kappa架构之后才意识到Kafka实际上已经成了整个数据平台的主存储这个认知转变直接决定了架构的生命力。另外一个非常实用的建议是不要一开始就把所有计算都迁到Kappa架构上来。我通常的改造路径是先选一条数据链路做试点把离线仓库作为结果存储把实时计算和重放计算都跑在Kafka Streams上验证数据和运维都稳定后再逐步切新链路。这样风险可控也便于团队积累经验。后续这块内容还可以扩展的方向一是Kappa架构中的状态管理怎么做结合Kafka Streams的状态存储和changelog实现精确一次语义二是Kappa和实时数仓的结合比如用Paimon或者Iceberg在流计算结果上做增量数据湖进一步扩大可分析的数据范围三是Kafka分区数扩展和存储均衡的自动化工具建设这是Kappa架构长期运行后必然要面对的问题。每个方向都能单独写一篇深度的实践文章等我把这轮的踩坑记录整理完再来聊。
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进