一、业务背景券商希望实时了解客户的交易和行为特征用于客户分层和精细化运营投顾服务推荐客户活跃度分析风险提示和合规监控流失客户识别例如系统需要识别出“客户张三最近 30 天频繁交易科技类股票风险等级为中高近 7 天活跃度下降但账户资产较高。”这就是一条比较完整的用户画像。二、数据来源券商的数据可能来自多个系统行情系统 - 股票行情、涨跌幅、成交量 交易系统 - 委托、成交、撤单、持仓、资金变动 App/网页 - 登录、浏览、搜索、收藏、自选股 客户系统 - 年龄、地区、客户等级、风险测评 客服系统 - 咨询记录、投诉记录、服务记录实时数据可以先进入 Kafka交易系统 / App / 行情系统 | v Kafka | v Flink三、要构建哪些画像标签可以将标签分成几类。1. 基础属性标签这类标签变化不频繁通常来自客户主数据或 CRM客户等级所在地区开户时间风险承受等级资产区间2. 活跃度标签例如最近登录时间近 7 天登录次数近 30 天交易天数最近一次交易距今天数是否为沉睡客户3. 交易行为标签例如近 7 天交易次数近 30 天成交金额买入金额与卖出金额平均持仓天数撤单率日内交易占比4. 投资偏好标签例如偏好股票、基金还是债券偏好的行业板块偏好的市值类型偏好高波动还是低波动资产是否偏好长期持有5. 风险和运营标签例如高频交易客户高资产低活跃客户风险等级与实际交易品种不匹配可能流失客户适合某类投教内容的客户这里需要注意券商画像不能只用于营销还必须结合适当性管理和合规要求。画像结果不能直接替代客户风险测评也不能绕过投资者适当性规则。四、实时处理流程以“近 7 天客户交易金额”为例成交事件进入 Kafka | v Flink 读取成交数据 | v 字段校验、去重、异常数据过滤 | v 按照 customer_id 分组 | v 基于事件时间做 7 天窗口聚合 | v 生成客户交易画像标签 | -- Doris / ClickHouse实时查询 -- Iceberg明细和历史数据沉淀 -- Kafka供其他系统继续消费 -- Redis低延迟画像查询五、一个具体例子假设 Kafka 中有成交事件{customer_id:C1001,order_id:O9001,symbol:600000.SH,industry:金融,side:BUY,amount:50000,event_time:2025-01-10 10:15:00}Flink 对客户进行统计客户 C1001 近 7 天成交次数18 次 近 7 天成交金额320000 元 近 30 天交易天数12 天 买入金额220000 元 卖出金额100000 元 偏好行业金融、科技 最近交易时间2025-01-10 10:15:00然后根据规则生成标签交易活跃度 高 交易频率 高频 资金规模 中等 投资偏好 股票型、金融行业偏好 客户状态 正常活跃六、Flink SQL 示例定义成交数据表CREATETABLEtrades(customer_id STRING,order_id STRING,symbol STRING,industry STRING,side STRING,amountDECIMAL(18,2),event_timeTIMESTAMP(3),WATERMARKFORevent_timeASevent_time-INTERVAL5SECOND,PRIMARYKEY(customer_id,order_id)NOTENFORCED)WITH(connectorkafka,topicbroker-trades,formatjson);统计每个客户近 7 天的交易情况可以使用 Hop 窗口SELECTcustomer_id,window_start,window_end,COUNT(*)AStrade_count,SUM(amount)AStotal_amount,SUM(CASEWHENsideBUYTHENamountELSE0END)ASbuy_amount,SUM(CASEWHENsideSELLTHENamountELSE0END)ASsell_amountFROMTABLE(HOP(TABLEtrades,DESCRIPTOR(event_time),INTERVAL1DAY,INTERVAL7DAY))GROUPBYcustomer_id,window_start,window_end;含义是每天更新一次统计每个客户最近 7 天的数据。也可以按照行业统计客户偏好SELECTcustomer_id,industry,SUM(amount)ASindustry_amountFROMtradesGROUPBYcustomer_id,industry;实际生产中通常还需要进一步排序取每个客户交易金额最高的行业作为主要偏好行业。七、客户流失标签怎么计算例如定义近 30 天没有交易 并且近 90 天有过交易 并且客户资产大于某个阈值可以形成“高资产低活跃”标签。处理逻辑是客户登录、交易事件 - Flink 更新客户最近行为时间 | v 定时或事件驱动判断是否超过 30 天未交易 | v 生成客户活跃状态标签如果客户重新交易Flink 再把他的状态从“可能流失”更新成“活跃”。这里可以使用 Keyed State 保存Keycustomer_id State - last_trade_time - last_login_time - trade_count_7d - total_amount_30d - preferred_industry八、架构中的各组件职责Kafka接收交易、登录、浏览等实时事件 Flink清洗、去重、关联、窗口统计、更新画像 MySQL/维表客户基础信息、风险等级等相对稳定数据 Iceberg保存客户行为明细和画像历史快照 Doris/ClickHouse支持运营和分析人员实时查询 Redis支持画像服务低延迟读取 Hive/Spark SQL做离线历史分析和标签回溯九、面试官可能追问的问题1. 交易数据重复怎么办使用order_id或trade_id做幂等键在 Flink 中去重下游写入也要设计幂等主键不能只依赖 Flink 内部 Exactly-Once。2. 交易数据乱序怎么办使用事件时间和 Watermark。比如允许 5 秒乱序晚到数据根据业务需求进行修正、旁路输出或丢弃。3. 为什么需要 CheckpointCheckpoint 保存 Flink 的状态和 Kafka 消费位点。任务失败后可以恢复客户统计状态避免从头计算。4. 客户数量很大状态会不会爆会因此需要使用 State TTL、窗口自动清理和合理的状态模型。Flink 状态只保存实时计算必要的信息历史明细写入 Iceberg不把所有历史交易永久放在状态中。5. 一个客户交易特别多出现数据倾斜怎么办先识别热点客户或异常 Key。对热点 Key 可以做分片聚合先按customer_id bucket做局部聚合再按customer_id做二次聚合。对异常客户也可以单独处理。6. 实时画像结果写入数据库失败怎么办使用支持幂等或事务的 Sink配合重试、死信队列和监控。对于可更新的画像结果可以使用customer_id tag_code作为唯一键重复写入不会产生重复记录。十、面试时可以直接这样总结我理解的券商用户画像系统是将交易、行情、App 行为和客户基础信息等数据接入 Kafka由 Flink 进行清洗、去重、维表关联和基于事件时间的窗口聚合按 customer_id 维护客户的活跃度、交易频率、资产变化和投资偏好等标签。实时结果可以写入 Doris 或 Redis 供查询明细和历史结果写入 Iceberg供 Hive 和 Spark SQL 做离线分析。系统需要重点处理数据乱序、迟到、重复、状态膨胀、数据倾斜、Checkpoint 恢复和下游幂等问题同时遵守券商的数据安全、权限和投资者适当性要求。