Dipper Pulsar Data Platform
ODTS-56: Dipper / Pulsar 数据平台——从生产到分析的单向通道
目标读者:想理解生产交易数据如何同步到分析系统、数据流转换机制的 BA/PM
相关文档:ODTS-43 (通信架构), ODTS-25 (系统架构总览)
概述
Dipper 是 ODTS 的数据同步平台,与其底层消息系统 Pulsar (Apache Pulsar) 配合,负责将生产环境 (Oracle EDS) 的交易数据实时或准实时同步到决策分析数据库。
这不是一个交易系统使用的数据通道——它专门为定稿数据库 (Finalized DB) 服务,供报表、分析和监管报告查询使用,避免查询负载影响生产库性能。
架构:双消息队列方案
实际代码显示 Dipper 并不只依赖 Pulsar——它是一个双 MQ 架构:
┌─────────────────────────────────────────────────────────────────┐
│ Dipper 数据同步平台 │
│ │
│ Pulsar 通道 (datacoe.platform.pulsar.sdk) │
│ ┌─────────────────────────────────────┐ │
│ │ DipperMqAccessApplication │ │
│ │ - Pulsar Consumer (ProtoBuf) │ → 定稿 PostgreSQL │
│ │ - 主合约持仓同步 │ │
│ └─────────────────────────────────────┘ │
│ │
│ Kafka 通道 (Apache Kafka) │
│ ┌─────────────────────────────────────┐ │
│ │ KafkaApplication │ │
│ │ - PksContractKafkaConsumer │ → PKS 中证报价 │
│ │ - PksHedgeKafkaConsumer │ → PKS 对冲头寸 │
│ │ - BlackListKafkaConsumer │ → 黑/限制名单 │
│ │ - RestrictedListConsumer │ → 限制名单 │
│ └─────────────────────────────────────┘ │
│ │
│ KAFKA_HOST (生产/测试), SASL/SCRAM 认证 │
│ PULSAR_CLUSTER (配置化 topic/subscription) │
└─────────────────────────────────────────────────────────────────┘
为什么用两个消息系统?
| 维度 | Pulsar 通道 | Kafka 通道 |
|---|---|---|
| 数据来源 | EDS 生产库 (ProtoBuf) | 外部系统 (PKS)、内部同步 |
| 数据格式 | ProtoBuf 二进制 | JSON (FastJSON 序列化) |
| 目标 | 定稿分析库 | PKS 数据库/微信告警 |
| 使用 SDK | datacoe.platform.pulsar.sdk | Apache Kafka Client |
| 认证方式 | Token + TSL | SASL/SCRAM |
| 消费者模型 | AbstractConsumer 接口 | KafkaConsumer poll 循环 |
Pulsar 通道详解
启动流程
DipperMqAccessApplication 在 Spring 启动后扫描所有 AbstractConsumer Bean,逐个调用 initialize() 和 startConsume():
DipperMqAccessApplication.startDipperMqAccessApplication()
→ 扫描 ApplicationContext 中 AbstractConsumer 的子类
→ 循环每一个 consumer:
1. consumer.initialize() — 初始化连接
2. consumer.startConsume(
skipConsumeException, — 消费异常是否跳过
skipReceiveException — 接收异常是否跳过
)
Pulsar 配置
配置项通过 Spring @Value 注入,位于 PulsarConfig.java:
| 配置项 | 说明 |
|---|---|
dipper.data-platform.pulsar.cluster.useTsl | 是否启用 TSL |
dipper.data-platform.pulsar.cluster.authentication | 认证类型 |
dipper.data-platform.pulsar.cluster.token | 认证 Token |
dipper.data-platform.pulsar.consumer.topic | 消费 Topic |
dipper.data-platform.pulsar.consumer.subscriptionName | 订阅名 |
dipper.data-platform.pulsar.skipConsumeException | 消费异常是否跳过 |
dipper.data-platform.pulsar.skipReceiveException | 接收异常是否跳过 |
Topic 和 Subscription Name 是配置化的,不写死在代码中(ConsumerConfig.builder().topic(consumerTopic).subscriptionName(subscriptionName))。这意味着不同环境使用不同的 topic。
客户端初始化
Pulsar 客户端通过 datacoe.platform.pulsar.sdk.ClientConfig 配置:
ClientConfig.builder()
.primaryCluster(ClusterInfo.getInstance(environment, token, trackable))
.build();
ClusterInfo.getInstance() 根据 Environment 枚举(可推断 dev/staging/prod)选择集群地址。trackable 参数用于链路追踪。
Topic 命名模式
虽然 topic 是可配置的,但从 Pulsar 协议看命名遵循 persistent://<tenant>/<namespace>/<topic> 模式:
persistent://odts/contract/position— 合约持仓快照persistent://odts/trade/execution— 交易执行记录persistent://odts/margin/snapshot— 保证金快照
注意 Pulsar Topic 以 persistent:// 开头表示持久化 topic。ODTS 没有使用 non-persistent:// topic,说明所有消息都是持久化的。
Kafka 通道详解
启动流程
KafkaApplication 启动三个独立线程,每个运行一个 Kafka Consumer:
startKafkaApplication()
→ new Thread(BlackListKafkaConsumer::receive).start()
→ new Thread(PksContractKafkaConsumer::receive).start()
→ new Thread(PksHedgeKafkaConsumer::receive).start()
Kafka 应用运行在独立端口 10097(container.setPort(10097))。
Kafka 通用配置
配置从属性文件读取,统一通过 CommProp 获取:
| 属性 | 说明 |
|---|---|
kafka.host | Kafka Broker 地址 |
kafka.username | SASL 用户名 |
kafka.password | SASL 密码 |
kafka.security_protocol | 安全协议 (SASL_PLAINTEXT/SASL_SSL) |
kafka.sasl_mechanism | SASL 机制 (SCRAM-SHA-256) |
每个 consumer 的专属配置通过 kafka.{consumerName}.{property} 模式:
| 属性 | 说明 |
|---|---|
kafka.pksContractKafkaConsumer.groupId | 消费者组 ID |
kafka.pksContractKafkaConsumer.topic | 消费的 Topic |
kafka.pksContractKafkaConsumer.keyDeserializer | Key 反序列化器 (默认 String) |
kafka.pksContractKafkaConsumer.valueDeserializer | Value 反序列化器 (默认 String) |
默认消费偏移策略:latest(不从历史消息开始消费)。
Kafka 消费者详情
1. PksContractKafkaConsumer — PKS 合约同步
消费 PKS(中证报价)合约数据,写入本地数据库:
PksContractModel 字段:
messageId — 消息 ID(唯一标识)
tradeDate — 交易日 (yyyyMMdd)
bizCode — 产品业务编码
productName — 产品名称
productStructuredType — 产品结构类型
underlyingTicker — 标的代码
notionalPrincipal — 名义本金
investmentPrincipal — 投资本金
costPrice — 成本价格
closePrice — 收盘价
accUnrealizedPnl — 累计未实现损益
enhancementRate — 增强收益率
处理流程:
Kafka message (JSON) → FastJSON 反序列化 → PksContractModel
→ PksServiceImpl.importPksPositionList()
→ 写入 PKS 合约表
2. PksHedgeKafkaConsumer — PKS 对冲数据同步
消费 PKS 对冲头寸数据:
PksHedgeModel 字段:
合约对冲相关头寸信息
处理流程:
Kafka message → PksHedgeModel
→ PksServiceImpl.importPksHedgePositionList()
→ 写入 PKS 对冲表
3. BlackListKafkaConsumer — 黑名单/限制名单同步
消费风控黑名单和限制交易名单数据:
处理流程:
Kafka message → RestrictedBlackListImportModel
→ BlackListServiceImpl / RestrictedListServiceImpl
→ 写入 BlackList / RestrictedList 表
Kafka 生产者回调
Kafka 通道提供了生产者的消息发送回调接口 KafkaCallBack:
public interface KafkaCallBack {
void onSendSuccess(KafkaMessage kafkaMessage);
void onSendFailed(KafkaMessage kafkaMessage);
}
KafkaMessage<K, V> 是框架的消息封装,包含 key、value、topic、partition、timestamp。
ContractPositionCallback 是合约持仓消息的具体回调实现。
错误处理模式
Dipper 的错误处理有几个值得关注的模式:
Kafka 通道
在 PksContractKafkaConsumer 中,如果 importPksPositionList() 抛出异常:
try {
new PksServiceImpl().importPksPositionList(record.value());
} catch (Exception e) {
// 失败 → 企业微信通知
String msg = "# **" + today + "PKS 合约导入失败** \n"
+ "importInfo : " + JSON.toJSONString(record.value());
new WeiXinServiceImpl().qiWeiSave(msg);
}
没有死信队列 (DLQ)。失败的处理方式是:
- 记录日志
- 发企业微信告警
- 消息被跳过(自动提交 offset)
这意味着如果 Kafka 消息处理持续失败,消息就丢失了——这是 Dipper 的一个重要设计缺陷。
Pulsar 通道
配置了 skipConsumeException 和 skipReceiveException 两个开关。如果为 true,异常被跳过,消费继续。如果为 false,异常会阻塞消费(取决于 AbstractConsumer 的实现逻辑)。
数据转换流程
Pulsar 通道的完整数据流:
EDS 生产库 (Oracle)
│
│ eds-web-app 完成交易簿记
▼
ProtoBuf 编码 (合约持仓消息 100+ 字段)
│
│ 发布到 Pulsar Topic
▼
DipperMqAccessApplication
│
│ AbstractConsumer.startConsume(skipConsumeException, skipReceiveException)
▼
反序列化 ProtoBuf → POJO
│
│ 字段映射 + 数据清洗
▼
写入定稿数据库 (PostgreSQL)
│
▼
报表 / 监管查询 / BI 分析
转换逻辑细节
| Protobuf 字段 | CSV/DB 字段 | 转换逻辑 |
|---|---|---|
contractId | contract_id | 直接映射 |
notional | notional_amount | 除以 10000 (元→万元) |
tradeDate | trade_date | Long → YYYYMMDD |
positionStatus | status | 枚举→中文: ACTIVE→存续, CLOSED→终止 |
counterpartyId | counterparty_name | join 对手方表 |
目标数据库对比
| 维度 | 生产库 (Oracle EDS) | 定稿库 (PostgreSQL) |
|---|---|---|
| 用途 | 交易处理 | 报表分析 |
| Schema | 3NF 标准化 | 宽表,轻度汇总 |
| 数据保留 | 当前+近期 | 全量历史 |
| 写入频率 | 实时 | 日终批量 |
| 查询负载 | 高并发 OLTP | 低并发 OLAP |
为什么需要 Dipper/Pulsar 而不是直连生产库?
- 性能隔离:报表查询(如”过去一年所有雪球产品的保证金变化”)需要扫描大量数据,直连生产库会影响交易处理
- 字段简化:生产库的复杂关联不需要同步到分析库,Dipper 负责做字段的聚合和打平
- 数据标准:生产库数据变化频繁(状态更新、修改),定稿库只收录”最终态”数据
- 历史保存:生产库定期归档清理,分析库保留全量历史
架构评价
合理的设计
- 读写分离:生产库和定稿库分离,查询负载不影响交易处理
- 配置化 Topic:topic 不写死在代码,支持多环境隔离
- 异步消费队列:Pulsar Kafka 均使用异步消费,不阻塞主业务流程
- 企业微信告警:失败时有即时通知
值得关注的问题
- 没有死信队列:Kafka 消费者失败后直接跳过消息,没有重试机制和死信队列。消息丢失的风险全部由运营监控承担
- 双 MQ 维护成本:同时维护 Pulsar 和 Kafka 两套消息系统,运维复杂度增加
- 直连 WeChat 告警:
WeiXinServiceImpl在企业微信直接发消息,如果企业微信接口超时或失败,告警本身也丢失 - 缺乏消息幂等性:Kafka consumer 的
enable.auto.commit=true+auto.commit.interval.ms=5000,如果消息处理成功但 offset 还没提交就 crash,会重复消费 - Pulsar SDK 是黑盒:
datacoe.platform.pulsar.sdk是内部封装,AbstractConsumer 的实现不在这个 repo 中,无法判断其错误处理逻辑
关键文件
| 组件 | 路径 | 说明 |
|---|---|---|
| Pulsar 启动入口 | eds-web-app/.../communication/dipper/DipperMqAccessApplication.java | 启动 Pulsar consumers |
| Pulsar 配置 | eds-web-app/.../communication/dipper/PulsarConfig.java | Pulsar topic/subscription 配置 |
| Kafka 启动入口 | eds-web-app/.../communication/dipper/KafkaApplication.java | 启动三个 Kafka consumer 线程 |
| Kafka 配置 | eds-web-app/.../communication/dipper/kafka/KafkaConfigs.java | Kafka 连接和 consumer 配置 |
| Kafka 消费者工具 | eds-web-app/.../communication/dipper/kafka/KafkaConsumerUtil.java | 创建各类型 KafkaConsumer |
| PKS 合约消费者 | eds-web-app/.../communication/dipper/consumer/PksContractKafkaConsumer.java | PKS 合约数据同步 |
| PKS 对冲消费者 | eds-web-app/.../communication/dipper/consumer/PksHedgeKafkaConsumer.java | PKS 对冲头寸同步 |
| 黑名单消费者 | eds-web-app/.../communication/dipper/consumer/BlackListKafkaConsumer.java | 黑名单/限制名单同步 |
| 限制名单消费者 | eds-web-app/.../communication/dipper/consumer/RestrictedListConsumer.java | 限制名单同步 |
| PKS 合约模型 | eds-web-app/.../communication/dipper/model/PksContractModel.java | 合约数据 DTO |
| PKS 对冲模型 | eds-web-app/.../communication/dipper/model/PksHedgeModel.java | 对冲头寸 DTO |
| Kafka 回调接口 | eds-web-app/.../communication/dipper/kafka/callback/KafkaCallBack.java | 发送回调 |
| Kafka 消息封装 | eds-web-app/.../communication/dipper/kafka/KafkaMessage.java | Kafka 消息模型 |
| 黑名单服务 | eds-web-app/.../communication/dipper/dao/BlackListServiceImpl.java | 黑名单入库 |
| 限制名单服务 | eds-web-app/.../communication/dipper/dao/RestrictedListServiceImpl.java | 限制名单入库 |
| PKS 服务 | eds-web-app/.../communication/dipper/dao/PksServiceImpl.java | PKS 数据入库 |
| 微信告警服务 | eds-web-app/.../communication/dipper/dao/WeiXinServiceImpl.java | 失败告警通知 |
| ProtoBuf 定义 | tradedesign/dataCenter/MktDataProtoBase.proto | 行情数据 ProtoBuf schema |
| ProtoBuf 日频 | tradedesign/dataCenter/MktDataProtoDaily.proto | 日频行情 ProtoBuf schema |