ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Kafka+Flink实时数据管道实践:从架构设计到压测排障

Kafka+Flink实时数据管道实践:从架构设计到压测排障 1. 项目背景与核心需求拆解1.1 圣保罗业务场景的特殊性先说项目背景。我们做的是一个面向巴西圣保罗市场的实时业务系统核心场景是支付链路中的风险控制和用户行为实时分析。圣保罗是南美最大的金融科技与电商聚集地Pix即时支付普及率极高用户对“点击到结果返回”的耐心窗口极短。业务方提的需求从来不是“能不能做”而是“这个数据延迟能不能压到秒级以内”。在这个环境下做数据流处理管道最大的挑战不是技术本身而是三件事吞吐的突发性、链路的稳定性、以及数据准确性和延迟之间的取舍。圣保罗的网络基础设施和国内一线城市相比有差距跨区域链路的抖动和时间不同步问题非常普遍任何一个环节设计不当都会直接反映到业务指标上。我们的实时管道要解决的核心问题可以概括为每秒处理数十万级事件、端到端延迟控制在3秒以内、在流量洪峰下不丢数据不阻塞、并且支持业务方自助接入新数据源。这几个约束放在一起就决定了不能只靠某个开源组件“开箱即用”必须做大量的工程化改造。1.2 需求指标拆解与约束分析项目启动第一天我们做的第一件事不是画架构图而是把非功能性需求逐条量化。这里有一个经验如果业务方只说“要实时”“要能扩展”后面一定会扯皮必须把指标落到数字上。最终我们和业务方达成一致的关键指标如下指标项目标值说明平均端到端延迟≤ 3秒从事件产生到下游可消费峰值吞吐≥ 50万 events/s按业务大促峰值放大1.5倍数据不丢率100%精确一次语义重点链路不允许丢数服务可用性≥ 99.95%按季度统计新增数据源接入周期≤ 3天业务自助或半自助接入这些指标里最折磨人的是“精确一次语义 低延迟 高吞吐”这个组合。很多团队在初期会用Kafka作为缓冲层但Kafka的at-least-once语义配合下游幂等可以做到不丢不重。然而在实时风控场景里重复事件会导致用户被误判所以最终我们选择了Flink的端到端精确一次语义通过checkpoint 下游幂等写入来实现。另一个重要的约束是成本。圣保罗机房资源价格不低领导不可能让你无脑堆机器。所以我们的设计目标里还有一条隐性指标在满足前四项的前提下把资源成本控制在预算线以下。这直接影响技术选型——比如能用Kafka Streams解决的就不上Flink后面会详细讲这个决策过程。2. 技术选型与整体架构设计2.1 为什么选 Kafka Flink而不是别的组合当时摆在我们面前的备选方案有三组Kafka Flink、Kafka Streams 单体应用、Pulsar 自研处理框架。我逐个说下取舍逻辑。Kafka Streams 的优势是部署简单、学习成本低和Kafka同源天然支持只回溯和状态存储。但它的扩展性受限于应用内部的线程模型在状态很大的场景下比如用户会话窗口状态管理能力不如Flink成熟。另一个硬伤是流式SQL支持有限我们团队里有一半人是从批处理转过来的纯写DSL的接受度不高。Pulsar 自研这套方案当时我们评估过后放弃了主要原因是团队对Pulsar的运维经验不足而且Pulsar的IO隔离优势在这个场景下不是刚需。我们是业务技术团队不是中间件团队选最成熟、社区资料最多的方案永远是第一原则。最终选定 Kafka Flink 的理由一句话就能讲清楚生态最全出问题能搜到答案团队上手快性能满足需求成本可控。Kafka负责削峰填谷和消息持久化Flink负责有状态计算和时间窗口处理。2.2 分层架构设计与职责边界整个管道按职责拆成四层每层只做一件事边界非常清晰第一层是接入层统一接收业务方的原始事件数据。业务方通过HTTP或Kafka Producer SDK上报事件接入层做格式校验、协议转换和初步的数据清洗。这一层我们用了独立的Kafka集群作为缓冲区避免业务方直接和Flink耦合。第二层是处理层也就是Flink集群。这里跑着三个核心作业实时特征计算作业、风险规则匹配作业、用户行为聚合作业。三个作业之间通过内部Kafka topic解耦任何一个作业重启都不会影响其他作业。第三层是存储层包含Redis实时特征缓存、ClickHouse明细和指标存储、MySQL元数据管理以及对象存储原始数据归档。存储层是根据访问模式做了分级的热数据进Redis温数据进ClickHouse冷数据进对象存储。第四层是服务层为业务方提供实时查询API和实时特征服务。这一层不直接读Kafka而是从Redis或ClickHouse读取计算好的结果保证查询延迟在毫秒级。这个分层架构最大的好处是故障隔离。处理层出问题不会影响接入层收数存储层慢查询不会拖垮处理层。代价是多了一跳网络开销换来的是整体的稳定性我认为值得。2.3 环境部署与集群规格参考这里分享一套经过压测验证的部署规格供参考。我们是三个Flink作业共用一个集群通过YARN Session模式跑每个作业独立提交互不干扰。集群配置如下Kafka集群3个Broker每个Broker 16核32GB内存4块NVMe SSD操作系统页缓存拉满Flink集群5个TaskManager每个TaskManager 16核64GB内存单个TaskManager配4个SlotZooKeeper复用已有集群不做独立部署因为用的Kafka版本支持KRaft模式所以后来去掉了ZK依赖这里不展开ClickHouse2个节点每个节点8核32GB主要用于近实时查询实际上线后我们发现瓶颈几乎永远不在CPU和内存上而在网络带宽和磁盘IO。尤其是Kafka的磁盘吞吐高峰期能到800MB/sSSD选型上千万不要省钱。另外TaskManager的堆内存要留足Flink的状态存在堆外内存RocksDB堆内存只存数据流处理和算子状态这个比例要提前压测调好。3. 管道核心实现与关键代码3.1 数据接入层的SDK设计与使用接入层我们给业务方提供了一个轻量级SDK核心思路是“只要你会发HTTP请求就能接入实时管道”。SDK内部将事件批量打包异步发送到Kafka然后返回一个消息ID给业务方业务方凭消息ID可以查状态。SDK的核心配置类代码如下public class EventCollectorConfig { // 接入端点生产环境走内网域名 private String endpoint; // 批量发送大小默认512条或2MB触发发送 private int batchSize 512; private int batchSizeBytes 2 * 1024 * 1024; // 发送间隔默认500ms兜底触发 private int lingerMs 500; // 重试次数和退避策略 private int retries 3; private long retryBackoffMs 1000; // 本地缓冲队列大小防止业务线程阻塞 private int bufferQueueSize 10000; }这里有一个很关键的工程细节SDK绝对不能阻塞业务主线程。我们内部维护了一个有界队列业务方调用send()方法只是往队列里丢数据后台有一个Sender线程负责批量发送。当队列满了的时候说明消费速度跟不上SDK会先尝试丢弃非关键事件只有关键事件才触发阻塞式重传。这个设计在流量突刺时保住了主链路的可用性。接入层收到数据后会做三件事校验JSON格式、补全服务端时间戳、按事件类型分发到对应的Kafka topic。格式校验这一环必须做否则脏数据流进Flink会导致作业反序列化失败严重时会导致整个作业重启。我们线上就遇到过有人传了字段类型不一致的数据排障花了两个多小时。3.2 Flink 核心作业的骨架实现Flink作业的骨架我们统一做了一个模板工程新作业只需要关注业务逻辑其他交给模板。以实时特征计算作业为例核心代码结构如下public class FeatureComputeJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime); env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000); env.getCheckpointConfig().setCheckpointTimeout(120000); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1); DataStreamRawEvent source env.addSource(new FlinkKafkaConsumer( input-topic, new RawEventDeserializationSchema(), kafkaProps )).assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractorRawEvent(Time.seconds(5)) { Override public long extractTimestamp(RawEvent element) { return element.getEventTime(); } }); DataStreamFeatureValue featureStream source .keyBy(RawEvent::getUserId) .process(new FeatureProcessFunction()) .name(feature-compute-processor); // 把特征值写入Redis和内部Kafka featureStream.addSink(new RedisFeatureSink(redisConfig)); featureStream.addSink(new FlinkKafkaProducer( feature-output-topic, new FeatureValueSerializationSchema(), kafkaProps )); env.execute(real-time-feature-compute-job); } }这段代码里有几个值得说的点第一为什么用EventTime而不是ProcessingTime因为数据跨机房传输Producer和Consumer所在机器时间可能存在几十到几百毫秒的偏差直接用ProcessingTime会导致窗口计算结果不准确。用EventTime配合Watermark可以容忍一定程度的乱序。我们设置的是5秒的乱序容忍度意味着事件最多迟到5秒还能被正确处理超过5秒就被丢弃或走侧输出流单独处理。第二Checkpoint间隔为什么设60秒这个值是压测试出来的。checkpoint太频繁会导致状态后端压力大影响吞吐太稀疏则故障恢复时丢的数据多。60秒配合最少间隔30秒既保证恢复点在30秒内又不会频繁做快照。第三为什么不用Flink SQL我们评估过SQL在Window聚合、Join操作上确实开发效率高但在自定义特征计算比如多条件叠加、时间衰减因子上表达能力受限。对于风控特征这种逻辑经常要微调的场景Process Function的灵活性是必要的。但纯Aggregation的场景比如用户行为计数我们内部还是会用SQL。3.3 状态管理RocksDB 还是 Heap状态管理是Flink工程实践里最容易被忽视的坑。我们第一版用的是Heap状态后端因为配置简单、访问快。但跑了两个月后发现Heap状态在作业重启时需要从Kafka回溯数据来恢复恢复时间长达20分钟。业务方直接炸了你这不是实时系统这是半实时系统。后来切到RocksDB状态后端设计上就是增量checkpoint恢复时间降到2分钟以内。代价是访问速度比Heap慢但通过合理设计key实际影响可控。RocksDB的关键配置如下RocksDBStateBackend rocksDBStateBackend new RocksDBStateBackend(hdfs:///flink/checkpoints, true); rocksDBStateBackend.setNumberOfTransferThreads(3); rocksDBStateBackend.setPredefinedOptions(PredefinedOptions.SPINNING_DISK_OPTIMIZED_HIGH_MEM);注意我这里开了增量checkpoint开关构造器第二个参数传true这是恢复速度大幅提升的关键。另外如果状态数据量不大几GB级别可以用SPINNING_DISK_OPTIMIZED_HIGH_MEM预选项让RocksDB更多地使用内存缓存减少磁盘IO。状态量超过100GB时再考虑默认配置。还有一个经验RocksDB的block cache大小和Flink托管内存是分开的。我们通过taskmanager.memory.managed.fraction设置托管内存占比默认是0.4但如果你用RocksDB建议调到0.5到0.6之间给RocksDB更多空间否则状态访问的随机读会成为瓶颈。4. 高可扩展性设计与压测验证4.1 分区策略从 KeyBy 到自定义分区器高可扩展的第一个核心是分区策略。Kafka topic的分区数是吞吐上限的决定因素Flink算子的并行度受限于上游Kafka分区数。如果分区数不够消费并发上不去再怎么加机器都是白搭。我们刚开始为每个topic只设计了12个分区压测发现单分区消费速率大约每秒5000条全量也就6万条每秒离50万的目标差得远。后来把分区数提高到64仍然不够最终定为128个分区。但分区数不是越多越好。Kafka单分区是有序性保证的有状态操作比如按用户聚合要求同一个用户的记录进入同一个分区。我们这里按user_id做keyBy如果直接使用系统默认的key取模分区器在key分布不均时会导致数据倾斜。比如巴西市场有几个大客户单个客户的流量能占到总量的30%如果不做处理就算有128个分区热点区域还是忙不过来。这里我们做了自定义分区器核心思路是两级分区先用一致性哈希把userId映射到一个虚拟桶比如1024个桶再把桶映射到分区。这样即使单个userId流量很大也能打散到多个分区同时同一个用户的数据仍然有序进入同一个分区。public class UserIdPartitioner extends Partitioner { private static final int VIRTUAL_BUCKETS 1024; Override public int partition(String topic, Object key, byte[] keyBytes, Object value, byte[] valueBytes, Cluster cluster) { ListPartitionInfo partitions cluster.partitionsForTopic(topic); int numPartitions partitions.size(); if (key null) { // 没有key的走轮询 return ThreadLocalRandom.current().nextInt(numPartitions); } String userId key.toString(); int bucket Math.floorMod(userId.hashCode(), VIRTUAL_BUCKETS); // 桶到分区的映射关系通过外部配置控制可以动态调整 return BucketToPartitionMapping.getPartition(bucket, numPartitions); } }这套方案的另一个好处是支持分区数扩缩容。当我们把topic分区数从128扩到256时只需要更新桶到分区的映射关系不用全量重算userId的hash值。虽然在扩容瞬间会有部分数据乱序但窗口计算允许5秒乱序这个影响可以接受。4.2 弹性伸缩感知流量洪峰的自动扩缩容圣保罗的流量特点是“白天平稳晚上和周末飙升”而且经常有大促活动。人工扩缩容显然不现实因为从发现流量涨到扩容完成至少需要15分钟流量洪峰可能只持续10分钟。这期间管道会被背压压垮。我们的方案是两层弹性伸缩第一层是Kafka消费者即Flink作业的并行度调整。Flink作业在YARN Session下运行TaskManager数量可以动态调整配合外部监控系统当消费延迟超过阈值时自动申请新的TaskManager并重启作业让作业恢复后使用更多并行度继续消费。这里有个前提你的分区数要足够多否则并行度再高也只能最多分到和分区数一样的订阅者。第二层是Kafka生产端的缓冲策略调整。SDK里面我们实现了一个动态batch机制——当检测到Kafka响应时间变长或发送失败率升高时自动降低batchSize减少单次发送的数据量降低Broker压力当链路空闲时再增大batchSize提高吞吐。这个策略用了最简单的PID控制思路不需要太复杂。弹性伸缩这里要特别提醒一个坑永远不要把自动重启和人工运维混在一起。我们第一版自动扩缩容是直接调YARN接口的结果有一次作业状态后端损坏自动重启反复失败形成了“启动-崩溃-再启动”的死循环。后来加了一个熔断器机制连续重启失败3次就停止自动操作转人工这个坑才算填上。4.3 压测方法与实践数据说到压测很多人会犯一个错误用Kafka的producer性能作为整个管道的压测目标结果压测报告很好看上线就被打回原形。我们的经验是压测一定要跑全链路从业务方SDK发送到Flink处理结束写入下游存储全链路每个环节都要监控。我们用了自己的压测工具基于Go写的模拟500个客户端并发发送事件。压测等级分三档日常流量10万 events/s、大促流量30万 events/s、极限压测50万 events/s持续30分钟。压测中暴露出来的几个问题和解决记录第一个是Kafka Broker的页缓存命中率。刚开始压测到20万的时候Kafka的磁盘读IO开始飙升延迟从5毫秒涨到200毫秒。排查发现是因为事件带的时间戳和当前时间差距太大Kafka不会自动清理超过保留期的segment导致消费者频繁向磁盘请求冷数据。优化方案是调整log.segment.bytes和log.retention.hours把segment切得更小256MB保留时间设短24小时冷数据快速淘汰。第二个是Flink的Watermark在压测下的表现。默认的BoundedOutOfOrderness是每秒计算一次Watermark在高吞吐下这会造成大量事件被标记为迟到。压测时我们把Watermark的计算间隔调到了200毫秒但这带来额外开销实际性价比不高。最终方案是接受少量迟到事件用侧输出流单独处理延迟指标不受影响。第三个是下游ClickHouse的写入冲突。高峰期Flink往ClickHouse写的并发过高导致ClickHouse的merge线程打满查询延迟飙升。后来加了写入缓冲层Flink先写到Kafka的中间topic然后一个单独的作业以可控并发写入ClickHouse。压测最终数据50万 events/s事件速率下端到端延迟平均1.8秒P99 3.2秒满足业务指标。这个数据为我们上线提供了决策依据。5. 线上问题与排查实录5.1 背压问题从现象到Flink Web UI定位上线第三周业务方反馈实时风控决策变慢用户支付页面出现转圈。我们第一反应是看Flink Web UI的背压情况。打开作业页面果然看到某个算子的背压指标BackPressure为HIGH这意味着下游处理速度跟不上上游生产速度数据在算子输入缓冲区和网络缓冲区堆积。再往下看瓶颈定位在特征聚合算子它的输入数据量是其他算子的3倍。定位到算子后我们检查了并行度配置。这个算子默认并行度是8但上游Kafka分区有128个导致一个subtask要处理16个分区的数据单点压力太大。把并行度调到24后背压消失延迟降回正常水平。这个问题的经验是Flink Web UI的背压颜色OK、LOW、HIGH不能只看最开始的几个算子要顺着链条往上游找第一个出现HIGH的算子那里才是瓶颈所在。另外并行度设置最好以Kafka分区数为基准并留出一定冗余一个subtask处理的分区数上限建议控制在6以内。5.2 数据倾斜大客户和热门Key引发的连锁反应数据倾斜是我们遇到的第二大问题。巴西市场有个头部电商客户他们的流量能占到全站30%这个客户的用户行为事件在keyBy之后集中在少数几个subtask上导致这几个节点CPU 100%其他节点空闲。初次尝试是把大客户单独分流到独立的Flink作业处理效果有但增加了运维复杂度。后来发现大客户的流量特征不是一直大而是有波峰比如促销时间段内增长10倍平时也还好单独开一个作业纯属浪费。最终采用了“动态拆分大key”方案在ProcessFunction里维护一个计数器当检测到某个userId的事件速率超过阈值时把这个userId的key加上随机后缀拆成多个虚拟key让它们分散到不同subtask。窗口聚合的结果通过合并算子做最终归并。这个方案在压测中验证有效大客户高峰期整体吞吐提升约40%。5.3 数据丢失排查被忽略的幂等机制有一次双十一大促后核对数据发现某个指标少了一小部分数据量。我们第一反应是Flink的checkpoint出了问题但检查checkpoint的history发现所有checkpoint都成功了。排查过程花了一整天最后把问题定位到接入层SDK的重试逻辑上。SDK默认重试3次当Kafka发生可重试异常比如网络抖动时重发消息。但下游业务方在消费事件时没有做幂等重发导致部分数据被当成重复数据丢弃。当时很崩溃因为问题不在管道本身而是业务方的消费逻辑。后来我们在接入层给每条事件加上全局唯一ID并在存储层加了去重机制不管是SDK重发还是Flink重启后的数据回放下游看到的都是一份数据。同时推动所有核心业务方在消费端做幂等处理。这个案例给我们的教训是端到端精确一次是系统性问题不是Flink单点能解决的上下游都要配合。6. 监控体系与运维保障6.1 管道的核心监控指标监控是数据管道稳定运行的底气。我们围绕管道设计了四个维度的监控指标全部接入Prometheus Grafana第一层是Kafka集群指标吞吐量、消费延迟Lag、分区的Leader分布、请求处理的P99延迟。消费延迟是最核心的指标我们设置了2个阈值超过5000条触发Warning超过50000条触发PagerDuty告警。第二层是Flink作业指标checkpoint状态、背压比例、算子处理速率、状态大小、Watermark推进情况。其中Watermark推进是一个容易被忽略的指标如果Watermark一直不涨说明没有新事件进入很可能是上游断流这个告警能比消费延迟更早发现问题。第三层是存储层指标Redis的命中率、ClickHouse的查询延迟和写入队列长度、MySQL的慢查询数。第四层是业务指标端到端事件延迟从事件时间到写入存储的时间、核心特征服务的可用率、API的P99延迟。这里要强调一个实践业务指标的监控和运维指标的监控要放在同一个团队和同一套体系里。如果只看技术指标你永远不知道管道“活着但没干正事”。比如Flink作业在跑但业务方的一个上游系统挂了导致事件没产生技术指标都正常业务却已经停了。所以我们把业务方事件上报量的日环比、时环比也加了监控一旦异常能立刻看出是哪条数据链路出了问题。6.2 全链路告警策略与值班响应告警策略上我们坚持“少而准”。初期告警规则太多导致大家收到告警都麻木了真正出问题的时候反而没人响应。后来做了收敛只保留三类告警影响可用性的作业挂掉、checkpoint连续失败、影响延迟的消费延迟超阈值、背压持续HIGH、影响数据准确性的数据量环比剧烈波动、Watermark长时间不更新。值班响应机制也做了分级处理P0级告警作业挂掉、大面积延迟要求15分钟内响应30分钟内处理P1级告警延迟超阈值但未影响业务要求1小时内响应P2级告警容量危险、潜在风险只需要工作日处理。这套机制运行了半年效果比较理想平均故障恢复时间从最初的45分钟降到了15分钟左右业务方对我们的信任度明显提升。6.3 后续演进方向管道上线稳定运行后团队内部也在规划下一阶段的方向这里简单分享几个思路也都是围绕当前管道的痛点一是引入Flink SQL化改造把一部分相对固定的聚合逻辑用SQL实现降低新需求的上线周期。我们已经在内部搭建了SQL校验和调试环境业界这块生态已经比较成熟。二是把实时特征服务做成独立的平台能力让业务方能自助配置特征计算逻辑不需要每次提需求等我们排期。三是探索用Kubernetes替代YARN作为Flink的部署底座。Kubernetes在弹性伸缩上比YARN灵活配合Pod水平自动伸缩可以做到更细粒度的资源调控。不过这个工作排期靠后毕竟系统稳定运行的时候动底层架构需要勇气和足够的理由。从我个人经验来说这类实时数据管道的建设最难的部分往往不在写代码而在于对业务的理解和运维体系的建设。代码写得再漂亮没有一套靠谱的监控和告警上线之后就是在裸奔。希望这篇分享能给正在做类似系统的朋友一些启发少走一些我们踩过的弯路。
返回列表