news 2026/8/29 22:44:51

交易类项目-flink

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
交易类项目-flink

一、先看整体链路

一笔交易相关事件,可能经过类似这样的链路:

行情源 / 交易柜台 / 订单系统 / 账户系统 | v Kafka | v Flink 清洗、关联、聚合、风控 | +-----------+-----------+-----------+ | | | | Redis Doris Iceberg Kafka 实时查询 实时分析 历史沉淀 下游消费

在真实系统中,核心交易链路通常更强调确定性、低延迟、强一致和故障隔离,不会把所有逻辑都塞进 Flink。Flink 更常见于核心交易系统旁路的数据处理链路,或者处理交易系统已经产生的事件。

二、行情实时处理

券商每天会接收大量行情消息,例如:

股票代码、交易所、最新价、买一卖一、成交量、成交额、时间戳

原始行情往往是逐笔成交或逐档盘口数据,频率非常高。Flink 可以做以下处理:

  1. 过滤无效行情、纠正格式和字段类型。
  2. 按股票代码和市场分区处理。
  3. 计算 1 分钟、5 分钟、15 分钟 K 线。
  4. 计算成交量、成交额、VWAP、涨跌幅等指标。
  5. 计算盘口深度、买卖盘不平衡度等实时指标。
  6. 生成行情告警,例如价格突破、成交量异常、涨跌幅超过阈值。

例如,每分钟生成一根 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 DeviceChanged

Flink 可以对这些事件进行实时规则判断,例如:

同一账户 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 做事件关联和实时聚合,但必须特别注意:

  1. 不能把行情价格和成交状态简单地当作同一种事件。
  2. 成交回报可能重复、乱序或延迟,需要幂等处理。
  3. 账户资产涉及金额,必须明确精度、币种、汇率和时间点。
  4. 关键账务余额应以权威账户系统为准,Flink 计算结果通常是查询视图或下游缓存,不能取代总账。
  5. 需要支持重放和对账,发现结果不一致后可以基于事件重新计算。

更稳妥的模式是:

核心账务系统保存权威结果 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 更像实时数据和实时状态计算引擎,负责把订单、成交、行情、账户行为快速转换成风控事件、资产视图、客户标签和实时指标;它通常不取代撮合、清算和权威账务系统。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/8/29 22:39:09

蓝桥杯单片机频率计设计:测频法与测周法融合及自动量程切换实战

1. 项目缘起与核心需求:为什么蓝桥杯单片机赛题偏爱频率计?如果你参加过蓝桥杯电子类单片机组(省赛或国赛)的竞赛,或者仔细研究过历年的真题,你会发现一个高频出现的“老朋友”——频率计。从早期的51单片机…

作者头像 李华
网站建设 2026/8/29 22:38:33

第四范式前端笔试复盘:从JS到算法,原理型选手的筛选

1. 这套笔试到底在考什么:题型分布与考察逻辑 2020年秋招那阵子,我收到第四范式的笔试链接时,第一反应是:一家做AI平台的公司,前端笔试会不会剑走偏锋,比如考一堆机器学习概念或者算法推导?真点…

作者头像 李华
网站建设 2026/8/29 22:38:22

markitdown:两条命令把 PPT 转成 AI 能直接读的 Markdown

markitdown:两条命令把 PPT 转成 AI 能直接读的 Markdown 【免费下载链接】markitdown Python tool for converting files and office documents to Markdown. 项目地址: https://gitcode.com/GitHub_Trending/ma/markitdown markitdown 是一个开源的 Python…

作者头像 李华
网站建设 2026/8/29 22:38:21

大扭矩电机驱动IC怎么选?以RMC2082为例讲透选型逻辑

上周有人拿着RMC2082这款电机来问我配什么驱动IC,我第一反应是:先看规格书。对方面不改色:“规格书没在手边,你就告诉我买哪种,别啰嗦。”我当场乐了——你不孤独,我几乎每个月都能遇到几个这样问的。 驱动…

作者头像 李华
网站建设 2026/8/29 22:37:37

提示工程完整指南:4个核心技巧10分钟写出稳定提示词

提示工程完整指南:4个核心技巧10分钟写出稳定提示词 【免费下载链接】Prompt-Engineering-Guide 🐙 Guides, papers, lessons, notebooks and resources for prompt engineering, context engineering, RAG, and AI Agents. 项目地址: https://gitcode…

作者头像 李华