一、先看整体链路
一笔交易相关事件,可能经过类似这样的链路:
行情源 / 交易柜台 / 订单系统 / 账户系统 | v Kafka | v Flink 清洗、关联、聚合、风控 | +-----------+-----------+-----------+ | | | | Redis Doris Iceberg Kafka 实时查询 实时分析 历史沉淀 下游消费在真实系统中,核心交易链路通常更强调确定性、低延迟、强一致和故障隔离,不会把所有逻辑都塞进 Flink。Flink 更常见于核心交易系统旁路的数据处理链路,或者处理交易系统已经产生的事件。
二、行情实时处理
券商每天会接收大量行情消息,例如:
股票代码、交易所、最新价、买一卖一、成交量、成交额、时间戳原始行情往往是逐笔成交或逐档盘口数据,频率非常高。Flink 可以做以下处理:
- 过滤无效行情、纠正格式和字段类型。
- 按股票代码和市场分区处理。
- 计算 1 分钟、5 分钟、15 分钟 K 线。
- 计算成交量、成交额、VWAP、涨跌幅等指标。
- 计算盘口深度、买卖盘不平衡度等实时指标。
- 生成行情告警,例如价格突破、成交量异常、涨跌幅超过阈值。
例如,每分钟生成一根 K 线:
开盘价 = 这一分钟的第一笔成交价 最高价 = 这一分钟的最高成交价 最低价 = 这一分钟的最低成交价 收盘价 = 这一分钟的最后一笔成交价 成交量 = 这一分钟所有成交量之和 成交额 = 这一分钟所有成交额之和这类计算本质上是按symbol分组的事件时间窗口聚合:
stream.keyBy(event->event.getSymbol()).window(TumblingEventTimeWindows.of(Time.minutes(1))).aggregate(newKlineAggregateFunction());如果行情存在乱序,就要结合事件时间和 Watermark。交易所消息的时间、消息进入 Kafka 的时间、Flink 接收时间不是一回事,不能随意使用处理时间替代事件时间。
但是,行情数据通常有一个特殊要求:行情量很大、实时性要求高、短时波动频繁。因此需要关注反压、分区数、热点股票、状态大小和下游写入能力。对于极低延迟的核心行情分发,Flink 也未必是唯一或最合适的组件,专用行情系统、内存计算组件或 C++/Java 低延迟服务可能更合适。Flink 更适合实时计算和分发衍生指标。
三、实时交易风控
这是互联网券商非常典型的 Flink 使用场景之一。
用户下单、撤单、成交、入金、出金、登录、设备变化等行为都可以形成事件:
OrderCreated OrderCanceled TradeFilled DepositCompleted WithdrawalRequested LoginSucceeded DeviceChangedFlink 可以对这些事件进行实时规则判断,例如:
同一账户 1 分钟内下单超过 100 次 同一设备关联多个账户 短时间内频繁撤单和重新下单 账户刚异地登录,随后立即发起大额交易 短时间内入金后快速出金 某账户的交易行为显著偏离历史模式典型处理方式是:
事件进入 Kafka -> Flink 按 account_id 或 device_id 分组 -> 维护时间窗口状态 -> 匹配风险规则 -> 输出风险事件 -> 风控服务、告警系统或人工审核平台处理例如统计账户最近 5 分钟的下单次数:
orders.keyBy(Order::getAccountId).window(SlidingEventTimeWindows.of(Time.minutes(5),Time.seconds(10))).process(newHighFrequencyOrderRule());这里的窗口含义是“最近 5 分钟”,每 10 秒重新评估一次。实际生产中也可能使用 Keyed State 加定时器实现更灵活的规则,因为每个规则的窗口长度、过期方式和输出策略并不完全相同。
需要区分两类风控:
- 交易前强实时风控:用户点击下单后,必须在极短时间内同步返回是否允许下单。这通常由交易柜台、风控服务和规则引擎完成,不能简单依赖一个异步 Flink 作业。
- 交易后实时监控:成交或行为发生后,持续检测异常模式、账户风险和市场风险。Flink 很适合做这一类流式检测。
因此,Flink 更常作为交易后监控、实时特征计算和风险事件发现组件;是否阻断交易,要看具体链路的延迟和一致性要求。
四、客户资产和交易画像
券商需要实时回答很多问题:
客户当前持有哪些资产? 最近 7 天交易金额是多少? 近 30 天交易次数和盈亏如何? 客户偏好港股、美股还是基金? 客户是否属于高频交易客户? 客户最近是否长期未交易?Flink 可以消费订单、成交、持仓变更、入金、出金、行情等事件,实时维护客户特征:
近 1 天 / 7 天 / 30 天交易次数 近 7 天交易金额 买入和卖出金额 不同市场的交易占比 最近交易时间 活跃天数 平均客单价 胜率或盈亏相关指标一个简化链路是:
成交事件 -> Flink 按 customer_id 分组 -> 计算窗口指标 -> 生成画像标签 -> 写入 Redis / Doris / ClickHouse例如:
7 天交易金额 > 100 万 -> 高净值活跃客户 30 天没有交易 -> 沉默客户 近 7 天交易次数快速上升 -> 活跃度上升 港股成交额占比 > 80% -> 港股偏好这里使用 Flink 的理由不是“7 天只能用 Flink”,而是客户事件不断发生,画像需要持续更新。如果只要求每天凌晨计算一次,Spark + Iceberg 完全可以胜任。
五、账户余额、持仓和资产快照
交易系统产生的成交事件可以用于构建下游实时资产视图:
成交回报 -> 更新可用持仓 -> 更新持仓数量 -> 关联最新行情 -> 计算市值 -> 计算账户资产快照例如:
持仓市值 = 持仓数量 * 最新价格 账户总资产 = 现金资产 + 各类持仓市值 + 其他资产这类场景可以用 Flink 做事件关联和实时聚合,但必须特别注意:
- 不能把行情价格和成交状态简单地当作同一种事件。
- 成交回报可能重复、乱序或延迟,需要幂等处理。
- 账户资产涉及金额,必须明确精度、币种、汇率和时间点。
- 关键账务余额应以权威账户系统为准,Flink 计算结果通常是查询视图或下游缓存,不能取代总账。
- 需要支持重放和对账,发现结果不一致后可以基于事件重新计算。
更稳妥的模式是:
核心账务系统保存权威结果 Flink 根据账务事件构建实时查询视图 定期与权威账务系统对账六、实时盈亏和风险指标
通过持仓事件和行情事件,Flink 可以计算:
持仓市值 浮动盈亏 当日盈亏 仓位比例 集中度 杠杆率 保证金占用 维持担保比例这类计算通常需要把两类流进行关联:
持仓流 + 行情流概念上类似:
持仓数量、成本价、账户信息 + 股票最新价、汇率、合约乘数 | v 实时资产指标Flink 可以用 Broadcast State 分发规则参数,也可以使用 Keyed State 保存账户或证券的当前状态。对于大规模账户,需要根据数据分布设计 Key:不能让所有账户都集中到少数几个并行实例。
如果是美股、港股、基金、期权等多个市场,还要处理不同交易时段、币种、汇率、节假日和合约规则。例如美股收盘而港股开盘时,资产快照不一定可以用同一套时间窗口直接计算。
七、交易事件清洗、去重和数据质量检查
Kafka 至少一次投递、消费者重启、上游重试都可能产生重复消息。因此 Flink 常用于:
字段校验 类型转换 非法金额过滤 交易状态校验 按业务主键去重 补充市场和账户维度 输出异常数据例如以trade_id去重:
同一个 trade_id 重复到达 -> Flink 状态中检查是否处理过 -> 第一次输出 -> 后续重复事件丢弃或转入审计流不过去重状态不能无限增长。需要根据业务确定:
- 去重主键是什么;
- 事件可能迟到多久;
- 去重状态保留多长时间;
- 过期后再次出现同一 ID 怎么处理;
- 是否需要把原始事件写入 Iceberg 以便审计和重算。
生产系统常用“业务幂等 + Flink 状态 + 下游幂等写入”共同保证结果可靠,不能只依赖 Flink 的 checkpoint 就认为业务天然不重复。
八、实时行情和交易告警
用户关心的提醒可能包括:
价格突破自选价 涨跌幅达到阈值 成交量突然放大 新股上市 持仓触及止盈止损条件 账户保证金比例过低 订单成交或部分成交Flink 可以消费行情和交易状态,匹配用户订阅条件,然后输出通知事件:
行情事件 -> 按 symbol 找到订阅用户 -> 判断用户条件 -> 去重和限频 -> Kafka / 推送服务 / 短信服务这里最容易被忽略的是“用户订阅条件”的动态变化。用户可能随时添加、修改或删除自选股。通常需要把规则变化作为另一条事件流,通过广播状态或外部规则服务同步到计算任务。
还要做告警限频,否则同一个价格条件在高频行情下可能连续触发数千次。常见控制包括:
同一用户、同一证券、同一规则在一段时间内只通知一次 价格必须先离开阈值,再次跨越时才重新触发 推送失败要重试,但不能无限重试九、运营指标和实时数据看板
券商运营团队可能实时查看:
在线用户数 下单人数 成交人数 订单成功率 订单延迟 行情延迟 入金和出金金额 各市场交易额 系统错误率Flink 可以按分钟或秒级窗口聚合这些指标,再写入 Doris、ClickHouse 或监控系统。它适合计算实时指标,但通常不负责最终可视化。
例如:
每 1 分钟统计不同市场的订单量和成交量 每 10 秒统计订单失败率 按渠道统计登录、开户、入金转化情况这类指标对迟到数据和修正结果的要求,往往低于账务数据,开发实现相对简单。
十、实时推荐和营销触达
当用户出现某种行为时,系统可以实时更新标签并触发后续动作:
新用户完成开户 -> Flink 识别开户完成事件 -> 更新客户生命周期状态 -> 触发新手任务或教育内容 客户连续多日查看某市场行情 -> 更新市场偏好标签 -> 推荐相关行情或研究内容这种场景通常不会由 Flink 直接决定所有营销内容,而是由 Flink 生成事件和特征,再交给推荐服务、营销平台或 CRM 系统执行。
十一、Flink 和 Spark 在券商类业务中的合理分工
可以这样理解:
Flink:处理“现在发生的事情” Spark:处理“历史数据重新计算” Iceberg:保存“可追溯的明细和历史结果” Redis:提供“极低延迟的当前状态查询” Doris / ClickHouse:提供“多维分析查询” Kafka:连接各个实时系统典型组合可能是:
Kafka ├── Flink │ ├── 实时行情指标 │ ├── 交易后风控 │ ├── 客户画像 │ ├── 实时资产视图 │ └── 告警事件 │ └── Iceberg └── 原始明细和历史事件 Spark ├── T+1 客户画像重算 ├── 历史盈亏修正 ├── 风控规则回溯验证 ├── 数据质量核对 └── 离线报表十二、哪些地方不应该直接依赖 Flink
以下系统通常需要更加严格的专用实现:
- 订单撮合引擎;
- 交易所连接和柜台核心链路;
- 权威资金账本;
- 最终清算和结算;
- 必须同步返回的交易前强校验;
- 需要严格审计和事务保证的核心账务写入。
Flink 可以消费这些系统发出的事件,做旁路计算、监控、视图构建和风险分析,但不能因为 Flink 支持 Exactly-Once,就把它等同于完整的金融账务系统。Exactly-Once 主要描述 Flink 处理和特定连接器的语义,不能自动保证整个端到端业务链路的业务一致性。
十三、一个更贴近实际的例子:实时维护客户近 7 天画像
成交回报进入 Kafka -> Flink 校验 trade_id、customer_id、amount -> 按 trade_id 去重 -> 使用 trade_time 设置事件时间和 Watermark -> 按 customer_id 分组 -> 维护近 7 天交易次数、金额、市场偏好 -> 输出客户画像标签 -> Redis 保存当前标签 -> Doris 保存可分析结果 -> Iceberg 保存原始明细 -> Spark 每日重算并校正历史结果这个例子里使用 Flink 的理由是“成交发生后需要快速更新”,不是因为“7 天窗口必须用 Flink”。如果业务只是每天生成一份客户报表,那么 Spark 会更简单。
最后的判断标准
对于这类业务,下面这些问题适合 Flink:
是否需要秒级或分钟级响应? 事件是否持续不断地产生? 是否要处理乱序、迟到和状态? 是否需要实时关联多个事件流? 是否需要实时触发告警或下游动作?如果多数答案是“是”,Flink 很有价值。
如果需求是:
每天凌晨统计昨天数据 基于 Iceberg 重算最近 30 天 生成历史报表 进行大规模回溯分析Spark 往往更合适。
一句话概括:在这类互联网券商里,Flink 更像实时数据和实时状态计算引擎,负责把订单、成交、行情、账户行为快速转换成风控事件、资产视图、客户标签和实时指标;它通常不取代撮合、清算和权威账务系统。