交易类项目-flink
一、先看整体链路一笔交易相关事件可能经过类似这样的链路行情源 / 交易柜台 / 订单系统 / 账户系统 | 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 ├── T1 客户画像重算 ├── 历史盈亏修正 ├── 风控规则回溯验证 ├── 数据质量核对 └── 离线报表十二、哪些地方不应该直接依赖 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 更像实时数据和实时状态计算引擎负责把订单、成交、行情、账户行为快速转换成风控事件、资产视图、客户标签和实时指标它通常不取代撮合、清算和权威账务系统。

相关新闻

知漫剧自动分集功能实践:长篇小说转AI漫剧完整流程

知漫剧自动分集功能实践:长篇小说转AI漫剧完整流程

把长篇小说改编成漫剧,最头疼的事:几十万字怎么拆成几十集?手动拆太慢太累,拆不好节奏就崩了。知漫剧(zz.jiaxunai.cn) 价格实惠,小白上手就能一键生成,平台提供全流程教学&#xff…

2026/8/29 22:44:54 阅读更多 →
scrcpy 安卓投屏教程:5 分钟把手机画面投到电脑并用键鼠控制

scrcpy 安卓投屏教程:5 分钟把手机画面投到电脑并用键鼠控制

scrcpy 安卓投屏教程:5 分钟把手机画面投到电脑并用键鼠控制 【免费下载链接】scrcpy Display and control your Android device 项目地址: https://gitcode.com/GitHub_Trending/sc/scrcpy scrcpy 把安卓手机的屏幕和声音投屏到电脑,让你用鼠标键…

2026/8/29 22:44:54 阅读更多 →
基于SpringBoot的门店物流管理系统的设计与实现(源码+讲解视频+LW)

基于SpringBoot的门店物流管理系统的设计与实现(源码+讲解视频+LW)

温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台官方提供的学长联系方式的名片! 温馨提示:本人主页置顶文章(点我)开头有 CSDN 平台…

2026/8/29 22:44:54 阅读更多 →

最新新闻

奇安信春招前端试卷拆解:安全厂商面试到底考什么

奇安信春招前端试卷拆解:安全厂商面试到底考什么

这套试卷在我电脑里躺了挺久,我本来没打算细看。但后来事情变得有意思了——我把这套“奇安信春招前端方向试卷2”发给几个准备跳槽的朋友,结果有人直接在微信里回我一句:“这题是正经前端题吗?怎么还有安全的东西?”我…

2026/8/29 23:28:34 阅读更多 →
金山办公2020校招前端笔试题解析:从JS基础到框架底层

金山办公2020校招前端笔试题解析:从JS基础到框架底层

每年七八月份都是校招笔试高峰,前端开发的岗位尤其卷。群里经常有学弟学妹拿着各类公司真题来问我,其中金山办公这套2020校招前端开发笔试题(一),被问到的频率意外地高。按理说年份也不算近,但这套题的考察…

2026/8/29 23:28:34 阅读更多 →
STM32N657 MIPI CSI-2摄像头驱动Bring-up实战与踩坑指南

STM32N657 MIPI CSI-2摄像头驱动Bring-up实战与踩坑指南

这段时间一直在折腾 STM32N657X0H3Q 的 MIPI CSI-2 Camera Driver Bring-up,从拿到样片到最终出图,前后花了差不多两周时间,中间踩了不少坑,也把整个 MIPI CSI-2 链路的细节摸了一遍。如果你正准备在新板子上点亮摄像头&#xff0…

2026/8/29 23:28:34 阅读更多 →
HTML5实战测验:从文档骨架到Canvas动效的完整指南

HTML5实战测验:从文档骨架到Canvas动效的完整指南

这套"HTML5测验一"不是书后习题那种填空选择,而是把日常开发里真正绕不过去的点拎出来过一遍:文档骨架、播放器兼容、Canvas绘图、综合动效。我见过太多人背了标签却写不出一个能在手机上正常跑的视频页,也见过初学者一碰到心形曲线…

2026/8/29 23:27:34 阅读更多 →
HTML5测验第五期:语义化、播放器兼容与Canvas特效实战

HTML5测验第五期:语义化、播放器兼容与Canvas特效实战

最近把HTML5测验系列做到了第五期,前四期我基本都在折腾基础标签、CSS协作和整体布局,这一期我特意换了方向。大家后台问得最多的几个话题,一个是不同浏览器对HTML5播放器的支持差异,一个是Canvas特效里那些看起来简单、实际容易翻…

2026/8/29 23:27:34 阅读更多 →
纯HTML5打造混合题型交互测验应用:第五版实战与兼容性指南

纯HTML5打造混合题型交互测验应用:第五版实战与兼容性指南

