尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

体育赛事数据系统架构设计与实战:实时链路从Kafka到Flink

体育赛事数据系统架构设计与实战:实时链路从Kafka到Flink 做体育赛事数据系统这件事听起来不就是“记个比分嘛”真正动手你才会发现这玩意儿是实时系统里最磨人的一类。比分只是最外层的结果底下还压着事件流、比赛时钟、状态机、多源数据比对、突发流量、历史数据归档一连串事。我这套系统从需求立项到上线跑了四五个月中间经历了完整的技术选型和两轮架构调整今天把决策过程、落地细节和踩过的坑一次性说清楚。文章主要覆盖实时数据管道、存储分层、对外API与WebSocket推送这几个核心模块适合正在设计赛事数据平台、体育数据服务API或者做直播数据系统的朋友参考。1. 先想清楚体育赛事数据系统到底在解决什么问题1.1 核心需求拆解体育赛事数据系统的表面需求很简单收集比赛数据处理后对外展示和提供查询。但一旦展开需求就变成了这样第一实时性要求极高。用户在看直播时盯着比分数据系统必须让比分误差控制在秒级以内最好能做到亚秒级。比如足球进球事件从现场裁判确认到数据源发出信号再到用户手机弹出推送这条链路里任何一环都不能拖后腿。第二支持多项目、多赛事类型。足球、篮球、网球、电竞都有各自的数据特征。足球有半场、补时、红黄牌篮球有节次、暂停、个人犯规网球有盘、局、分。如果为每个项目单独设计一套数据结构后期维护会非常痛苦必须抽出一套通用的赛事数据模型。第三稳定性要求极其苛刻。真正的比赛不会等你修bug晚上八点开球的德甲比赛不会因为你的服务挂了就推迟。数据系统必须以“比赛期间不能宕机”作为最低底线而不是“尽量可用”。第四对外服务能力。系统本身不只是一块给运营看的可视化大屏还要作为数据中心向媒体、合作方、内部业务方提供数据API。这些接口的QPS压力、鉴权逻辑、数据一致性都是需求里绕不开的部分。1.2 赛事数据和其他实时系统的差异很多人觉得赛事数据系统和股票行情、物流轨迹差不多都是“实时数据的采集与分发”其实差异很大。股票行情的数据源是确定的、结构化的交易所统一发布数据质量有保障。赛事数据不一样数据源可能是现场人工录入、半自动传感器、官方数据服务商不同项目的信号质量参差不齐甚至会出现同一场比赛两个数据源给出的比分不一样的情况。物流轨迹的状态机很简单待揽收、运输中、已签收基本是单向推进的。比赛状态机复杂得多未开始、进行中、中场休息、已结束、技术暂停、延期、取消、腰斩而且状态可以回退——比如篮球第四节还剩两分钟时因技术原因中断之后恢复比赛状态就要从“已结束”退回“进行中”。这一点非常坑后面我会展开讲。比赛的时间语义也很特殊。服务器收到一条“足球比赛第67分钟进球”的数据这里的第67分钟是比赛时钟不是系统时钟。系统必须自己维护每个比赛当前进行到第几分钟再结合数据源的时间戳做判断绝不能拿服务器本地时间当作比赛时间。1.3 系统的目标形态在动手选型之前我给系统画了这样一个目标形态一条实时的数据管道从多个数据源接收原始比赛事件一个统一的数据核心维护所有比赛的状态、比分、事件和统计信息一组对外服务通过API和WebSocket为业务方提供数据和推送一个分析侧的数据仓库沉淀历史比赛记录供数据分析和商业智能使用。也就是说这是一个典型的“输入多源、输出多端”的中间平台。想清楚目标形态之后选型就变得有的放矢了。2. 技术选型为什么最终是这套组合技术选型的本质不是选最好的工具而是选最适合这套系统特性的组合。我在选型前明确了几条硬约束团队主力语言是Java和Go预算在可控范围部署环境是自建Kubernetes集群数据量级在每天几万个事件峰值请求在每秒几万左右。在这些约束下做对比分析结论会清晰很多。2.1 消息队列Kafka几乎是唯一解数据管道的第一站必然是消息队列。我在Kafka、RabbitMQ、Pulsar之间做了对比侧重点不在基准性能而在赛事数据场景真正需要的几个能力。评估维度KafkaRabbitMQPulsar吞吐能力极高中等极高数据回溯重新消费天然支持按offset回放较弱依赖消息过期策略天然支持分层存储一致性语义强大搭配Flink生态一般较强运维复杂度中等低高组件多生态适配Flink/Spark/ClickHouse全覆盖主要面向传统微服务相对较新RabbitMQ的路由灵活、AMQP协议在传统业务里很强但吞吐量不是赛事场景的瓶颈问题在于它的数据保存和回溯能力偏弱。赛事数据有一个很讨厌的特性晚到事件。比如一个进球事件因为现场信号问题迟到了两分钟才发出来这期间系统已经基于不完整数据发出了快照。要修正它消息队列必须能保存足够长时间的数据供重新消费这在RabbitMQ里做起来很别扭。Kafka的日志就是它的核心优势——数据按分区顺序存储保留期可以按天配置消费者可以随时从任意offset重新读取。这正好匹配赛事数据“修正、补发、重放”的需求。Pulsar在分层存储和跨区域复制上确实有亮点但部署运维复杂度明显高团队没有专职人去维护它最终还是选了Kafka。具体配置上我用了三个Kafka节点副本因子3。Topic按项目类型拆分football_events、basketball_events、match_snapshot、match_status_change分区数以6为核心后续按流量再扩。赛事场景的写入峰值其实是可预估的——同一时间进行中的比赛数量有限单场比赛事件频率最高也就是每几百毫秒一条Kafka本身的压力不大真正考验它的是消费者的链路设计。2.2 流处理为什么必须上Flink数据到了Kafka之后要经过处理解析、清洗、归一化、状态计算、规则匹配。这个环节我对比了Flink和Spark Streaming。Spark Streaming基于微批吞吐量极高但延迟在低延迟场景下并不理想。赛事场景中进球事件从数据源发出到入库端到端延迟要求控制在1到3秒内Spark Streaming要把微批窗口调得非常小才能满足反而丧失了对批处理的优化优势。Flink是真正的流式计算引擎事件驱动天然支持event time和watermark。赛事数据场景里这两样东西是刚需event time让系统能正确处理迟到的比赛事件watermark让系统知道该继续等多久再输出结果。比如一个进球事件迟到了Flink可以根据事件自带的时间戳把它计算到正确的比赛时间槽里而不是按到达时间处理。我选Flink还有一层考虑状态管理。比赛状态机、比分快照、事件去重这些都需要状态存储Flink的keyed state配合RocksDB可以很方便地维护“每场比赛的当前状态”而不需要每来一条消息都去查一次MySQL。这个能力让实时链路的代码简洁很多也天然规避了实时计算与数据库之间高频交互的性能问题。2.3 存储体系热、温、冷三层分离赛事数据系统的存储天然分三类三类特性差异很大绝不能塞进同一个数据库里。热数据比赛进行中的实时状态和最近几分钟的快照。特征是写频繁、读更频繁、只关心最新值。用户看直播时最常访问的就是当前比分和最近事件列表。Redis在这里是绝对主力我用String结构存比赛状态JSONList结构存最近事件流Hash结构存比分字段TTL设为比赛结束后24小时防止垃圾数据堆积。温数据已结束比赛的详细记录和事件列表。需要支持按球队、赛事、日期维度查询写入后基本不变更。MySQL完全胜任但要注意按时间做分表策略。我用了按年份分表每张表存当年全部比赛数据定期归档。冷数据历史赛季汇总和数据仓库。这些数据用来做统计分析、趋势研究查询模式是复杂的聚合查询。迁到ClickHouse之后性能提升非常明显——千万级行数的聚合查询从MySQL的几十秒降到百毫秒级后期所有分析报表都跑在ClickHouse上。2.4 API与推送通道一份数据多种出口对外接口这层我做了两个出口。查询类API用RESTful JSON就够了媒体方和业务方对接成本最低写清楚OpenAPI文档用签名认证控制访问权限。实时推送用WebSocket比分和事件一旦变化就主动推给订阅方避免前端轮询带来的压力和延迟。内部服务之间则是gRPC。主要考虑是性能好、连接复用而且proto文件能约束接口结构避免前后端各维护一套文档导致的口径不一致。这里有个经验不要把内部的gRPC接口直接暴露给外部客户。外部需要的是一套稳定的、版本化的REST API内部接口可以随便改外部接口要尽量向后兼容。3. 架构落地从骨架到细节选型做完接下来是架构落地。整体分四层数据采集层、实时处理层、存储层、对外服务层。下面从下往上讲每一层的落地细节。3.1 数据采集层多源接入与数据归一化数据采集层要解决的第一件事是多源适配。我的数据源有三类官方数据服务商提供的实时XML/JSON流合作场馆回传的数据接口以及少量的人工录入终端。三种数据源的格式完全不同但到了系统内部必须统一成一种内部事件格式。我定义了一套内部事件模型核心字段如下event_id全局唯一的原始事件ID用于去重match_id比赛IDsport_type项目类型event_type事件类型goal、card、foul、score_change等event_time比赛时钟时间第几分钟/第几秒occur_at事件发生的UTC绝对时间payload事件负载JSON各项目自定义数据采集服务的职责就是拉取外部数据源原始数据 → 解析并映射到内部事件模型 → 做基础合法性校验比赛ID是否存在、时间戳是否合理、事件ID是否重复→ 写入Kafka对应Topic。这里最考验工程能力的不是“写通”而是校验逻辑。比如足球比赛里出现了一个“第120分钟进球”事件这合法吗要看比赛有没有加时。篮球里“第三节的比分为200比50”是不是数据源故障了这些规则如果不在入口拦截后面实时计算层就会算出各种诡异的统计结果。3.2 实时处理层Flink作业的核心设计实时处理层是整套系统的大脑我拆成了三个Flink作业事件清洗作业从Kafka读取原始事件做格式校验、协议转换、去重然后写回Kafka的clean_topic。这个作业不维护状态纯粹做ETL。比赛状态计算作业核心作业。用match_id作为key按事件时间排序处理每场比赛的事件。每处理一条事件就更新一次比赛状态比分是多少、当前是第几分钟、最近发生了哪些事件。状态保存在Flink的keyed state里定期快照到Redis作为对外数据源。这个作业直接用event time做窗口解决迟到事件的问题。统计聚合作业独立于比赛状态之外专门计算球员、球队的统计数据射门次数、控球率、犯规次数等结果写Redis和ClickHouse。Flink作业在Kubernetes上以原生模式运行每个作业的并行度根据赛事密度动态调整。重要比赛期间并行度提到最高无赛期降到1甚至暂停作业节省资源。扩缩容通过修改并行度参数后重启Job实现整个过程十分钟内完成。3.3 存储层三级存储的具体用法前面说了热、温、冷三层存储这里讲具体用法。Redis存的是“对外可读的实时状态”。我把每一场比赛整理成一个统一的视图结构包含基础信息球队、赛事、开球时间、当前状态状态码、当前比分、比赛时钟、最近事件列表、关键统计数据。这个视图由Flink计算作业在每场比赛状态变更时更新到Rediskey规则是match_view:{match_id}TTL设为比赛结束后48小时。MySQL存的是“有业务含义的正式数据”。比赛开始前先建一条比赛主记录状态为“未开始”。比赛结束后把最终比分、最终结果、赛事归属等字段更新到位。这里不存明细事件明细事件都在ClickHouse。MySQL承担业务数据主记录职责是给后台管理系统、订单系统、报表系统用的。ClickHouse存的是“全量事件明细和统计数据”。每天凌晨从Redis和Kafka导入前一天的比赛快照和事件明细按比赛写入MergeTree表。ClickHouse的查询性能对这套数据量来说绰绰有余我甚至不需要二级索引单靠分区和排序键就能跑得很流畅。为了控制成本只保留最近两个赛季的明细数据更早的压缩归档到对象存储走离线查询。3.4 对外服务层REST API与WebSocket推送的实现对外服务层是直接面向客户的设计上必须干净。REST API主要提供三类接口比赛列表查询按日期、赛事、球队筛选单场比赛详情比分、状态、事件列表、统计数据历史数据查询球队历史战绩、交锋记录、统计排名接口全部无状态基于token鉴权限制每用户QPS。考虑到外部客户可能频繁轮询比分接口接口层做了统一缓存读Redis缓存的比赛视图命中就直接返回未命中再读MySQL。缓存策略是“比赛状态一旦变化立即更新Redis”因此比分接口的读取一定是最终一致的而且几乎没有数据库压力。WebSocket推送通道处理“主动变化”。用户订阅一个match_id或赛事ID服务端在该比赛状态变化时推送消息。实现上我用了Redis的Pub/Sub功能Flink作业每次更新比赛视图时同时向Redis频道match_update发布一条消息推送服务订阅该频道收到消息后再推给对应的WebSocket连接。这个方案的好处是Flink作业完全不需要关心有多少个推送节点推送服务自身的扩展也不会影响实时计算链路。推送的格式要轻量不要每次都推整场的完整数据。我推的是增量事件{ type: goal, match_id: 100234, event_id: evt_889123, payload: { team: home, player_id: p_3021, minute: 67 } }前端拿到增量事件后自己合并本地的比赛状态这样可以大幅减少网络传输和前端渲染压力。4. 实操细节测试环境里永远学不到的东西架构设计是一回事真正在项目里跑起来是另一回事。这一节我把实操中最重要的几个细节单独拎出来讲每一个都对应我在生产环境里实实在在摔过的跟头。4.1 比赛状态机的设计是全局最难的部分比赛状态机是整个系统最容易翻车的地方。我一开始的设计是线性推进的未开始 → 进行中 → 已结束看起来够简单了吧上线不到两周就被打脸。先是篮球比赛因为场馆技术问题中断了20分钟恢复后继续打然后是足球比赛因为天气原因暂停半小时后又恢复最狠的一次是一场网球赛数据源先报了“已结束”结果半小时后又报了“恢复比赛”最后才真正结束。后来我把状态机改成了带“可回退”的模型核心状态如下SCHEDULED已安排未开始LIVE进行中INTERRUPTED中断可恢复SUSPENDED暂停预期马上恢复FINISHED已结束CANCELLED取消终态POSTPONED延期重新安排我把FINISHED设计成不可逆状态任何来自数据源的“从FINISHED回到LIVE”的事件先进入人工审核队列由运营确认后手动纠正而不是自动回退。原因很简单比赛结束后系统会触发一串副作用——生成战绩记录、更新排名、通知下游结算如果状态自动回退这些副作用就全乱了。宁可人工干预也不能自动翻车。此外每个状态变更都要附带原因字段和时间字段比如LIVE要记录“恢复时间为本地时间14:32原因为暂停结束”。这些元信息在后期排障时就是救命稻草能让你快速还原一场比赛当天到底经历了什么。4.2 迟到的比赛事件系统的“时空悖论”刚才提到Flink的event time这里展开讲一个具体场景。足球比赛第80分钟有一个进球但是数据源因为现场网络故障直到比赛第83分钟的哨响之后才把这条事件发出来。也就是说Kafka里已有的数据是“第82分钟比分1比1”新来了一条“第80分钟比分变成1比2”的事件。如果不用event time按到达时间顺序处理系统就会在“第83分钟”这个时间点突然把比分改成1比2然后在场的观众和直播平台的比分就会“穿越”——明明现在已经是1比2但比赛事件流里还记录着“第80分钟进了球”逻辑错乱。用Flink的event time加watermark就顺理成章了事件按自身的比赛时间排序第80分钟的进球虽然到达得晚但在事件流里会被排在第80分钟的位置处理结果会回溯修正当时输出的快照。Flink的watermark机制允许系统在“等待迟到事件”和“及时输出结果”之间找到平衡点我配置的allowed lateness是2分钟超过这个时间到达的事件交给侧输出流做人工核对而不是直接丢弃。这里有个直观的教训处理迟到事件靠的是事件本身的时间戳而不是到达时间。数据源给不给时间戳在选型阶段就要确认好。如果你的数据源连比赛事件时间都给不出来那Flink的event time优势就完全发挥不出来整个选型逻辑都要重写一遍。4.3 缓存与热点千万用户看同一场比赛怎么办赛事数据系统和其他系统最大的不同在于流量集中度。一场热门的焦点战可能有几千万观众同时盯着比分接口或者直播频道的推送通道。所有流量都打在同一个match_id上这是典型的“热点key”问题。我在Redis层面做了三个处理。第一比分接口的读缓存尽量短。比分更新频率低通常每秒一次算顶天了但读取频率极高。我在Redis里存的比赛视图更新时直接覆盖写读取时不需要做额外一致性校验因为比赛状态是最新的就是对的。TTL设短一点比如30秒万一Flink作业挂了导致视图没更新缓存也会自动过期让请求回源MySQL不会长时间卡在脏数据上。第二热点key加随机后缀。这是一个老方案但对赛事场景非常有效。比如某个接口QPS冲到每秒几万单个key的读压力太大就把比赛视图复制成多份加上1到5的随机后缀把读请求打散到不同副本上。注意这个方案只适用于只读且允许短暂不一致的场景比分接口正好满足。第三WebSocket推送的订阅压力管理。几百万在线用户订阅同一场决赛推送服务每来一条进球事件就要扇出几百万条WS消息这个扇出在单机上是扛不住的。我的做法是把订阅关系放在Redis的Hash里key是match_idfield是连接IDvalue是节点地址推送服务定期从Redis拉取本节点的订阅列表然后在本节点扇出。这样就实现了按比赛维度的推送负载均衡一个WebSocket节点只负责自己连接的订阅者不会出现所有流量压到某一台机器的情况。4.4 多项目数据模型的抽象与取舍前面提过我定义了一套内部事件模型实际用下来这套抽象在“通用性”和“易用性”之间是有博弈的。足球的payload字段是“进球球员、助攻球员、进球方式”篮网的payload是“得分球员、得分方式、助攻球员”两者有交集但不完全一样。我最终采用了“通用外壳 项目特定payload”的策略所有项目共用事件ID、比赛ID、时间戳这些外壳字段具体的业务数据全部塞进JSON payload。Flink计算作业在解析事件时根据sport_type选择对应的解析器把JSON payload映射成项目特定的内部对象。这个取舍的好处是新增一个体育项目时只需要写好数据源适配器、payload解析器、统计聚合规则核心的实时处理链路完全不用改。坏处是JSON字段没有强类型约束前后端出现字段名不一致时只能靠联调发现。为了弥补我给每个项目维护了一份payload schema规范文档并且在新项目接入时做一次全量样本的字段校验把风险控制在接入阶段。5. 踩坑实录那些测试环境永远试不出来的问题这一节是全文最有价值的部分。下面每个问题我都标记了严重程度和出现背景方便你判断如果自己做同类系统该优先防哪些雷。5.1 Kafka分区倾斜同一场比赛的所有事件压在一个分区问题表现Kafka某Topic的部分分区消费延迟高消费者lag持续增长另外一些分区几乎闲置。原因排查当时我把Topic的分区key设成了match_id目标是把同一场比赛的事件都分到同一个分区以保序。但热门比赛产生的事件量是其他比赛的几十倍全部压在一个分区里消费者处理不过来自然就lag了。解决办法不再简单用match_id做分区key改成“match_id在分区数范围内取模 按事件号分段”。具体来说我把每个match_id的事件再划分为多个事件段每段就是一个分区key这样单场比赛的事件会分散到多个分区同时同一段内仍然保序。代价是Flink端需要多维护一层“按段聚合”的逻辑但换来了分区负载均衡。经验总结保序和均匀在Kafka里经常打架面对热点场景不能无脑用业务ID做key必须考虑热点分区的单分区吞吐瓶颈。5.2 Flink状态无限增长没设置TTL的惨痛教训问题表现运行两周后Flink作业的状态存储RocksDB大小翻了几倍Checkpoint耗时越来越长最终Job出现周期性反压。原因排查我的keyed state里存了每一个事件的历史列表但忘了给状态设置TTL。一场比赛结束后状态没有被清理长期躺在RocksDB里占空间。解决办法给状态配置StateTtlConfig按“比赛结束后24小时”自动过期。设置完之后RocksDB大小立刻稳定下来Checkpoint时间回到正常水平。经验总结Flink的状态管理很强大但必须对状态的生存周期有清晰规划。没有TTL的keyed state最后一定会变成checkpoint里的定时炸弹。5.3 比分接口在大赛开始时刻的缓存穿透问题表现某次大型赛事开幕当天比分接口响应时间从10ms飙到500msMySQL的CPU冲到80%。原因排查接口逻辑按“Redis未命中 → 回源MySQL”设计。那天全平台用户几乎在同一时间刷新页面而Redis里的比赛视图TTL是30秒大量请求集中在TTL过期的瞬间同时穿透到MySQL把数据库打坏了。解决办法做了两级缓存。一级是Redis的比赛视图二级是本地进程内缓存TTL分别设成30秒和10秒。请求先查本地缓存没命中再查Redis最后才回源MySQL。同时给回源路径加了一个简单的“单飞”锁——同一场比赛的穿透请求只允许一个去查MySQL其余等待第一份结果。这几层叠加之后比分接口的P99响应时间稳定在30ms以内数据库压力完全可控。经验总结热点数据的访问链路每一层都要考虑防穿透不要觉得Redis缓存已经挡住绝大多数流量就高枕无忧。5.4 WebSocket掉线、重连与消息补偿问题表现用户观看直播时频繁出现比分突然跳变少了几条事件。原因排查WebSocket连接在移动网络下经常断掉重连而重连之后服务端只发“当前最新状态”不会补发断开期间漏掉的事件前端直接拿最新状态覆盖本地状态导致漏掉的事件永远不被感知。解决办法重连时前端带上自己本地已收到的最新事件IDlast_event_id推送服务从这个事件ID之后开始补发增量事件前端收到增量后逐条合并再更新最新状态。同时本地状态和最新状态之间做一次diff校验如果发现差异就拉一次全量快照兜底。经验总结实时推送系统里增量、快照、断点续传这三件套缺一不可。只推最新状态会让用户看到跳变只推增量会让断线用户永远缺事件。5.5 多数据源比分冲突以谁为准问题表现同一场比赛数据源A显示2比1数据源B显示1比1两边都坚称自己权威。原因排查赛事数据领域天然存在双源竞争因为数据源的分工不同——A可能是官方数据源B是合作场馆比赛结算时官方数据源的优先级应该更高。解决办法在系统里给每个数据源配置了数据源优先级。正常时以高优先级数据源为准低优先级数据源事件只做参考当两个数据源事件存在冲突时系统自动采用高优先级数据源的事件同时把低优先级数据源的事件标记为conflict进入待审核队列。人工审核后可以选择保留或纠正每次纠正都会生成一条审计日志。这套机制上线后因为数据源冲突导致的赛事数据事故基本清零。经验总结多源数据系统的核心能力不是“选一个最好的源”而是“管理好源与源之间的冲突”。谁高谁低、怎么审计、怎么回滚都要在设计阶段就定好规则。6. 最后再聊两句经验系统上线稳定运行大半年回过头看这套选型和架构我最大的体会是三句话。第一句是实时系统的坑绝大多数不在实时计算本身而在边界条件。状态回退、迟到事件、多源冲突、缓存穿透这些全是数据和时间在边界上的摩擦。Flink、Kafka、Redis这些工具都很成熟真正拉开系统差距的是你对比赛业务本身的理解深度。第二句是选型时一定要考虑团队运维能力。Pulsar再先进没人会维护就是负资产ClickHouse再快DBA不熟就是定时炸弹。我的选型逻辑很简单——每个组件都要有至少一个人能独立负责故障排查这是我列技术栈清单时的硬约束。第三句是别把外部数据源当“唯一真相”。数据源厂商的能力和稳定性参差不齐你必须在自己的系统里建立一套独立的数据校验和审计机制。现在每场比赛结束后系统会自动跑一遍完整性和一致性校验发现异常就发告警这套机制帮我们拦截了很多外部数据源自己的bug。如果你也在设计赛事数据系统建议从“数据源接入、实时链路、热点访问、状态机”这四个维度着手先跑通主链路再逐步优化。这套架构不复杂但每一层都有值得打磨的细节。我自己还在持续迭代冷数据分析和自动化数据质量检测的部分后面有新的沉淀再单独成文分享。
返回列表