基于ApachePaimon全增量一体的有序流读海兰寰宇轨迹分析服务技术选型与架构演进
以下文章来源于 Apache Paimon,作者谢琛羚

Apache Paimon
Apache Paimon 中文社区官微,PMC 成员维护
摘要
海兰寰宇是一家为客户提供海洋智慧监测服务的公司。雷达网产品是一种基于雷达技术的海洋监测系统,主要用于海域监测和预警。该产品可以实现对海洋中的目标物进行追踪和识别,由此可得到船只的航行轨迹。雷达网产品的应用领域广泛,包括但不限于海上风电、港口码头、海上交通、海洋环境保护等领域。通过对海洋中的目标物进行追踪和识别,可以有效防范海上事故的发生,同时对于港口码头的管理和运营也有很大的帮助。其中大数据业务有很大一部分是对船舶轨迹数据做分析,对船舶一系列违规行为做实时报警。除了实时报警,还有对历史轨迹数据进行行为分析的需求。
本文主要讲述 Apache Paimon 是怎么通过历史数据流读的方式来实现多船历史数据轨迹分析的需求,以及技术选型与架构演进的过程。希望对有历史数据流读需求的 Paimon 用户提供一定的参考和思路。
01 业务效果



02 业务需求
我们有大量的船舶轨迹数据。原始轨迹数据由两种传感器收集:雷达和 AIS。雷达数据根据雷达设备探测,持续跟踪后,分配一个类似 UUID 的 ID。如果跟踪丢失,再次出现会生成一个新的 ID。因此雷达数据 ID 生命周期比较短。AIS 数据会有一个固定的 MMSI 号码作为唯一标识,类似汽车的车牌号。并且 MMSI 号码可以随意篡改,就像汽车的套牌。但与汽车套牌不同的是,套牌成本比汽车低得多。因此,AIS 数据中 MMSI 的生命周期长,但存在套牌问题。对这两种数据进行多传感器数据融合处理后,形成了融合轨迹数据。
有了轨迹融合数据之后,我们有了一些实时轨迹分析/报警的需求,通过对轨迹数据的计算,识别出一些船只的违规、违法或有一定特点的行为,并且实时报警。根据同时分析的船只数量分类,可分为单船轨迹分析、双船轨迹分析和多船轨迹分析。
我们前后开发的实时轨迹分析程序数量有数十个。由于以下两个原因,导致我们想要有实时程序对应的离线程序,可以对历史数据做相应的轨迹行为分析。一个原因是,实时程序是逐步开发上线的,上线之后才累积计算结果,上线之前的历史数据是没有结果的。另一个原因是,很多报警程序有丰富的规则设置项,不同的设置会有不同的结果,有拿不同规则重跑数据的需求。
船舶轨迹也是一种时序数据,所以在离线计算的时候,数据的排序很重要。单船分析比较简单,只要对数据根据 id 做 group by 或者分区,然后在分区或分组内根据时间排序,再计算即可。多船分析的时候则要求将 id/mmsi 不同的轨迹数据一起计算,并且计算时需要将多船之间的轨迹时间对齐。因此,多船历史轨迹数据分析,无法使用普通的批处理技术,例如 keyby/group by window 等方法来实现。除了要能支持离线的多船轨迹分析,我们还想实时程序和离线程序可以共用大部分代码。所以我们其实是需要一个能对历史数据快速全局排序或者能流读历史数据并且保证顺序的技术来实现多船历史轨迹分析的功能。
03 技术选型与架构演进
以下架构主要涉及船舶轨迹查询、轨迹数仓、离线轨迹分析,实时报警等业务。轨迹查询业务经历从 ElasticSearch 切换到 StarRocks/Doris 的过程,中间调研过 ClickHouse、GeoMesa、HBase。离线轨迹分析经历从 Hive 到 StarRocks/Doris 再到 Paimon 的过程。本文主要介绍离线分析数据源架构迭代到 Paimon 的历程。解决了 Hive 数据时效性的问题,数据无法根据时间全局排序的问题,无法分析长时间段数据问题,离线轨迹分析种类远少于实时报警的问题。
架构 1.0