先说结论:这个“HTML5测验五”不是一份静态试卷页面,是一套纯前端交互测验应用的第五个迭代版本。这版做完,我最大的感受是小项目也有小项目的讲究——一旦把题目类型扩展成单选、多选、判断和听力题混合,同时还要处理不同浏览器里…

2026/8/29 23:27:34 阅读更多 →

日新闻

etc目录下的profile.d文件目录设置环境变量和全局脚本shell

etc目录下的profile.d文件目录设置环境变量和全局脚本shell

一、设置环境变量etc目录下的profile.d文件目录 /etc/profile.d1、编写 vi test.sh文件内容# jdk变量 export ZHK_HOME/root export PATH$PATH:$ZHK_HOME/test # 可以取出来ZHK_HOME变量给ZZZ_HOME赋值 export ZZZ_HOME${ZHK_HOME}/test2、刷新 执行source /etc/profile 命令使…

2026/8/29 0:00:24 阅读更多 →
【JavaScript】内存管理-垃圾回收机制-内存泄露

【JavaScript】内存管理-垃圾回收机制-内存泄露

内存管理 C 语言这样的底层语言一般都有底层的内存管理接口,比如 malloc()和free()。 而 JavaScript 是在创建变量(对象,字符串等)时自动进行了分配内存,并且在不使用它们时“自动”释放。释放的过程称为垃圾回收。 整…

2026/8/29 0:00:24 阅读更多 →
Labgrid-MCP:为嵌入式硬件实验室接入AI Agent操控能力

Labgrid-MCP:为嵌入式硬件实验室接入AI Agent操控能力

Labgrid-MCP 的目标是把 MCP(Model Context Protocol)能力延伸到真实嵌入式硬件实验室:AI Agent 通过一个标准化的 MCP Server,就能查看目标板状态、控制上电断电、复位开发板、读取串口日志,甚至执行镜像刷写。对于经…

2026/8/29 0:00:24 阅读更多 →

周新闻

[光学原理与应用-521]:对光的错误理解与纠偏

[光学原理与应用-521]:对光的错误理解与纠偏

首先光是一种能量的载体和形态,宏观上观察到的光是由无数个微观的光量子组成的,每个光子在产生的瞬间,其在真空的空间中以确定不变的速度沿着一个初始的方向一直向前,在微观层面,每个光量子的运动轨迹是以波函数所展现…

2026/8/29 18:08:35 阅读更多 →
SIP通话转接原理与REFER方法实战解析

SIP通话转接原理与REFER方法实战解析

1. 通话转接不是“挂断再拨号”,而是SIP会话的动态重定向你有没有遇到过这样的场景:客服坐席A正在和客户通电话,突然需要把这通对话无缝转给专家坐席B,客户完全感知不到中间的断连——既没听到忙音,也没被要求重新拨号…

2026/8/28 23:05:07 阅读更多 →
Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

Kolla-ansible单节点OpenStack部署实战:从环境准备到排坑指南

1. 为什么选择Kolla-ansible来部署单节点OpenStack?如果你正在寻找一种能把OpenStack从“概念”快速变成“可用的实验环境”的方法,那么Kolla-ansible几乎是当前最主流、最省心的选择。我见过太多人卡在手动编译依赖、配置服务、处理版本冲突的泥潭里&am…

2026/8/28 19:47:53 阅读更多 →

月新闻

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南

免费解锁百度网盘SVIP加速:macOS用户必备的下载提速终极指南 【免费下载链接】BaiduNetdiskPlugin-macOS For macOS.百度网盘 破解SVIP、下载速度限制~ 项目地址: https://gitcode.com/gh_mirrors/ba/BaiduNetdiskPlugin-macOS 还在为百度网盘macOS版的龟速下…

2026/8/29 4:34:53 阅读更多 →
终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换

终极ncmdump指南:3分钟实现网易云NCM音乐解密与格式转换 【免费下载链接】ncmdump 项目地址: https://gitcode.com/gh_mirrors/ncmd/ncmdump 还在为网易云音乐下载的NCM格式文件无法在其他播放器播放而烦恼吗?ncmdump解密工具帮你轻松解决这个困…

2026/8/28 17:43:04 阅读更多 →
HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

HarmonyOS 应用开发《掌上英语》第81篇: 智能体卡片:为英语学习 App 打造桌面级学习助手

AgentCard 智能体卡片:为英语学习 App 打造桌面级学习助手适用平台:HarmonyOS 7.0 (API 26 Beta)一、引言 HarmonyOS 7.0(API 26 Beta)新增了 AgentCard 智能体卡片能力,这是继 HMAF(鸿蒙智能体框架&#x…

2026/8/29 2:05:18 阅读更多 →