Learning
VOL. VII · NO. 106 · OTC Derivatives · 19 JUL 2026

Dipper Pulsar Data Platform

OTC 衍生品 · 19 JUL 2026 · 10 min read · 1,682 words
· · ·

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 数据库/微信告警
使用 SDKdatacoe.platform.pulsar.sdkApache Kafka Client
认证方式Token + TSLSASL/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.hostKafka Broker 地址
kafka.usernameSASL 用户名
kafka.passwordSASL 密码
kafka.security_protocol安全协议 (SASL_PLAINTEXT/SASL_SSL)
kafka.sasl_mechanismSASL 机制 (SCRAM-SHA-256)

每个 consumer 的专属配置通过 kafka.{consumerName}.{property} 模式:

属性说明
kafka.pksContractKafkaConsumer.groupId消费者组 ID
kafka.pksContractKafkaConsumer.topic消费的 Topic
kafka.pksContractKafkaConsumer.keyDeserializerKey 反序列化器 (默认 String)
kafka.pksContractKafkaConsumer.valueDeserializerValue 反序列化器 (默认 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 通道

配置了 skipConsumeExceptionskipReceiveException 两个开关。如果为 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 字段转换逻辑
contractIdcontract_id直接映射
notionalnotional_amount除以 10000 (元→万元)
tradeDatetrade_dateLong → YYYYMMDD
positionStatusstatus枚举→中文: ACTIVE→存续, CLOSED→终止
counterpartyIdcounterparty_namejoin 对手方表

目标数据库对比

维度生产库 (Oracle EDS)定稿库 (PostgreSQL)
用途交易处理报表分析
Schema3NF 标准化宽表,轻度汇总
数据保留当前+近期全量历史
写入频率实时日终批量
查询负载高并发 OLTP低并发 OLAP

为什么需要 Dipper/Pulsar 而不是直连生产库?

  1. 性能隔离:报表查询(如”过去一年所有雪球产品的保证金变化”)需要扫描大量数据,直连生产库会影响交易处理
  2. 字段简化:生产库的复杂关联不需要同步到分析库,Dipper 负责做字段的聚合和打平
  3. 数据标准:生产库数据变化频繁(状态更新、修改),定稿库只收录”最终态”数据
  4. 历史保存:生产库定期归档清理,分析库保留全量历史

架构评价

合理的设计

  • 读写分离:生产库和定稿库分离,查询负载不影响交易处理
  • 配置化 Topic:topic 不写死在代码,支持多环境隔离
  • 异步消费队列:Pulsar Kafka 均使用异步消费,不阻塞主业务流程
  • 企业微信告警:失败时有即时通知

值得关注的问题

  1. 没有死信队列:Kafka 消费者失败后直接跳过消息,没有重试机制和死信队列。消息丢失的风险全部由运营监控承担
  2. 双 MQ 维护成本:同时维护 Pulsar 和 Kafka 两套消息系统,运维复杂度增加
  3. 直连 WeChat 告警WeiXinServiceImpl 在企业微信直接发消息,如果企业微信接口超时或失败,告警本身也丢失
  4. 缺乏消息幂等性:Kafka consumer 的 enable.auto.commit=true + auto.commit.interval.ms=5000,如果消息处理成功但 offset 还没提交就 crash,会重复消费
  5. 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.javaPulsar topic/subscription 配置
Kafka 启动入口eds-web-app/.../communication/dipper/KafkaApplication.java启动三个 Kafka consumer 线程
Kafka 配置eds-web-app/.../communication/dipper/kafka/KafkaConfigs.javaKafka 连接和 consumer 配置
Kafka 消费者工具eds-web-app/.../communication/dipper/kafka/KafkaConsumerUtil.java创建各类型 KafkaConsumer
PKS 合约消费者eds-web-app/.../communication/dipper/consumer/PksContractKafkaConsumer.javaPKS 合约数据同步
PKS 对冲消费者eds-web-app/.../communication/dipper/consumer/PksHedgeKafkaConsumer.javaPKS 对冲头寸同步
黑名单消费者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.javaKafka 消息模型
黑名单服务eds-web-app/.../communication/dipper/dao/BlackListServiceImpl.java黑名单入库
限制名单服务eds-web-app/.../communication/dipper/dao/RestrictedListServiceImpl.java限制名单入库
PKS 服务eds-web-app/.../communication/dipper/dao/PksServiceImpl.javaPKS 数据入库
微信告警服务eds-web-app/.../communication/dipper/dao/WeiXinServiceImpl.java失败告警通知
ProtoBuf 定义tradedesign/dataCenter/MktDataProtoBase.proto行情数据 ProtoBuf schema
ProtoBuf 日频tradedesign/dataCenter/MktDataProtoDaily.proto日频行情 ProtoBuf schema