如图所示为数据架构 1.0 架构图,由于我们的产品围绕船舶轨迹,数据和业务与一般数仓都不同,数仓的分层不是标准分层。架构比较接近 Lambda 架构。1.0 架构中的 ODS 层其实应该是 DWD 层,因为原始轨迹点的数据经过加工之后,整个系统中使用的都是加工之后的数据,所以原始数据并未保存,而是将处理之后的数据作为 DWD 层。
流处理:
原始轨迹数据发送至 Kafka 后,Flink 消费 Kafka 的数据进行数据预处理,主要是对原始数据进行清洗、过滤脏数据以及计算生成部分用于后续业务需求的通用参数指标,将预处理后的数据 Sink 至 Kafka。预处理后的数据主要用于以下几个场景:
- 通过 Flume 实时同步至 HDFS,后续加载至 Hive 中,供离线批处理场景使用;
- 针对不同的船只行为报警模型参数,利用 Flink 消费 Kafka 的数据,实时计算分析出满足不同报警模型的行为数据,推送至 Kafka,产生实时报警数据;
- Flink 实时消费轨迹数据,利用在线轨迹压缩算法,对轨迹数据进行压缩处理,并将压缩后的轨迹数据 Sink 至 ElasticSearch,供后续轨迹查询使用;
- 其他一些实时业务场景,包括:全目标实时轨迹、实时目标搜索等。
批处理:
实时数据通过 Flume 同步至 HDFS 后,通过 T+1 的方式,每天定时将数据加载至 Hive 表中。Hive ODS 层存储 Kafka 中同步过来的原始 JSON 格式数据,不对数据格式进行处理,DWD 层存储 ODS 层进行数据解析及去重后的标准化数据。批处理主要有以下场景:
- 离线轨迹分析,与实时报警类似,只不过是针对历史数据进行行为分析。由 DolphinScheduler 调度对应 Flink/Spark 离线轨迹分析任务,以 Hive DWD 层的数据做为 Source,生成分析结果直接写入 MySQL,供后端业务查询;
- 针对业务需求开发的指标统计功能,比如:目标数据量统计、船舶类型统计、报警数量统计、轨迹密度统计等。
存在的问题:
- Hive DWD 层为 T+1 数据,离线轨迹分析不支持当天数据。
- 实时报警有数十种,对应的离线轨迹分析只支持少数几种,原因有两个,一是复杂报警程序难以用 Spark SQL 批处理实现,二是 Flink on Hive 数据全局排序性能低下,超过一天的时间段,任务就没法正常执行完。
- 实时报警程序为 Flink,部分离线轨迹分析为 Spark,代码不能复用。
- 船舶轨迹存储在 Elasticsearch 中,轨迹数据量巨大,只能拿 HDD 存储,Elasticsearch 用 HDD 存储时性能很差,不能存储完整轨迹,只能存储经压缩算法压缩后的轨迹,压缩后的轨迹,查询时间段超过一个月也会经常超时。只能查询几天的轨迹。
转化为技术需求:
- 需要一个可对海量历史数据快速全局排序并的存储,或者可以按写入顺序流读的存储并且支持 Flink 读取。
- 需要一个基于 HDD 磁盘可快速查询海量轨迹数据的数据库,基于 LSM-Tree 的数据库。
架构 1.1 调研

调研 1.1 主要想解决轨迹查询性能问题,将轨迹存储由 ElasticSearch 切换为 HBase/GeoMesa。1.0 只存储了压缩之后的轨迹,经过测试 HBase/GeoMesa 查询性能很好,可以存储全量未压缩的轨迹数据,也能很快出查询结果。由于 HBase 热点问题需要调整 Row Key 加盐等,同时 api 难用,sql 查询还需要借助 Phoenix,运维使用比较麻烦所以放弃。GeoMesa 有一个优势是支持 geo 索引,作为数仓可以加速 geo 计算,但当时还不支持 Flink,并且运维麻烦、api 也不易使用,所以也放弃了。
根据这次调研,得出结论,基于 LSM-Tree 的数据库确实可以解决海量轨迹数据查询的性能问题,需要找到一种基于 LSM-Tree 且支持 SQL 的数据库。
架构 1.2

由于轨迹数据上游经常升级程序版本,需要保留数据排除各种数据问题,将原有 ODS 变更为 DWD 层,将最原始数据保存下来为新的 ods 层。
实时链路中各层数据都保存在 Kafka 中,数仓由 Hive 换为 StarRocks,轨迹数据通过 Routine Load 实时同步到 StarRocks。轨迹查询与离线轨迹分析的数据源统一成了 StarRocks。
选定 StarRocks 之前还调研过 ClickHouse 和 Doris,Doris/StarRocks 运营和使用比 ClickHouse 简单,StarRocks 当时 sql 函数比 Doris 更丰富,稳定性也好一些,所以最终选定了 StarRocks。
解决问题:
- Spark 离线轨迹分析的性能问题,性能是 Hive 的 5-15 倍。
- Flink 离线轨迹分析数据全局排序性能问题,依赖数据全局有序的程序可以正常运行。
- 轨迹查询性能满足产品要求。
存在的问题:
- Spark 离线轨迹分析程序运行实例达到一定数量,会将 StarRocks io 全部占满,导致轨迹查询服务超时不可用。
- Flink 离线轨迹分析要求全局排序的程序只能选取 1 天时间段,时间段超过 1 天数据全局排序,StarRocks 会 OOM。
转化为技术需求:
- 实时和离线资源隔离,避免离线任务影响实时业务。
- 放弃全局排序,需要一个写入有序,可按写入顺序流读历史数据的存储。
架构 2.0

新增 Apache Paimon 作为离线轨迹数据源,StarRocks 单纯作为轨迹查询数据源。Paimon 的轨迹表为 append only 模式。
解决问题:
- 所有的 Flink 实时报警程序均可以只改 source 和 sink 部分的代码就转化为离线轨迹分析程序。离线轨迹分析程序分析数据的时间段可以任选。致此 Paimon 满足我们离线轨迹分析程序的全部技术需求。
04 Paimon 流读历史数据改进
使用过程中有一些不支持的地方,社区的之信都帮忙一一开发新功能解决。
- 我们测试时往 Paimon 写实时数据的时候,同时还会拉取历史数据往同一张 paimon 表里面写,发现读取时数据会在分区之间乱序。
| 参数名 | 默认值 | Type | 描述 |
|---|---|---|---|
| scan.plan-sort-partition | false | Boolean | 是否按分区字段对计划文件进行排序 |
使用示例:
sql
SELECT * FROM dwd_extended_trajectory /*+ OPTIONS('scan.plan-sort-partition'='true') */
- 历史流读可自动停止 flink 程序
查询时指定 scan.bounded.watermark,写入 paimon 的 source 表指定 watermark
| 参数名 | 默认值 | Type | 描述 |
|---|---|---|---|
| scan.bounded.watermark | (none) | Long | 有界流模式的结束条件“水印”,当遇到较大的水印快照时,流阅读将结束。 |
使用示例:
source 表:
sql
CREATE TEMPORARY TABLE `dwd_extended_trajectory_kafka` (
`lastTm` bigint,
`lastTmLTZ` AS TO_TIMESTAMP_LTZ(lastTm, 3),
WATERMARK FOR lastTmLTZ AS lastTmLTZ
)
读取 paimon:
sql
SELECT * FROM dwd_extended_trajectory /*+ OPTIONS('scan.bounded.watermark'='1681027200000') */
Paimon 默认是全增量⼀体流读,⾸先读取全量数据,当全量阶段读取完后再流读增量数据。如果全量阶段没有读完,⽽全量数据包含了⼤于指定 scan.bounded.watermark 的数据,那么也并不会⽴即停⽌读取数据,直到全量阶段读完了,遇到了新的 snapshot,才会停⽌读取数据。所以如果想精准的停⽌流读,可以在 Select Query 中也指定下 Filter 条件,这样数据读到指定的 watermark 时,就不会继续再读取数据了。
sql
SELECT * FROM dwd_extended_trajectory /*+ OPTIONS('scan.bounded.watermark'='1681027200000') */ WHERE dt >= '2023-04-08' AND lastTm <= 1681027200000
-
将时间戳转换后设置 watermark,发现 watermark 隔 8 小时,后转换 TIMESTAMP_LTZ 类型时发现 paimon 字段不支持 TIMESTAMP_LTZ 类型。
-
多并行度读取 paimon watermark 对齐
流读历史数据时,多个 subtask 之间消费速率不同,会导致 flink window 中数据迟到,timer 也不正常。
增加以下配置开启 watermark 对齐功能:
| 参数名 | 默认值 | Type | 描述 |
|---|---|---|---|
| scan.watermark.alignment.group | (none) | String | 定义用于对齐水印源的组名 |
| scan.watermark.alignment.update-interval | 1 s | Duration | 多久进行一次水印对齐 |
| scan.watermark.alignment.max-drift | (none) | Duration | 最大允许 subtask 间水印间隔 |
并且在 paimon 表中预先定义 watermark
使用示例:
sql
CREATE TABLE IF NOT EXISTS `$tableName` (
`lastTmLTZ` TIMESTAMP_LTZ(3),
WATERMARK FOR lastTmLTZ AS lastTmLTZ
) PARTITIONED BY (dt)
读取表数据时,通过 Dynamic Option 指定上述三个参数,控制水印对齐。
sql
SELECT * FROM dwd_extended_trajectory /*+ OPTIONS('scan.watermark.alignment.group'='alignment-group','scan.watermark.alignment.max-drift'='10s','scan.watermark.alignment.update-interval'='10ms') */ WHERE dt >= '2023-04-08' AND lastTm <= 1681027200000
Paimon watermark 对齐需依赖于 flink 的 FLIP-182 FLIP-217
其中 FLIP-182 支持 source subtask 之间对齐,不支持单个 subtask 之间不同 splits 之间的对齐,FLIP-182 从 flink 1.15 开始支持,FLIP-217 可以支持单个 subtask 之间不同 splits 之间的对齐,FLIP-217 从 flink 1.17 开始支持。所以 1.15-1.16 使用 watermark 对齐的功能需要 source 并行度大于或等于 paimon 表 bucket 数,flink 1.17 后则不需要关注这个问题。
FLIP-182: Support watermark alignment of FLIP-27 Sources
FLIP-217: Support watermark alignment of source splits
其余问题:
- 流读历史数据时,程序异常重启 restore 时,如果刚好碰到 snapshot 过期,会导致找不到 bucket 内文件,导致一直无法重启。
增加 consumer-id 配置,可以阻止在读的 snapshot 过期。
但之信老师说 consumer-id 会导致 watermark 无法严格对齐,flink 1.18 可能会修复这个问题。
- 小文件问题,append-only 表每个分区最后一段数据,每个 bucket 内有几十个小文件
增加配置 full-compaction.delta-commits,几次成功 checkpoint 后触发 full-compaction。这样除了 snapshot 未过期的部分小文件外,不再有小文件。
05 流量历史数据技术需求
- 可读取全量任意历史数据
- 读取时全局有序(按写入顺序或指定字段顺序或海量数据可全局排序)
- 不同分区之间读取顺序按分区排序,不按数据写入时间排序,同一分区同一 bucket 内按写入顺序读取,可随意添加或重写分区数据,对整体排序不影响,不用按不同的时间段分不同的表。
- 支持 watermark 对齐,datastream 中 window 与 timer 正常
- 支持有界流,程序跑完指定数据可自动退出停止
未解决问题
- state 数据量暴增的问题
因为 flink state ttl 是基于 proccess time,而流读历史数据的速率是实时消费的 10-100 倍,导致 ttl 时间增大到 10-100 倍,导致 state 暴增,极其影响程序性能。
flink 社区有一个 ttl 支持 event time 语义的 issues,但已经很久没更新了。
[FLINK-12005] [State TTL] Event time support - ASF JIRA
要解决这个问题,就得自己控制 state 数据的过期。一个是注册定时器,定时清除 state 内过期数据,但是大量定时器也导致程序性能低下。
还有一个方法是根据消费速率对于现实时间的倍数来粗略放小 ttl 倍数,让 state 中过期数据快速过期,但这种做法不能保证数据准确,state 可能提前过期或延迟过期。
06 不满足需求的技术路线
在选用 Paimon 之前,我们还调研过 Apache Kafka Tiered Storage 、Apache Pulsar、Apache Iceberg 等由于以下原因,最终选择了 Paimon。
Apache Kafka Tiered Storage
- 官方分层存储还在开发未发布,Confluent 需要收费
- 数据顺序完全依赖于写入顺序,没法往同一个 topic 中补录或重写历史数据,Paimon 可基于分区名排序,处理历史数据更加灵活
Apache Pulsar
- 时间支持 event time 事件时间和 publish time 数据写入时间,publish time 不能自定义,消费历史数据时只能指定 publish time,导致如果 pulsar 内数据不是实时写入而是导入历史数据,离线任务的开始时间则没法在 pulsar 内定位到对应 event time 数据的 offset
- 存在和 kafka 一样的问题,不能灵活补录重写历史数据
Apache Iceberg
- Apache Iceberg 可通过时间旅行的方式流读历史数据,时间旅行需要指定 snapshot id,snapshot 会过期,不能流读过期之前的数据,而且需要自己维护 snapshot 与数据事件时间的映射关系才能根据时间条件快速定位到 snapshot id
07 展望
尝试解决 state ttl 的问题。
调研通过 Starrocks/Doris 读取 Apache Paimon 数据,看能否用 Paimon 做一些轨迹方面的业务。现在目前我们只使用了 Append Only Table,调研一下 Primary Key Table 能否用到一些现有业务或拓展一些新的业务。
最后,衷心感谢 Apache Paimon 开源社区的大力支持,使我们能够高效地实现项目需求。通过与 Apache Paimon 社区的紧密沟通,我们在开发过程中遇到的问题都得到了完美的解决方案。后面我们也将继续与社区保持紧密联系,为社区的发展做出我们的努力。
写在文末:
Apache Paimon (incubating) 是一项流式数据湖存储技术,可以为用户提供高吞吐、低延迟的数据摄入、流式订阅以及实时查询能力。业务中存在海量数据,有构建数据湖需求的同学,欢迎试用和讨论。

end
